Commit 06a71f39 by huangfusuper

任务表扫描线程重构

parent f741bebc
package com.byit.factory.material;
import com.byit.thread.BaseThreadRunHelper;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.ToString;
import javax.sql.DataSource;
/**
* @author huangfu
*/
@AllArgsConstructor
@Data
@ToString
public class ThreadRunMaterial {
/**
* 此线程
*/
private Thread thread;
/**
* 业务规则
*/
private BaseThreadRunHelper threadRunHelper;
/**
* 循环一轮睡眠时长
*/
private Long sleepTime;
}
package com.byit.thread.helper;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodeNameEnum;
import com.byit.enums.NodePropertyEnum;
import com.byit.enums.NodeRunStatusPropertyEnum;
import com.byit.model.JobTask;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskSchedule;
import com.byit.model.RunRecording;
import com.byit.service.JobTaskRunLogService;
import com.byit.service.JobTaskService;
import com.byit.service.NodeDependencyService;
import com.byit.service.RunRecordingService;
import com.byit.service.mapservice.RunRecordingAndJobTaskService;
import com.byit.service.mapservice.TaskAndLogServer;
import com.byit.service.mapservice.TaskAndScheduleService;
import com.byit.thread.BaseThreadRunHelper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Component;
import javax.sql.DataSource;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Collectors;
/**
* 任务表操作 向排期表添加任务节点并执行
* @author huangfu
*/
@Component
@Slf4j
public class TaskThreadRunHelper extends BaseThreadRunHelper {
/**
* 读取任务节点的预读
*/
private static final long PRE_READ_MS = 7000;
private static final String LOCK_NAME = "job_task_lock";
private final DataSource dataSource;
private final JobTaskService jobTaskService;
/**
* 运行实例和task的组合操作
*/
private final RunRecordingAndJobTaskService runRecordingAndJobTaskService;
/**
* 运行记录表信息操作
*/
private final RunRecordingService runRecordingService;
/**
* 节点依赖查询操作
*/
private final NodeDependencyService nodeDependencyService;
/**
* 日志节点操作
*/
private final JobTaskRunLogService jobTaskRunLogService;
/**
* 任务表和日志表操作
*/
private final TaskAndLogServer taskAndLogServer;
/**
* 任务表和排期表的组合操作
*/
private final TaskAndScheduleService taskAndScheduleService;
public TaskThreadRunHelper(DataSource dataSource, JobTaskService jobTaskService,
RunRecordingAndJobTaskService runRecordingAndJobTaskService,
RunRecordingService runRecordingService, NodeDependencyService nodeDependencyService,
JobTaskRunLogService jobTaskRunLogService, TaskAndLogServer taskAndLogServer,
TaskAndScheduleService taskAndScheduleService) {
this.dataSource = dataSource;
this.jobTaskService = jobTaskService;
this.runRecordingAndJobTaskService = runRecordingAndJobTaskService;
this.runRecordingService = runRecordingService;
this.nodeDependencyService = nodeDependencyService;
this.jobTaskRunLogService = jobTaskRunLogService;
this.taskAndLogServer = taskAndLogServer;
this.taskAndScheduleService = taskAndScheduleService;
}
@Override
public Long start() {
long nowTime = System.currentTimeMillis();
//开始寻找此时 不是暂停状态,而且七秒内即将运行的任务 而且还不是暂停的节点
List<JobTask> jobTasks = jobTaskService.findJobTaskByTriggerNextTimeLessThanEqual(nowTime + PRE_READ_MS);
if(CollectionUtil.isNotEmpty(jobTasks)) {
List<JobTaskSchedule> jobTaskSchedules = new ArrayList<JobTaskSchedule>(15);
//遍历七秒内将要运行的节点数据
for(JobTask jobTask : jobTasks){
log.debug("任务:{},开始运行", jobTask);
/**
* 判断节点状态
* 1.虚节点状态,虚节点状态是映射了一个工作流,需要将该节点映射的工作流下所由的几点拉取到任务表
* 2.普通节点也有两种状态:
* I.开始节点:开始节点不需要验证上级工作流,直接放行执行
* II.正常节点:正常节点需要验证上级节点,首先判断自己是否收弱引用,如果是弱引用那么需要判断
* 上级节点是否已经全部都执行完了,执行完后不论成功与否都执行,同时工作流的运行结果
* 只与end节点关联
*/
if(FlowPropertyEnum.IS_INNER.getCode().equals(jobTask.getIsVirtual())) {
//虚节点处理操作
innerNodeOperating(jobTask);
}else{
//处理开始节点
if(NodeNameEnum.START_NODE.getNodeName().equals(jobTask.getNodeName())){
startNodeOperating(jobTask,jobTaskSchedules);
}else{
//处理普通节点
nodeOperating(jobTask,jobTaskSchedules);
}
}
}
//执行保存到排表 删除任务表操作
taskAndScheduleService.saveScheduleAndDeleteTask(jobTaskSchedules);
}else{
return PRE_READ_MS;
}
return NOT_WAIT_TIME;
}
/**
* 普通节点操作
* @param thisJobTask 当前的任务节点
* @param jobTaskSchedules 排期集合
*/
private void nodeOperating(JobTask thisJobTask,List<JobTaskSchedule> jobTaskSchedules){
//查询该节点的依赖节点
List<Integer> dependIdByNodeId = nodeDependencyService.findDependIdByNodeId(thisJobTask.getNodeId());
//这里返回的是上级节点的日志执行情况 把运行中的数据给过滤掉了
List<JobTaskRunLog> jobTaskRunLogList = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(dependIdByNodeId, thisJobTask.getRunId());
if (CollectionUtil.isNotEmpty(jobTaskRunLogList)) {
//判断父类节点是否已经全部完成,只需要判断依赖节点的数目和查询出来的日志数据是否相同
if(dependIdByNodeId.size() == jobTaskRunLogList.size()){
//过滤失败的节点
List<JobTaskRunLog> errorJobLog = jobTaskRunLogList.stream().
filter(jobTaskRunLog -> (NodeRunStatusPropertyEnum.RUN_FAILURE.getCode().equals(jobTaskRunLog.getRunCode())
|| NodeRunStatusPropertyEnum.RE_RUN_FAILURE.getCode().equals(jobTaskRunLog.getRunCode())
|| NodeRunStatusPropertyEnum.PARENT_NODE_FAILED.getCode().equals(jobTaskRunLog.getRunCode())))
.collect(Collectors.toList());
//判断剩余执行次数是否为0
if (parentNodeErrorCount(errorJobLog)) {
log.debug("--------------【{}的上级节点的失败节点已经全部重试完毕】------------------",thisJobTask);
//该节点如果为弱引用
if (NodePropertyEnum.WEAK_NODE.getCode().equals(thisJobTask.getSuperSuccessRun())) {
log.info("-------------【查询到有弱引用节点】-----------------");
//执行代码
runJobTask(thisJobTask,jobTaskSchedules);
}else{
//如果有失败的节点 就把该节点置为失败
if(CollectionUtil.isNotEmpty(errorJobLog)){
//删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(thisJobTask);
}else{
//执行代码
runJobTask(thisJobTask,jobTaskSchedules);
}
}
}
}
}
}
/**
* start节点的操作
* @param thisJobTask 当前的任务节点
* @param jobTaskSchedules 排期集合
*/
private void startNodeOperating(JobTask thisJobTask,List<JobTaskSchedule> jobTaskSchedules){
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(thisJobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
//更改运行记录为运行中
String runId = thisJobTask.getRunId();
Integer flowId = thisJobTask.getFlowId();
RunRecording runRecordingByFlowIdAndRunId = runRecordingService.findRunRecordingByFlowIdAndRunId(flowId, runId);
runRecordingByFlowIdAndRunId.setFlowStatus(FlowPropertyEnum.FLOW_RUN_ING.getCode());
runRecordingService.updateRunRecordingById(runRecordingByFlowIdAndRunId);
}
/**
* 内嵌节点处理操作
* @param thisJobTask 当前的任务节点
*/
private void innerNodeOperating(JobTask thisJobTask){
try {
runRecordingAndJobTaskService.saveRunRecordingAndTask(thisJobTask);
log.info("----------------【虚节点保存成功,删除虚节点】--------------------");
jobTaskService.removeMythJobTaskById(thisJobTask.getId());
} catch (Exception e) {
log.error("--------------------虚节点处理出现异常{}------------------",e.getMessage());
}
}
@Override
public DataSource getDataSource() {
return dataSource;
}
@Override
public String getLockName() {
return LOCK_NAME;
}
/**
* 判断失败节点的重试次数是不是为0
* @param errorJobLog 上级节点的全部失败节点
* @return
*/
private boolean parentNodeErrorCount(List<JobTaskRunLog> errorJobLog){
if(CollectionUtil.isEmpty(errorJobLog)){
return true;
}
for (JobTaskRunLog jobTaskRunLog : errorJobLog) {
//失败重试次数大于0 而且错误原因不是上级节点执行失败
if(jobTaskRunLog.getFailedRemainingCount()>0 && !("6".equals(jobTaskRunLog.getRunCode()))){
log.debug("----------------【{}节点没有重试完毕】----------------",jobTaskRunLog);
return false;
}
}
return true;
}
/**
* 运行符合条件的人物节点 将task节点转换成排期节点 保存到集合
* @param jobTask 任务节点
* @param jobTaskSchedules 排期集合
*/
private void runJobTask(JobTask jobTask,List<JobTaskSchedule> jobTaskSchedules){
//到这里 父类节点一定是全部都执行成功了,或者是弱节点!
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
}
}
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