Commit ae96fb05 by guominglei

修改虚节点的不跟随主工作流调度时间BUG

parent 09c58301
......@@ -89,7 +89,7 @@ public class RunNodeServiceImpl implements RunNodeServer, ApplicationEventPublis
jobTaskService.saveJobTasks(jobTasks);
log.info("-------------【开始修改工作流{}的下次运行时间,以及各种状态】---------------",flow);
try {
flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date(flow.getTriggerNextTime())).getTime());
flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
} catch (ParseException e) {
flow.setTriggerNextTime(999999999999999999L);
......
......@@ -88,7 +88,8 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
jobTaskRunLog.setIsVirtual(jobTask.getIsVirtual());
jobTaskRunLog.setMapFlowId(jobTask.getMapFlowId());
jobTaskRunLog.setRunCode(NodeRunStatusPropertyEnum.RUN_ING.getCode());
jobTaskRunLog.setRunType(jobTask.getJobType());
//目前都设置为执行机运行
jobTaskRunLog.setRunType("1");
Date thisDate = new Date();
jobTaskRunLog.setTriggerTime(thisDate);
jobTaskRunLog.setStartTime(thisDate);
......@@ -96,7 +97,6 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
jobTaskRunLogService.saveJobTaskRunLog(jobTaskRunLog);
log.info("-----------【虚节点保存运行记录】--------------");
//保存进运行记录表
Long triggerNextTime = virFlow.getTriggerNextTime();
RunRecording runRecording = RunRecording.builder()
.runId(jobTask.getRunId())
.flowId(mapFlowId)
......@@ -108,7 +108,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
.alarmEmail(virFlow.getAlarmEmail())
.alarmlAction(virFlow.getAlarmlAction())
.priority(virFlow.getPriority())
.triggerTime(equals?mainRecording.getTriggerTime(): triggerNextTime)
.triggerTime(equals?jobTask.getTriggerTime():virFlow.getTriggerNextTime())
.principal(virFlow.getPrincipal())
.startTime(new Date())
.flowNodeCount(virFlow.getFlowNodeCount())
......@@ -131,11 +131,18 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
}
BeanUtils.copyProperties(node, task);
task.setTriggerTime(mainRecording.getTriggerTime());
if (equals) {
task.setTriggerTime(mainRecording.getTriggerTime());
task.setTriggerTime(jobTask.getTriggerTime());
} else if (virFlag) {
task.setTriggerTime(triggerNextTime);
task.setTriggerTime(virFlow.getTriggerNextTime());
}else {
task.setTriggerTime(node.getTriggerNextTime());
try {
node.setTriggerNextTime(new CronExpression(node.getNodeCron()).getNextValidTimeAfter(new Date()).getTime());
} catch (ParseException e) {
e.printStackTrace();
}
//TODO 添加节点的下次运行时间修改
}
task.setRunId(jobTask.getRunId());
task.setTriggerStatus("1");
......@@ -147,9 +154,9 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
jobTaskService.removeMythJobTaskById(jobTask.getId());
log.info("-----------saveRunRecordingAndTask end【虚节点保存服务】--------------");
if(StringUtils.isNotBlank(virFlow.getFlowCron())){
virFlow.setTriggerNextTime(new CronExpression(virFlow.getFlowCron()).getNextValidTimeAfter(new Date(virFlow.getTriggerNextTime())).getTime());
}
virFlow.setTriggerNextTime(new CronExpression(virFlow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
flowService.updateByIdSelective(virFlow);
}
}
......@@ -161,15 +168,13 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
String reRunId = jobTask.getReRunId();
List<JobTaskRunLogWithBLOBs> reNodeLog = jobTaskRunLogService.findJobTaskRunLogWithBLOBsByFlowIdAndRunId(mapFlowId, reRunId);
for (JobTaskRunLogWithBLOBs jobTaskRunLogWithBLOBs : reNodeLog) {
//将旧的运行日志修改为运行中
jobTaskRunLogWithBLOBs.setRunCode(ExecuteStatusEnum.RUNING.getCode());
jobTaskRunLogWithBLOBs.setOperator(jobTask.getOperator());
jobTaskRunLogWithBLOBs.setScheduleType(jobTask.getScheduleType());
jobTaskRunLogWithBLOBs.setReRunId(jobTask.getReRunId());
jobTaskRunLogWithBLOBs.setReRunId(jobTask.getRunId());
//修改日志
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLogWithBLOBs);
//将日志修改为task节点
JobTask logConvertTask = logConvertTask(jobTaskRunLogWithBLOBs);
logConvertTask.setRunId(jobTask.getRunId());
logConvertTask.setReRunId(jobTask.getReRunId());
logConvertTask.setOperator(jobTask.getOperator());
logConvertTask.setScheduleType(jobTask.getScheduleType());
......@@ -179,11 +184,22 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
RunRecording runRecordingByFlowIdAndRunId = runRecordingService.findRunRecordingByFlowIdAndRunId(mapFlowId, reRunId);
runRecordingByFlowIdAndRunId.setFlowStatus("2");
runRecordingService.updateRunRecordingById(runRecordingByFlowIdAndRunId);
runRecordingByFlowIdAndRunId.setRunId(jobTask.getRunId());
runRecordingByFlowIdAndRunId.setReRunId(reRunId);
runRecordingByFlowIdAndRunId.setOperator(jobTask.getOperator());
runRecordingByFlowIdAndRunId.setScheduleType(jobTask.getScheduleType());
runRecordingByFlowIdAndRunId.setRecordingId(null);
RunRecording runRecording = new RunRecording();
BeanUtils.copyProperties(runRecordingByFlowIdAndRunId, runRecording);
runRecording.setFailFast(RunRecordingEnum.FAIL_FAST_NO.getCode());
runRecording.setRunId(jobTask.getRunId());
runRecording.setReRunId(reRunId);
//设置为未开始
runRecording.setFlowStatus(RunRecordingEnum.FLOW_STATUS_NOT_RUN.getCode());
runRecording.setFlowRunResult(null);
runRecording.setEndTime(null);
runRecording.setStartTime(new Date());
runRecording.setTriggerTime(jobTask.getTriggerTime());
runRecording.setOperator(jobTask.getOperator());
runRecording.setScheduleType(jobTask.getScheduleType());
runRecording.setRecordingId(null);
//新增一条重跑的运行记录
runRecordingService.saveRunRecording(runRecordingByFlowIdAndRunId);
}
......
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