Commit ad0ffe29 by huangfusuper

修复补批工作流执行重复的BUG

parent e5e1e24a
...@@ -42,7 +42,7 @@ public class RunRecordingServiceImpl implements RunRecordingService { ...@@ -42,7 +42,7 @@ public class RunRecordingServiceImpl implements RunRecordingService {
/** /**
* 是否存在运行的工作流实例 * 是否存在运行的工作流实例
* @param flowId 工作流ID * @param flowId 工作流ID
* @return 是否在正在运行中 * @return true有正在运行中的数据 反之没有
*/ */
@Override @Override
public boolean findRunRecordingIsRunning(Integer flowId) { public boolean findRunRecordingIsRunning(Integer flowId) {
......
package com.byit.thread.helper; package com.byit.thread.helper;
import com.byit.model.RunRecording;
import com.byit.model.WaitingRecord; import com.byit.model.WaitingRecord;
import com.byit.service.RunRecordingService;
import com.byit.service.WaitingRecordService; import com.byit.service.WaitingRecordService;
import com.byit.service.mapservice.RunRecordingAndJobTaskService; import com.byit.service.mapservice.RunRecordingAndJobTaskService;
import com.byit.thread.BaseThreadRunHelper; import com.byit.thread.BaseThreadRunHelper;
...@@ -28,14 +30,16 @@ public class MakeUpFlowThreadRunHelper extends BaseThreadRunHelper { ...@@ -28,14 +30,16 @@ public class MakeUpFlowThreadRunHelper extends BaseThreadRunHelper {
private final WaitingRecordService waitingRecordService; private final WaitingRecordService waitingRecordService;
private final RunRecordingAndJobTaskService runRecordingAndJobTaskService; private final RunRecordingAndJobTaskService runRecordingAndJobTaskService;
private final RunRecordingService runRecordingService;
public MakeUpFlowThreadRunHelper(DataSource dataSource, WaitingRecordService waitingRecordService, RunRecordingAndJobTaskService runRecordingAndJobTaskService) { public MakeUpFlowThreadRunHelper(DataSource dataSource, WaitingRecordService waitingRecordService, RunRecordingAndJobTaskService runRecordingAndJobTaskService, RunRecordingService runRecordingService) {
this.dataSource = dataSource; this.dataSource = dataSource;
this.waitingRecordService = waitingRecordService; this.waitingRecordService = waitingRecordService;
this.runRecordingAndJobTaskService = runRecordingAndJobTaskService; this.runRecordingAndJobTaskService = runRecordingAndJobTaskService;
this.runRecordingService = runRecordingService;
} }
@Override @Override
...@@ -51,8 +55,12 @@ public class MakeUpFlowThreadRunHelper extends BaseThreadRunHelper { ...@@ -51,8 +55,12 @@ public class MakeUpFlowThreadRunHelper extends BaseThreadRunHelper {
Map<Integer, WaitingRecord> nextRunFlow = allByTriggerTime.stream() Map<Integer, WaitingRecord> nextRunFlow = allByTriggerTime.stream()
.collect(Collectors.toMap(WaitingRecord::getFlowId, Function.identity(), BinaryOperator.minBy(Comparator.comparingInt(WaitingRecord::getWaitOrder)))); .collect(Collectors.toMap(WaitingRecord::getFlowId, Function.identity(), BinaryOperator.minBy(Comparator.comparingInt(WaitingRecord::getWaitOrder))));
nextRunFlow.forEach((key,value) ->{ nextRunFlow.forEach((key,value) ->{
boolean runRecordingIsRunning = runRecordingService.findRunRecordingIsRunning(value.getFlowId());
if(!runRecordingIsRunning){
runRecordingAndJobTaskService.updateRunRecordingAndTask(value); runRecordingAndJobTaskService.updateRunRecordingAndTask(value);
}
}); });
return UNIVERSAL_WAIT_TIME; return UNIVERSAL_WAIT_TIME;
......
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