Commit 185bea57 by guominglei

修改重跑的实现方式

parent 20959ae7
......@@ -482,25 +482,22 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//获取当前时间
String reRunId = UUID.randomUUID().toString().replace("-","");
Long triggerTime = System.currentTimeMillis();
JobTask jobTask = new JobTask();
BeanUtils.copyProperties(jobTaskRunLog, jobTask);
jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId);
jobTask.setReRunId(jobTaskRunLog.getRunId());
WaitingTask waitingTask = new WaitingTask();
BeanUtils.copyProperties(jobTaskRunLog, waitingTask);
waitingTask.setTriggerTime(triggerTime);
waitingTask.setReRunId(jobTaskRunLog.getRunId());
//设置为重跑
jobTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
jobTask.setOperator(userName);
jobTask.setLogId(null);
waitingTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
waitingTask.setOperator(userName);
//查询当前节点的依赖节点
List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(node.getNodeId());
if (dependNodeIdList != null && dependNodeIdList.size() > 0){
jobTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
waitingTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
}
List<JobTask> jobTaskList = new ArrayList<>();
jobTaskList.add(jobTask);
List<WaitingTask> waitingTaskList = new ArrayList<>();
waitingTaskList.add(waitingTask);
//校验通过,开始设置重跑
//判断重跑机制(单节点重跑,节点及下游重跑
......@@ -508,17 +505,26 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//查询依赖本节点的节点,并添加到集合中
List<Integer> subNodeIdList = nodeDependencyMapper.findSubNodeList(node.getNodeId());
if (subNodeIdList != null && subNodeIdList.size() > 0){
addDependNode(reRunId, runInfo.getRunId(), triggerTime, jobTaskList, subNodeIdList, userName);
addDependNode(runInfo.getRunId(), triggerTime, waitingTaskList, subNodeIdList, userName);
}
}
WaitingRecord waitingRecord = new WaitingRecord();
RunRecording newRunRecording = buildRunRecording(runRecording, triggerTime, reRunId, flow, userName);
runRecording.setFlowStatus(RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode());
runRecording.setFlowRunResult(null);
runRecordingMapper.updateState(runRecording);
BeanUtils.copyProperties(newRunRecording, waitingRecord);
Integer order = waitingRecordMapper.findOrderByFlowId(flow.getFlowId());
if (order == null){
order = 0;
}
waitingRecord.setWaitOrder(++order);
Integer waitId = waitingRecordMapper.insertSelective(waitingRecord);
//生成新的工作流实例
runRecordingMapper.saveRunRecording(newRunRecording);
jobTaskMapper.saveJobTasks(jobTaskList);
waitingTaskList.forEach(waiting -> {
waiting.setWaitId(waitingRecord.getWaitId());
waitingTaskMapper.insertSelective(waiting);
});
}
private RunRecording buildRunRecording(RunRecording runRecording, Long triggerTime, String reRunId, Flow flow, String userName) {
......@@ -567,72 +573,66 @@ public class ApiFlowServiceImpl implements ApiFlowService {
Long triggerTime = System.currentTimeMillis();
String reRunId = UUID.randomUUID().toString().replace("-","");
List<JobTask> jobTaskList = new ArrayList<>();
List<JobTaskRunLogWithBLOBs> jobTaskRunLogList = jobTaskRunLogMapper.findJobTaskRunLogWithBLOBsByFlowIdAndRunId(flow.getFlowId(), runInfo.getRunId());
ValidationUtil.dataNotNull(jobTaskRunLogList, "该工作流没有在调度上的任务");
String userName = currentUserUtils.account();
WaitingRecord waitingRecord = new WaitingRecord();
//生成新的实例
RunRecording newRunRecording = buildRunRecording(runRecording, triggerTime, reRunId, flow, userName);
BeanUtils.copyProperties(newRunRecording, waitingRecord);
Integer order = waitingRecordMapper.findOrderByFlowId(flow.getFlowId());
if (order == null){
order = 0;
}
waitingRecord.setWaitOrder(++order);
Integer waitId = waitingRecordMapper.insertSelective(waitingRecord);
//生成新的工作流实例
runRecordingMapper.saveRunRecording(newRunRecording);
jobTaskRunLogList.forEach(jobTaskRunLog -> {
JobTask jobTask = new JobTask();
BeanUtils.copyProperties(jobTaskRunLog, jobTask);
jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId);
jobTask.setReRunId(runInfo.getRunId());
WaitingTask waitingTask = new WaitingTask();
BeanUtils.copyProperties(jobTaskRunLog, waitingTask);
waitingTask.setTriggerTime(triggerTime);
waitingTask.setReRunId(jobTaskRunLog.getRunId());
//设置为重跑
jobTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
jobTask.setOperator(userName);
jobTask.setLogId(null);
waitingTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
waitingTask.setOperator(userName);
waitingTask.setReRunId(jobTaskRunLog.getRunId());
//查询当前节点的依赖节点
List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(jobTaskRunLog.getNodeId());
if (dependNodeIdList != null && dependNodeIdList.size() > 0){
jobTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
waitingTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
}
jobTaskList.add(jobTask);
waitingTask.setWaitId(waitingRecord.getWaitId());
waitingTaskMapper.insertSelective(waitingTask);
});
//生成运行实例
RunRecording newRunRecording = buildRunRecording(runRecording, triggerTime, reRunId, flow, userName);
runRecording.setFlowStatus(RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode());
runRecording.setFlowRunResult(null);
runRecordingMapper.updateState(runRecording);
//将工作流下的节点改为运行中
jobTaskRunLogMapper.updateByFlowIdAndRunId(flow.getFlowId(), runInfo.getRunId());
//生成新的工作流实例
runRecordingMapper.saveRunRecording(newRunRecording);
jobTaskMapper.saveJobTasks(jobTaskList);
}
}
private void addDependNode(String reRunId, String runId, Long triggerTime, List<JobTask> jobTaskList, List<Integer> subNodeIdList, String userName) {
private void addDependNode(String runId, Long triggerTime, List<WaitingTask> waitingTaskList, List<Integer> subNodeIdList, String userName) {
subNodeIdList.forEach(childNodeId -> {
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, childNodeId);
JobTask jobTask = new JobTask();
BeanUtils.copyProperties(jobTaskRunLog, jobTask);
jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId);
jobTask.setReRunId(runId);
WaitingTask waitingTask = new WaitingTask();
BeanUtils.copyProperties(jobTaskRunLog, waitingTask);
waitingTask.setTriggerTime(triggerTime);
waitingTask.setReRunId(runId);
//设置为重跑
jobTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
jobTask.setOperator(userName);
jobTask.setLogId(null);
waitingTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
waitingTask.setOperator(userName);
//查询当前节点的依赖节点
List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(jobTaskRunLog.getNodeId());
if (dependNodeIdList != null && dependNodeIdList.size() > 0){
jobTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
waitingTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
}
//将节点对应的日志改为运行中
jobTaskRunLog.setRunCode(NodeRunStatusPropertyEnum.RUN_ING.getCode());
jobTaskRunLogMapper.updateJobTaskRunLog(jobTaskRunLog);
jobTaskList.add(jobTask);
waitingTaskList.add(waitingTask);
//查询依赖于当前节点的下级节点
List<Integer> childNodeIdList = nodeDependencyMapper.findSubNodeList(childNodeId);
if (childNodeIdList != null && childNodeIdList.size() > 0){
addDependNode(reRunId, runId, triggerTime, jobTaskList, childNodeIdList, userName);
addDependNode(runId, triggerTime, waitingTaskList, childNodeIdList, userName);
}
});
}
......
......@@ -206,8 +206,6 @@ public interface RunRecordingMapper {
*/
RunRecording findNewStatus(Integer flowId);
int updateState(RunRecording runRecording);
Integer findFlowNum(@Param("startTime")Date startTime,
@Param("endTime")Date endTime,
@Param("flowIdList")List<Integer> flowIdList);
......
......@@ -606,9 +606,5 @@
where run_id = #{runId,jdbcType=VARCHAR}
and flow_name = #{flowName,jdbcType=VARCHAR}
</update>
<update id="updateState" parameterType="com.byit.model.RunRecording">
update run_recording set flow_status = #{flowStatus},flow_run_result = #{flowRunResult}
where recording_id = #{recordingId,jdbcType=INTEGER}
</update>
</mapper>
\ No newline at end of file
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