Commit 5f50cb3d by huangfusuper

增加时间参数可动态更改

parent 78e7837d
......@@ -3,6 +3,7 @@ package com.byit.api;
import com.byit.dto.executor.KillDto;
import com.byit.dto.plugin.CollectData;
import com.byit.dto.plugin.StopFlowParam;
import com.byit.dto.specials.RepairFlow;
import com.byit.dto.web.ResponseResult;
import com.byit.model.RunRecording;
import com.byit.model.vo.RunRecordingVo;
......@@ -174,13 +175,13 @@ public class ApiFlowController {
/**
* 补批工作流
*
* @param param
* @param repairFlow
* @return
*/
@PostMapping("/repairFlow")
@ApiOperation("补批工作流")
public ResponseResult repairFlow(String param) {
apiFlowService.repairFlow(param);
public ResponseResult repairFlow(@RequestBody RepairFlow repairFlow) {
apiFlowService.repairFlow(repairFlow);
return ResponseResult.ok("SUCCESS");
}
......
......@@ -3,6 +3,7 @@ package com.byit.service;
import com.byit.dto.plugin.CollectData;
import com.byit.dto.plugin.PluginFlow;
import com.byit.dto.plugin.StopFlowParam;
import com.byit.dto.specials.RepairFlow;
import com.byit.model.RunRecording;
import com.byit.model.vo.RunRecordingVo;
......@@ -77,7 +78,7 @@ public interface ApiFlowService {
* 补批工作流
* @param param
*/
void repairFlow(String param);
void repairFlow(RepairFlow repairFlow);
/**
......
......@@ -9,6 +9,7 @@ import com.byit.enums.NodeTypeEnum;
import com.byit.enums.PlaceholderEnum;
import com.byit.enums.ScheduleTypeEnum;
import com.byit.job.utils.CurrentUserUtils;
import com.byit.job.utils.DateUtil;
import com.byit.job.utils.PlaceholderUtils;
import com.byit.mapper.RunRecordingMapper;
import com.byit.mapper.WaitingRecordMapper;
......@@ -38,9 +39,7 @@ import java.util.stream.Collectors;
@Service
public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
private final static SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMdd");
public static final String PRE = "${";
public static final String SUFFER = "}";
@Value("${server.port}")
private Integer serverPort;
......@@ -87,6 +86,10 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
//补批的类型 跑不跑本节点 startNodeTaskName
String supplementStatus = specialJobParam.getSupplementStatus();
ValidationUtil.dataNotBank(supplementStatus, "补批的类型不允许为空");
String timeFormatName = specialJobParam.getTimeFormatName();
ValidationUtil.dataNotBank(timeFormatName, "补批的时间类型不允许为空");
String nowTimeFormatName = specialJobParam.getNowDateTimeFormatName();
ValidationUtil.dataNotBank(nowTimeFormatName, "nowDate时间类型不允许为空");
//获取公共参数
Map<String, String> publicParam = specialJobParam.getPublicParam();
//补批时间
......@@ -131,7 +134,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
List<WaitingRecord> waitingRecords = buildWaitingRecordList(flowByName, runId, nodeRelyDtoSet.size(), operator, repairTimeList);
waitingRecords.forEach(waitingRecordMapper::insertSelective);
List<RunRecording> runRecordings = buildRunRecordingList(waitingRecords);
List<WaitingTask> waitingTaskList = buildWaitingTaskResult(waitingRecords, nodeRelyDtoSet,specialJobParam.getPublicParam());
List<WaitingTask> waitingTaskList = buildWaitingTaskResult(waitingRecords, nodeRelyDtoSet,specialJobParam.getPublicParam(),timeFormatName, nowTimeFormatName);
runRecordings.forEach(runRecordingMapper::saveRunRecording);
waitingTaskList.forEach(waitingTaskMapper::insertSelective);
}
......@@ -146,7 +149,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
* @param nodeRelyDtoSet 节点对象
* @return 返回一个等待队的全部节点
*/
private List<WaitingTask> buildWaitingTask(WaitingRecord waitingRecord, Set<NodeRelyDto> nodeRelyDtoSet,Map<String,String> publicParam) {
private List<WaitingTask> buildWaitingTask(WaitingRecord waitingRecord, Set<NodeRelyDto> nodeRelyDtoSet,Map<String,String> publicParam, String timeFormat, String nowDateFormat) {
return nodeRelyDtoSet.stream().map(nodeRelyDto -> {
Node node = nodeRelyDto.getNode();
String relyId = nodeRelyDto.getRelyId();
......@@ -163,31 +166,28 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
waitingTask.setScheduleType(ScheduleTypeEnum.REPAIR.getCode());
waitingTask.setWaitId(waitingRecord.getWaitId());
try {
RunParamWrapped runParamWrapped = JSON.parseObject(waitingTask.getRunParam(), RunParamWrapped.class);
runParamWrapped.setPublicParamMap(publicParam);
try {
SimpleDateFormat sdf = new SimpleDateFormat(timeFormat);
Date repairDate = sdf.parse(repeatTime);
ValidationUtil.isTrueValidation(!repairDate.before(new Date()), "只能补过去时间的批次!");
} catch (ParseException e) {
log.error("补批日期不符合规范,例:20200101");
ValidationUtil.isTrueValidation(true, "补批日期不符合规范,例:20200101");
log.error("补批日期不符合规范,例:{}",timeFormat);
ValidationUtil.isTrueValidation(true, "补批日期不符合规范,例:"+timeFormat);
}
RunParamWrapped runParamWrapped = JSON.parseObject(waitingTask.getRunParam(), RunParamWrapped.class);
runParamWrapped.setPublicParamMap(publicParam);
//补批只替换不是java的节点
if (!NodeTypeEnum.JAVA.getCode().equals(waitingTask.getJobType())) {
if (StringUtils.isNotEmpty(waitingTask.getRunParam())) {
String param = PlaceholderUtils.paramPlaceholder(runParamWrapped.getPrivateParam(), repeatTime);
String param = PlaceholderUtils.formatBizDateParam(runParamWrapped.getPrivateParam(), PlaceholderEnum.DATE_PLACEHOLDER.getName() ,repeatTime, 0);
String dateFormat = DateUtil.dateFormat(new Date(), nowDateFormat);
param = PlaceholderUtils.formatBizDateParam(param, PlaceholderEnum.DATE_PLACEHOLDER.getName() ,dateFormat, 0);
//设置替换完成后的参数
runParamWrapped.setPrivateParam(param);
}
//获取原始命令
String runCommand = waitingTask.getRunCommand();
//替换参数命令
String replaceRunCommand = runCommand.replace(PRE + PlaceholderEnum.DATE_PLACEHOLDER.getName() + SUFFER, repeatTime);
waitingTask.setRunCommand(replaceRunCommand);
}
waitingTask.setRunParam(JSON.toJSONString(runParamWrapped));
return waitingTask;
......@@ -201,10 +201,11 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
* @param nodeRelyDtoSet 节点集合
* @return 全部的节点信息
*/
private List<WaitingTask> buildWaitingTaskResult(List<WaitingRecord> waitingRecords, Set<NodeRelyDto> nodeRelyDtoSet,Map<String,String> publicPram) {
private List<WaitingTask> buildWaitingTaskResult(List<WaitingRecord> waitingRecords, Set<NodeRelyDto> nodeRelyDtoSet,
Map<String,String> publicPram, String timeFormat, String nowDateFormat) {
List<WaitingTask> waitingTaskListResult = new ArrayList<>(8);
waitingRecords.forEach(waitingRecord -> {
List<WaitingTask> waitingTaskList = buildWaitingTask(waitingRecord, nodeRelyDtoSet,publicPram);
List<WaitingTask> waitingTaskList = buildWaitingTask(waitingRecord, nodeRelyDtoSet, publicPram, timeFormat, nowDateFormat);
waitingTaskListResult.addAll(waitingTaskList);
});
return waitingTaskListResult;
......
......@@ -9,6 +9,7 @@ import com.byit.dto.StatisticsConditionDto;
import com.byit.dto.api.DeleteDto;
import com.byit.dto.executor.RunParamWrapped;
import com.byit.dto.plugin.*;
import com.byit.dto.specials.RepairFlow;
import com.byit.enums.*;
import com.byit.enums.plugin.PluginNodeTypeEnum;
import com.byit.job.utils.CronExpression;
......@@ -52,7 +53,6 @@ import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
@Service
@Transactional(rollbackFor = Exception.class)
public class ApiFlowServiceImpl implements ApiFlowService {
private final static SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMdd");
public static final String PRE = "${";
public static final String SUFFER = "}";
......@@ -1283,25 +1283,25 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
@Override
public void repairFlow(String param) {
ValidationUtil.dataNotBank(param, "请求参数不允许为空!");
JSONObject jsonObject = JSON.parseObject(param);
public void repairFlow(RepairFlow repairFlow) {
ValidationUtil.dataNotNull(repairFlow, "请求参数不允许为空!");
//获取工作空间名称
String workspaceName = jsonObject.getString("workspaceName");
String workspaceName = repairFlow.getWorkspaceName();
ValidationUtil.dataNotBank(workspaceName, "工作空间名称不允许为空!");
//获取工作流名称
String flowName = jsonObject.getString("flowName");
String flowName = repairFlow.getFlowName();
ValidationUtil.dataNotBank(flowName, "工作流名称不允许为空!");
String nodeNames = jsonObject.getString("nodeNames");
ValidationUtil.dataNotBank(nodeNames, "补批节点不允许为空!");
List<String> nodeNameList = Arrays.asList(nodeNames.split(","));
List<String> nodeNameList = repairFlow.getNodeNames();
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(nodeNameList), "补批节点不允许为空!");
//获取补批的日期
String repairTimes = jsonObject.getString("repairTimes");
ValidationUtil.dataNotBank(repairTimes, "补批日期不允许为空!");
List<String> repairTimeList = Arrays.asList(repairTimes.split(","));
List<String> repairTimeList = repairFlow.getRepairTimes();
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(repairTimeList), "补批日期不允许为空!");
String timeFormatName = repairFlow.getTimeFormatName();
ValidationUtil.dataNotBank(timeFormatName, "补批时间格式不允许为空!");
String nowDateTimeFormatName = repairFlow.getNowDateTimeFormatName();
ValidationUtil.dataNotBank(nowDateTimeFormatName, "nowDate时间格式不允许为空!");
//开始校验
Workspace workspace = workspaceMapper.getByName(workspaceName);
ValidationUtil.dataNotNull(workspace, workspaceName + "工作空间不存在");
......@@ -1315,7 +1315,8 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//设置触发时间
Long triggerTime = System.currentTimeMillis();
Map<String, List<WaitingTask>> waitTaskMap = buildTask(nodeList, nodeNameList, repairTimeList, triggerTime, flowName);
Map<String, List<WaitingTask>> waitTaskMap = buildTask(nodeList, nodeNameList, repairTimeList, triggerTime, flowName,
timeFormatName, nowDateTimeFormatName);
//将时间排序
Collections.sort(repairTimeList);
//查看当前排队的工作流最大排队序号
......@@ -1362,7 +1363,8 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
private Map<String, List<WaitingTask>> buildTask(List<Node> nodeList, List<String> nodeNameList, List<String> repairTimeList, Long triggerTime, String flowName) {
private Map<String, List<WaitingTask>> buildTask(List<Node> nodeList, List<String> nodeNameList, List<String> repairTimeList,
Long triggerTime, String flowName, String timeFormatName, String nowDateTimeFormatName) {
Map<String, List<WaitingTask>> result = new HashMap<>();
String usernName = currentUserUtils.account();
List<WaitingTask> waitingTaskList = new ArrayList<>();
......@@ -1393,13 +1395,14 @@ public class ApiFlowServiceImpl implements ApiFlowService {
for (String repairTime : repairTimeList) {
List<WaitingTask> waitingTaskDateList = new ArrayList<>();
try {
SimpleDateFormat sdf = new SimpleDateFormat(timeFormatName);
Date repairDate = sdf.parse(repairTime);
ValidationUtil.isTrueValidation(!repairDate.before(new Date()), "只能补过去时间的批次!");
} catch (ParseException e) {
log.error("补批日期不符合规范,例:20200101");
ValidationUtil.isTrueValidation(true, "补批日期不符合规范,例:20200101");
log.error("补批日期不符合规范,例:{}", timeFormatName);
ValidationUtil.isTrueValidation(true, "补批日期不符合规范,例:" + timeFormatName);
}
//TODO 替换参数
waitingTaskList.forEach(waitingTask -> {
//补批只替换不是java的节点
if (!NodeTypeEnum.JAVA.getCode().equals(waitingTask.getJobType())) {
......@@ -1407,12 +1410,11 @@ public class ApiFlowServiceImpl implements ApiFlowService {
String runParamSource = waitingTask.getRunParam();
RunParamWrapped runParamWrapped = JSON.parseObject(runParamSource, RunParamWrapped.class);
//替换私有参数
String param = PlaceholderUtils.paramPlaceholder(runParamWrapped.getPrivateParam(), repairTime);
String param = PlaceholderUtils.formatBizDateParam(runParamWrapped.getPrivateParam(), PlaceholderEnum.DATE_PLACEHOLDER.getName(), repairTime, 0);
String dateFormat = com.byit.job.utils.DateUtil.dateFormat(new Date(), nowDateTimeFormatName);
param = PlaceholderUtils.formatBizDateParam(param, PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName(), dateFormat, 0);
runParamWrapped.setPrivateParam(param);
param = JSON.toJSONString(runParamWrapped);
String runCommand = waitingTask.getRunCommand();
runCommand.replace(PRE + PlaceholderEnum.DATE_PLACEHOLDER.getName() + SUFFER, repairTime);
waitingTask.setRunCommand(runCommand);
waitingTask.setRunParam(param);
}
}
......
......@@ -75,18 +75,19 @@ public class ApiNodeServiceImpl implements ApiNodeService {
ValidationUtil.dataNotBank(runNode.getNodeId(), "节点id不允许为空!");
String runId = "REAL:EXEC:" + runNode.getNodeId();
ValidationUtil.dataNotBank(runNode.getNodeName(), "节点名称不允许为空!");
//ValidationUtil.dataNotBank(runNode.getRunCmd(), "运行命令不允许为空!");
ValidationUtil.dataNotBank(runNode.getJobType(), "节点类型不允许为空!");
Long triggerTime = System.currentTimeMillis();
long triggerTime = System.currentTimeMillis();
JobTaskSchedule schedule = new JobTaskSchedule();
//获取运行参数
String runParam = runNode.getRunParam();
RunParamWrapped runParamWrapped = new RunParamWrapped();
//包装公共参数和私有参数
runParamWrapped.setPrivateParam(runParam);
runParamWrapped.setPublicParamMap(runNode.getPublicParam());
//开始构建 schedule
schedule.setRunParam(JSON.toJSONString(runParamWrapped));
schedule.setNodeId(Integer.valueOf(runNode.getNodeId()));
schedule.setRunId(runId);
schedule.setNodeName(runNode.getNodeName());
......@@ -101,46 +102,32 @@ public class ApiNodeServiceImpl implements ApiNodeService {
schedule.setScheduleType(ScheduleTypeEnum.REAL.getCode());
schedule.setHandlerName(runNode.getHandlerName());
String runCommand = schedule.getRunCommand();
if (StringUtils.isNotBlank(runCommand)) {
runCommand = runCommand.replace("${"+ PlaceholderEnum.DATE_PLACEHOLDER.getName() +"}", DateUtil.dateLessDayStr(new Date(),"yyyyMMdd",1));
schedule.setRunCommand(runCommand);
}
//构建日志节点
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
BeanUtils.copyProperties(schedule, jobTaskRunLog);
jobTaskRunLog.setStartTime(new Date());
jobTaskRunLog.setTriggerTime(new Date());
jobTaskRunLog.setRunParams(runNode.getRunParam());
jobTaskRunLog.setRunCommand(runNode.getRunCmd());
//初步保存日志节点 当前日志为不完善的日志信息
jobTaskRunLogService.saveJobTaskRunLog(jobTaskRunLog);
//获取日志id
Integer logId= jobTaskRunLog.getLogId();
//设置日志id
schedule.setLogId(logId);
//构建等待日志
RunLog runLog = RunLog.builder().isEnd(false).runLog("等待服务器分配资源").build();
//生成本次运行标识
String runKey = IDGenerationStrategy.keyGenerationStrategy(logId);
//推送日志
stringRedisTemplate.opsForList().rightPush(runKey, JSON.toJSONString(runLog, WriteClassName));
TimerTask timerTask;
if (NodeTypeEnum.JAVA.getCode().equals(runNode.getJobType())) {
timerTask = new JavaNodeExecutorTask(schedule);
}else{
String privateParam = runParamWrapped.getPrivateParam();
String formatParam = PlaceholderUtils.formatParam(privateParam);
runParamWrapped.setPrivateParam(formatParam);
schedule.setRunParam(JSON.toJSONString(runParamWrapped));
timerTask = new ScriptExecutorJobTask(schedule);
}
TimerTask timerTask = new ScriptExecutorJobTask(schedule);
//添加到调度轮
WorkRoulette.addJob(timerTask, triggerTime);
RunLog runLogEnd = RunLog.builder().isEnd(false).runLog("服务器分配资源成功").build();
stringRedisTemplate.opsForList().rightPush(runKey, JSON.toJSONString(runLogEnd, WriteClassName));
return runKey;
}
......@@ -236,7 +223,7 @@ public class ApiNodeServiceImpl implements ApiNodeService {
javaTask.setJobName(javaTask.getTaskName());
}
javaTask.setTriggerTime(0L);
if (null != javaTask.getAlarmlAction() && "0".equals(javaTask.getRepeatCount())){
if (!"0".equals(javaTask.getAlarmlAction())){
ValidationUtil.dataNotBank(javaTask.getAlarmEmail(), "设置为告警时告警邮箱不允许为空");
}
TimerTask timerTask = new JavaTaskJobTask(javaTask);
......
package com.byit.service.mapservice.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.fastjson.JSON;
import com.byit.dto.executor.RunParamWrapped;
import com.byit.enums.EmailEnum;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.PlaceholderEnum;
import com.byit.enums.RunRecordingEnum;
import com.byit.job.utils.CronExpression;
import com.byit.job.utils.DateUtil;
import com.byit.job.utils.PlaceholderUtils;
import com.byit.model.*;
import com.byit.service.*;
......@@ -55,11 +59,19 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
jobTaskService.removeMythJobTaskById(jobTask.getId());
//添加任务日志
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
//执行快速失败时 将参数替换强行替换进去
int cutBackDay = 1;
if (jobTask.getScheduleType() != 1) {
cutBackDay = 0;
}
//执行快速失败时 将参数替换强行替换进去 替换的是私有参数
String runParam = jobTask.getRunParam();
if (StringUtils.isNotBlank(runParam)) {
String newParam = PlaceholderUtils.formatParam(runParam);
jobTaskRunLog.setRunParams(newParam);
RunParamWrapped runParamWrapped = JSON.parseObject(runParam, RunParamWrapped.class);
String privateParam = runParamWrapped.getPrivateParam();
String newParam = PlaceholderUtils.formatBizDateParam(privateParam, PlaceholderEnum.DATE_PLACEHOLDER.getName(), cutBackDay);
newParam = PlaceholderUtils.formatBizDateParam(newParam, PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName(), 0);
runParamWrapped.setPrivateParam(newParam);
jobTaskRunLog.setRunParams(JSON.toJSONString(runParamWrapped));
}
BeanUtils.copyProperties(jobTask, jobTaskRunLog);
......
......@@ -59,7 +59,6 @@ public class JavaTaskJobTask implements TimerTask {
hasRelatedFlow = javaJobConfDTO.getHasRelatedFlow();
}
//保存到日志
Integer logId = saveLog(javaTask);
GetRegConfig getRegConfig = SpringUtil.getBean(GetRegConfig.class);
......
......@@ -13,6 +13,7 @@ import com.byit.enums.task.RunResultEnum;
import com.byit.enums.task.RunTypeEnum;
import com.byit.factory.PluginClientFactory;
import com.byit.factory.PluginSpringClientFactory;
import com.byit.job.utils.DateUtil;
import com.byit.job.utils.PlaceholderUtils;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.model.JobTaskSchedule;
......@@ -98,6 +99,13 @@ public class ScriptExecutorJobTask implements TimerTask {
if(runParamWrapped != null) {
//获取到私有参数
String privateParam = runParamWrapped.getPrivateParam();
//替换运行参数中的时间参数
if (ScheduleTypeEnum.NORMAL.getCode().equals(mythJobTaskSchedule.getScheduleType())) {
privateParam = PlaceholderUtils.formatBizDateParam(privateParam, PlaceholderEnum.DATE_PLACEHOLDER.getName(), 1);
privateParam = PlaceholderUtils.formatBizDateParam(privateParam, PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName(), 0);
runParamWrapped.setPrivateParam(privateParam);
mythJobTaskSchedule.setRunParam(JSON.toJSONString(privateParam));
}
//获取到公有参数
Map<String,String> publicParam = runParamWrapped.getPublicParamMap();
//将私有参数转换为对应的参数DTO
......@@ -106,17 +114,14 @@ public class ScriptExecutorJobTask implements TimerTask {
if(scriptParamAndPlaceholderDto != null) {
Map<String, String> param = scriptParamAndPlaceholderDto.getParam();
if(CollectionUtil.isNotEmpty(param)){
command = PlaceholderUtils.commandReplace(command,param );
command = PlaceholderUtils.commandReplace(command,param);
}
}
//公共参数不为空的时候
if(CollectionUtil.isNotEmpty(publicParam)) {
command = PlaceholderUtils.commandReplace(command,publicParam );
}
//替换运行参数中的时间参数
if (ScheduleTypeEnum.NORMAL.getCode().equals(mythJobTaskSchedule.getScheduleType())) {
mythJobTaskSchedule.setRunParam(PlaceholderUtils.formatParam(privateParam));
}
}
return command;
}
......@@ -151,8 +156,6 @@ public class ScriptExecutorJobTask implements TimerTask {
if (mythJobTaskSchedule.getScheduleType().equals(4)) {
StringRedisTemplate stringRedisTemplate = (StringRedisTemplate) SpringUtil.getBean("stringRedisTemplate");
RunLog runLog = RunLog.builder().isEnd(true).isSuccess(false).runLog("执行资源异常" + MythLogUtils.getMessage(e)).build();
//stringRedisTemplate.convertAndSend("REAL:EXEC:" + mythJobTaskSchedule.getLogId() , JSON.toJSONString(runLog, WriteClassName));
//使用redis 想队尾push一个值
stringRedisTemplate.opsForList().rightPush("REAL:EXEC:" + mythJobTaskSchedule.getLogId() , JSON.toJSONString(runLog, WriteClassName));
}
......
......@@ -85,18 +85,14 @@ public class JobTaskThreadRunHelper extends BaseThreadRunHelper {
* 2.普通节点也有两种状态:
* I.开始节点:开始节点不需要验证上级工作流,直接放行执行
* II.正常节点:正常节点需要验证上级节点,首先判断自己是否收弱引用,如果是弱引用那么需要判断
* 上级节点是否已经全部都执行完了,执行完后不论成功与否都执行,同时工作流的运行结果
* 只与end节点关联
* 上级节点是否已经全部都执行完了,执行完后不论成功与否都执行该节点!
*
*/
if(FlowPropertyEnum.IS_INNER.getCode().equals(jobTask.getIsVirtual())) {
try {
if (nodeVerification.superiorNodeStatus(jobTask)) {
log.debug("检测到虚节点:{},切满足执行条件,开始运行", jobTask);
Integer scheduleType = jobTask.getScheduleType();
/*if (ScheduleTypeEnum.REPEAT.getCode().equals(scheduleType)){
log.debug("虚节点{},是重跑状态", jobTask)
runRecordingAndJobTaskService.updateRunRecordingAndSaveTask(jobTask)
}else */
if (ScheduleTypeEnum.REPAIR.getCode().equals(scheduleType) ||
ScheduleTypeEnum.REPEAT.getCode().equals(scheduleType)){
log.debug("虚节点{},是重跑或者补批状态", jobTask);
......
......@@ -62,14 +62,6 @@ public class ScheduleThreadRunHelper extends BaseThreadRunHelper {
log.debug("------排期表查询到有需要存在的节点--------");
//循环遍历添加任务
jobTaskSchedules.forEach(mythJobTaskSchedule ->{
//获取命令
String command = mythJobTaskSchedule.getRunCommand();
//除去JAVA任务
if(!NodeTypeEnum.JAVA.getCode().equals(mythJobTaskSchedule.getJobType())){
assert command != null;
command = command.replace("${"+ PlaceholderEnum.DATE_PLACEHOLDER.getName()+"}", DateUtil.dateLessDayStr(new Date(),"yyyyMMdd",1));
mythJobTaskSchedule.setRunCommand(command);
}
Long triggerTime = mythJobTaskSchedule.getTriggerTime();
TimerTask timerTask = null;
......
......@@ -11,6 +11,11 @@ import java.util.UUID;
*/
public class IDGenerationStrategy {
/**
* 工作流的runId生成策略
* @param port 端口号
* @return 生成的id
*/
public static String runIdGenerationStrategy(int port){
String runIdPre = UUID.randomUUID().toString().replace("-","");
String ipPort = IpUtil.getIpPort(port);
......@@ -22,6 +27,11 @@ public class IDGenerationStrategy {
return runIdPre;
}
/**
* 立即运行节点的runid生成策略
* @param logId 日志节点
* @return 生成的id
*/
public static String keyGenerationStrategy(Integer logId){
return "REAL:EXEC:" + logId;
}
......
......@@ -153,6 +153,9 @@ public class DateUtil {
* @return 计算后的时间字符串
*/
public static String dateLessDayStr(Date date,String dateFormat,Integer day){
if(StringUtils.isBlank(dateFormat)) {
dateFormat = "yyyyMMdd";
}
LocalDateTime localDateTime = dateToLocalDateTime(date);
LocalDateTime calculationDate = localDateTime.minusDays(day);
return calculationDate.format(DateTimeFormatter.ofPattern(dateFormat));
......
......@@ -89,35 +89,6 @@ public class PlaceholderUtils {
return sb.toString();
}
/**
* 脚本时间参数替换
* @param runParam 运行参数
* @param thisDate 补批时间
* @return 替换后的参数
*/
public static String paramPlaceholder(String runParam, String thisDate){
ScriptParamAndPlaceholderDto scriptParamAndPlaceholderDto = JSON.parseObject(runParam, ScriptParamAndPlaceholderDto.class);
if(scriptParamAndPlaceholderDto == null) {
return null;
}
Map<String, String> placeholder = scriptParamAndPlaceholderDto.getPlaceholder();
Map<String, String> param = scriptParamAndPlaceholderDto.getParam();
//替换时间参数
if (CollectionUtil.isNotEmpty(placeholder) && placeholder.containsKey(PlaceholderEnum.DATE_PLACEHOLDER.getName())) {
placeholder.put(PlaceholderEnum.DATE_PLACEHOLDER.getName(),thisDate);
}
//替换时间占位符
if(CollectionUtil.isNotEmpty(param) && param.containsKey(PlaceholderEnum.DATE_PLACEHOLDER.getName())){
param.put(PlaceholderEnum.DATE_PLACEHOLDER.getName(),thisDate);
}
return JSON.toJSONString(scriptParamAndPlaceholderDto);
}
/**
* 初始化脚本信息
* @param command 基础命令
......@@ -135,18 +106,47 @@ public class PlaceholderUtils {
* @param runParam 运行参数
* @return 参数字符串
*/
public static String formatParam(String runParam){
public static String formatBizDateParam(String runParam,String key, Integer cutBackDay){
return formatBizDateParam(runParam, key, null, cutBackDay);
}
public static String formatBizDateParam(String runParam,String key, String value, Integer cutBackDay){
if (StringUtils.isBlank(runParam)) {
return null;
}
ScriptParamAndPlaceholderDto scriptParamAndPlaceholderDto = JSON.parseObject(runParam,ScriptParamAndPlaceholderDto.class);
//获取命令参数
Map<String, String> param = scriptParamAndPlaceholderDto.getParam();
if(CollectionUtil.isNotEmpty(param) && param.containsKey(PlaceholderEnum.DATE_PLACEHOLDER.getName())){
String calculationDate = DateUtil.dateLessDayStr(new Date(), "yyyyMMdd", 1);
param.put(PlaceholderEnum.DATE_PLACEHOLDER.getName(),calculationDate);
Map<String, String> placeholder = scriptParamAndPlaceholderDto.getPlaceholder();
if(CollectionUtil.isNotEmpty(param) && param.containsKey(key)){
if(StringUtils.isNoneBlank(value)) {
placeholder.put(key,value);
}else{
String dateFormat = param.get(key);
String calculationDate = DateUtil.dateLessDayStr(new Date(), dateFormat, cutBackDay);
param.put(key,calculationDate);
scriptParamAndPlaceholderDto.setParam(param);
}
}
//获取脚本占位符
if(CollectionUtil.isNotEmpty(placeholder) && param.containsKey(key)){
if(StringUtils.isNoneBlank(value)) {
placeholder.put(key,value);
}else{
String dateFormat = param.get(key);
String calculationDate = DateUtil.dateLessDayStr(new Date(), dateFormat, cutBackDay);
placeholder.put(key,calculationDate);
scriptParamAndPlaceholderDto.setPlaceholder(placeholder);
}
}
return JSON.toJSONString(scriptParamAndPlaceholderDto);
}
}
......@@ -16,8 +16,9 @@ import java.util.Map;
@NoArgsConstructor
@AllArgsConstructor
public class ScriptParamAndPlaceholderDto implements Serializable {
private static final long serialVersionUID = -6901015072573970975L;
/**
* 参数的处理
* 参数的处理 命令参数
*/
private Map<String,String> param;
/**
......
package com.byit.dto.specials;
import lombok.Data;
import java.io.Serializable;
import java.util.List;
/**
* 补批参数
* @author huangfu
*/
@Data
public class RepairFlow implements Serializable {
private static final long serialVersionUID = -1984385012809451217L;
/**
* 工作空间名称
*/
private String workspaceName;
/**
* 工作流名称
*/
private String flowName;
/**
* 要补批的节点名称
*/
private List<String> nodeNames;
/**
* 补批的时间
*/
private List<String> repairTimes;
/**
* 时间类型格式
*/
private String timeFormatName;
/**
* nowDate时间格式
*/
private String nowDateTimeFormatName;
}
......@@ -18,6 +18,7 @@ import java.io.Serializable;
@AllArgsConstructor
@NoArgsConstructor
public class SpecialJavaNode implements Serializable {
private static final long serialVersionUID = 6183766311141765112L;
/**
* 工作流名称
*/
......
......@@ -51,5 +51,14 @@ public class SpecialJobParam {
* 时间参数
*/
private List<String> repairTimeList;
/**
* 时间类型格式
*/
private String timeFormatName;
/**
* nowDate时间格式
*/
private String nowDateTimeFormatName;
}
......@@ -7,7 +7,8 @@ package com.byit.enums;
public enum PlaceholderEnum {
DATE_PLACEHOLDER("biz_date"),
BIZ_SCRIPT_FILE("biz_file")
BIZ_SCRIPT_FILE("biz_file"),
NOW_DATE_PLACEHOLDER("now_date"),
;
private String name;
......
......@@ -5,6 +5,7 @@ import com.alibaba.fastjson.JSON;
import com.byit.dto.api.JavaCallbackLogDto;
import com.byit.dto.executor.PluginBeanJobInfo;
import com.byit.dto.plugin.*;
import com.byit.dto.specials.RepairFlow;
import com.byit.dto.specials.SpecialJobParam;
import com.byit.dto.web.ResponseResult;
import com.byit.executor.handler.interfaces.IJobHandler;
......@@ -619,20 +620,11 @@ public class JobUtils {
/**
* 补批工作流
*
* @param workspaceName 工作空间名称
* @param flowName 工作流名称
* @param nodeNames 节点名称 多个节点以","隔开
* @param repairTimes 补批日期 多个日期以","隔开
* @return
* @param repairFlow 补批参数
* @return 结果
*/
public static ResponseResult repairFlow(String workspaceName, String flowName, String nodeNames, String repairTimes) {
Map<String, Object> param = new HashMap<>();
param.put("workspaceName", workspaceName);
param.put("flowName", flowName);
param.put("nodeNames", nodeNames);
param.put("repairTimes", repairTimes);
String response = createHttpRequest(REQUEST_REPAIRFLOW, "param=" + JSON.toJSONString(param));
public static ResponseResult repairFlow(RepairFlow repairFlow) {
String response = createHttpRequest(REQUEST_REPAIRFLOW, JSON.toJSONString(repairFlow, WriteClassName));
log.debug("--------------------补批工作流接口调用成功,结果为:{}------------------------", response);
return JSON.parseObject(response, ResponseResult.class);
}
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment