Commit fc2bff87 by huangfusuper

开发等待工作流扫描任务

parent f8911063
package com.byit.thread.helper;
import com.byit.model.WaitingRecord;
import com.byit.service.WaitingRecordService;
import com.byit.service.mapservice.RunRecordingAndJobTaskService;
import com.byit.thread.BaseThreadRunHelper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import javax.sql.DataSource;
import java.util.*;
import java.util.concurrent.TimeUnit;
import java.util.function.BinaryOperator;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
* 补批工作流扫描
* 作用:补批 重跑之类的工作流等待进程
* @author huangfu
*/
@Component
@Slf4j
public class MakeUpFlowThreadRunHelper extends BaseThreadRunHelper {
private final int PRE_TIME = 30;
private static final String LOCK_NAME = "make_up_lock";
private final DataSource dataSource;
private final WaitingRecordService waitingRecordService;
private final RunRecordingAndJobTaskService runRecordingAndJobTaskService;
public MakeUpFlowThreadRunHelper(DataSource dataSource) {
public MakeUpFlowThreadRunHelper(DataSource dataSource, WaitingRecordService waitingRecordService, RunRecordingAndJobTaskService runRecordingAndJobTaskService) {
this.dataSource = dataSource;
this.waitingRecordService = waitingRecordService;
this.runRecordingAndJobTaskService = runRecordingAndJobTaskService;
}
@Override
public Long start() {
return null;
/**
* 30分钟的预读时间
*/
long preTestTime = System.currentTimeMillis()+ TimeUnit.HOURS.toMillis(PRE_TIME);
List<WaitingRecord> allByTriggerTime = waitingRecordService.findAllByTriggerTime(preTestTime);
/**
* 筛选数据,筛选出一组工作流上最小的任务流
*/
Map<Integer, WaitingRecord> nextRunFlow = allByTriggerTime.stream()
.collect(Collectors.toMap(WaitingRecord::getFlowId, Function.identity(), BinaryOperator.minBy(Comparator.comparingInt(WaitingRecord::getWaitOrder))));
nextRunFlow.forEach((key,value) ->{
runRecordingAndJobTaskService.updateRunRecordingAndTask(value);
});
return UNIVERSAL_WAIT_TIME;
}
@Override
......@@ -37,3 +68,4 @@ public class MakeUpFlowThreadRunHelper extends BaseThreadRunHelper {
return LOCK_NAME;
}
}
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