Commit 758cc4dd by huangfusuper

节点与节点之间的强依赖与弱依赖关系的配置及开发

parent 314ad159
...@@ -55,6 +55,7 @@ public class RunNodeServiceImpl implements RunNodeServer { ...@@ -55,6 +55,7 @@ public class RunNodeServiceImpl implements RunNodeServer {
List<JobTask> jobTasks = new ArrayList<>(32); List<JobTask> jobTasks = new ArrayList<>(32);
nodes.forEach(node -> { nodes.forEach(node -> {
JobTask jobTask = new JobTask(); JobTask jobTask = new JobTask();
jobTask.setTriggerStatus("1");
BeanUtils.copyProperties(node,jobTask); BeanUtils.copyProperties(node,jobTask);
if(isScheduleFollow){ if(isScheduleFollow){
jobTask.setTriggerTime(flow.getTriggerNextTime()); jobTask.setTriggerTime(flow.getTriggerNextTime());
......
...@@ -10,7 +10,6 @@ public interface TaskAndLogServer { ...@@ -10,7 +10,6 @@ public interface TaskAndLogServer {
/** /**
* 添加失败日志,删除任务表的任务 * 添加失败日志,删除任务表的任务
* @param jobTask * @param jobTask
* @param runCode
*/ */
void addRunLogAndRemoveTask(JobTask jobTask,String runCode); void addRunLogAndRemoveTask(JobTask jobTask);
} }
...@@ -27,7 +27,7 @@ public class TaskAndLogServerImpl implements TaskAndLogServer { ...@@ -27,7 +27,7 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
} }
@Override @Override
public void addRunLogAndRemoveTask(JobTask jobTask,String runCode) { public void addRunLogAndRemoveTask(JobTask jobTask) {
Date thisDate = new Date(); Date thisDate = new Date();
//删除任务节点 //删除任务节点
jobTaskService.removeMythJobTaskById(jobTask.getId()); jobTaskService.removeMythJobTaskById(jobTask.getId());
...@@ -44,7 +44,7 @@ public class TaskAndLogServerImpl implements TaskAndLogServer { ...@@ -44,7 +44,7 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
jobTaskRunLog.setHandlerName(jobTask.getHandlerName()); jobTaskRunLog.setHandlerName(jobTask.getHandlerName());
jobTaskRunLog.setIsVirtual(jobTask.getIsVirtual()); jobTaskRunLog.setIsVirtual(jobTask.getIsVirtual());
jobTaskRunLog.setMapFlowId(jobTask.getMapFlowId()); jobTaskRunLog.setMapFlowId(jobTask.getMapFlowId());
jobTaskRunLog.setRunCode(runCode); jobTaskRunLog.setRunCode("2");
jobTaskRunLog.setRunMsg("上级节点执行失败"); jobTaskRunLog.setRunMsg("上级节点执行失败");
jobTaskRunLog.setRunParams(jobTask.getRunParam()); jobTaskRunLog.setRunParams(jobTask.getRunParam());
jobTaskRunLog.setRunCommand(jobTask.getRunCommand()); jobTaskRunLog.setRunCommand(jobTask.getRunCommand());
......
...@@ -21,8 +21,10 @@ import java.sql.Connection; ...@@ -21,8 +21,10 @@ import java.sql.Connection;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.SQLException; import java.sql.SQLException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/** /**
* @program: byit-myth-job->JobScheduleHeloer * @program: byit-myth-job->JobScheduleHeloer
...@@ -126,18 +128,22 @@ public class JobScheduleHelper{ ...@@ -126,18 +128,22 @@ public class JobScheduleHelper{
preparedStatement.execute(); preparedStatement.execute();
//行锁已经加上 后续处理 //行锁已经加上 后续处理
long nowTime = System.currentTimeMillis(); long nowTime = System.currentTimeMillis();
//开始寻找此时 不是暂停状态,而且七秒内即将运行的任务 //开始寻找此时 不是暂停状态,而且七秒内即将运行的任务 而且还不是暂停的节点
List<JobTask> jobTasks = jobTaskService.findJobTaskByTriggerNextTimeLessThanEqual(nowTime + PRE_READ_MS); List<JobTask> jobTasks = jobTaskService.findJobTaskByTriggerNextTimeLessThanEqual(nowTime + PRE_READ_MS);
if(CollectionUtil.isNotEmpty(jobTasks)){ if(CollectionUtil.isNotEmpty(jobTasks)){
List<JobTaskSchedule> jobTaskSchedules = new ArrayList<JobTaskSchedule>(15); List<JobTaskSchedule> jobTaskSchedules = new ArrayList<JobTaskSchedule>(15);
jobTasks.forEach(jobTask -> { //遍历七秒内将要运行的节点数据
//判断是否是暂停状态 如果是暂停状态那么直接跳过这个节点 不执行 for(JobTask jobTask : jobTasks){
String triggerStatus = jobTask.getTriggerStatus(); /**
if(JobTriggerStatusEnums.STOP.getCode().equals(triggerStatus)){ * 判断节点状态
log.debug("--------------【{}节点处于暂停状态,跳过该节点】------------------"); * 1.虚节点状态,虚节点状态是映射了一个工作流,需要将该节点映射的工作流下所由的几点拉取到任务表
return; * 2.普通节点也有两种状态:
} * I.开始节点:开始节点不需要验证上级工作流,直接放行执行
* II.正常节点:正常节点需要验证上级节点,首先判断自己是否收弱引用,如果是弱引用那么需要判断
* 上级节点是否已经全部都执行完了,执行完后不论成功与否都执行,同时工作流的运行结果
* 只与end节点关联
*/
//虚节点的状态 //虚节点的状态
if("0".equals(jobTask.getIsVirtual())){ if("0".equals(jobTask.getIsVirtual())){
try { try {
...@@ -159,43 +165,37 @@ public class JobScheduleHelper{ ...@@ -159,43 +165,37 @@ public class JobScheduleHelper{
RunRecording runRecordingByFlowIdAndRunId = runRecordingService.findRunRecordingByFlowIdAndRunId(flowId, runId); RunRecording runRecordingByFlowIdAndRunId = runRecordingService.findRunRecordingByFlowIdAndRunId(flowId, runId);
runRecordingByFlowIdAndRunId.setFlowStatus(FlowPropertyEnum.FLOW_RUN_ING.getCode()); runRecordingByFlowIdAndRunId.setFlowStatus(FlowPropertyEnum.FLOW_RUN_ING.getCode());
runRecordingService.updateRunRecordingById(runRecordingByFlowIdAndRunId); runRecordingService.updateRunRecordingById(runRecordingByFlowIdAndRunId);
}else { }else{
//查询该节点的依赖节点 //查询该节点的依赖节点
List<Integer> dependIdByNodeId = nodeDependencyService.findDependIdByNodeId(jobTask.getNodeId()); List<Integer> dependIdByNodeId = nodeDependencyService.findDependIdByNodeId(jobTask.getNodeId());
//这里返回的是上级节点的日志执行情况 //这里返回的是上级节点的日志执行情况 把运行中的数据给过滤掉了
List<JobTaskRunLog> jobTaskRunLogList = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(dependIdByNodeId, jobTask.getRunId()); List<JobTaskRunLog> jobTaskRunLogList = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(dependIdByNodeId, jobTask.getRunId());
//如果上级节点为空的情况,就证明他的伤及节点一个都没执行成功,直接到下一个节点 if (CollectionUtil.isNotEmpty(jobTaskRunLogList)) {
if (CollectionUtil.isEmpty(jobTaskRunLogList)){ //判断父类节点是否已经全部完成,只需要判断依赖节点的数目和查询出来的日志数据是否相同
return; if(dependIdByNodeId.size() == jobTaskRunLogList.size()){
} //该节点如果为弱引用
//如果不为空 证明是有父类节点的,那么就要判断 if ("1".equals(jobTask.getSuperSuccessRun())) {
for (JobTaskRunLog jobTaskRunLog :jobTaskRunLogList){ //执行代码
//如果该节点不是成功的或者补批成功的 那么就跳过 runJobTask(jobTask,jobTaskSchedules);
if(!("1".equals(jobTaskRunLog.getRunCode()) || "3".equals(jobTaskRunLog.getRunCode()))){ }else{
if(!("0".equals(jobTaskRunLog.getRunCode()))){ //过滤失败的节点
//记录状态码 List<JobTaskRunLog> errorJobLog = jobTaskRunLogList.stream().
String runCode = jobTaskRunLog.getRunCode(); filter(jobTaskRunLog -> (!("1".equals(jobTaskRunLog.getRunCode()) || "3".equals(jobTaskRunLog.getRunCode()))))
log.info("-------------------【查询到有失败的节点,运行状态为:{}】-------------------",runCode); .collect(Collectors.toList());
/** //如果有失败的节点 就把改节点置为失败
* 判断父类节点是否全部都执行完毕了 if(CollectionUtil.isNotEmpty(errorJobLog)){
*/
if(parentIsEnd(jobTaskRunLogList)){
//删除这个数据 并且添加到日志 //删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask,runCode); taskAndLogServer.addRunLogAndRemoveTask(jobTask);
}else{
//执行代码
runJobTask(jobTask,jobTaskSchedules);
} }
} }
return;
} }
} }
//到这里 父类节点一定是全部都执行成功了!
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
} }
} }
}
});
if(CollectionUtil.isNotEmpty(jobTaskSchedules)) { if(CollectionUtil.isNotEmpty(jobTaskSchedules)) {
jobTaskScheduleService.saveAllData(jobTaskSchedules); jobTaskScheduleService.saveAllData(jobTaskSchedules);
...@@ -465,4 +465,11 @@ public class JobScheduleHelper{ ...@@ -465,4 +465,11 @@ public class JobScheduleHelper{
} }
return flag; return flag;
} }
private void runJobTask(JobTask jobTask,List<JobTaskSchedule> jobTaskSchedules){
//到这里 父类节点一定是全部都执行成功了!
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
}
} }
...@@ -56,7 +56,7 @@ ...@@ -56,7 +56,7 @@
, ,
<include refid="Blob_Column_List" /> <include refid="Blob_Column_List" />
from job_task from job_task
where trigger_time <![CDATA[ <= ]]> #{triggerNextTime,jdbcType=BIGINT} where trigger_time <![CDATA[ <= ]]> #{triggerNextTime,jdbcType=BIGINT} and trigger_status != '0'
</select> </select>
<!--根据id查询--> <!--根据id查询-->
<select id="findJobTaskById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs"> <select id="findJobTaskById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
......
...@@ -48,7 +48,8 @@ ...@@ -48,7 +48,8 @@
select <include refid="Base_Column_List" /> select <include refid="Base_Column_List" />
from job_task_run_log from job_task_run_log
where where
run_id = #{runId} run_code != '0'
and run_id = #{runId}
and node_id in and node_id in
<foreach item="nodeId" collection="nodeIds" open="(" separator="," close=")"> <foreach item="nodeId" collection="nodeIds" open="(" separator="," close=")">
#{nodeId} #{nodeId}
......
...@@ -19,7 +19,7 @@ public class DemoJob extends BaseJobHandler { ...@@ -19,7 +19,7 @@ public class DemoJob extends BaseJobHandler {
System.out.println("--------------任务就这样运行了 addJob----------------"+s ); System.out.println("--------------任务就这样运行了 addJob----------------"+s );
TimeUnit.SECONDS.sleep(10); TimeUnit.SECONDS.sleep(10);
if("1".equals(s)){ if("1".equals(s)){
throw new RuntimeException("我出异常了,信不信由你,我是爸爸"); throw new RuntimeException("我出异常了,信不信由你");
} }
return ReturnResult.SUCCESS; return ReturnResult.SUCCESS;
} }
......
...@@ -83,15 +83,17 @@ public class TestAddFlow { ...@@ -83,15 +83,17 @@ public class TestAddFlow {
pluginNode2.setConfig(pluginNodeConfig2); pluginNode2.setConfig(pluginNodeConfig2);
pluginNode2.setDependNodeNameList(Collections.singletonList("start")); pluginNode2.setDependNodeNameList(Collections.singletonList("start"));
//报错节点
PluginNode pluginNode3 = new PluginNode(); PluginNode pluginNode3 = new PluginNode();
PluginNodeConfig pluginNodeConfig3 = new PluginNodeConfig(); PluginNodeConfig pluginNodeConfig3 = new PluginNodeConfig();
pluginNode3.setName("中间节点2"); pluginNode3.setName("中间节点2");
pluginNode3.setDesc("中间节点2"); pluginNode3.setDesc("中间节点2");
pluginNode3.setType("node"); pluginNode3.setType("node");
pluginNode3.setAuthor("皇甫"); pluginNode3.setAuthor("皇甫");
pluginNode3.setSuperSuccessRun("0");
pluginNode3.setJobType("JAVA"); pluginNode3.setJobType("JAVA");
pluginNode3.setHandlerName("addJob"); pluginNode3.setHandlerName("addJob");
pluginNode3.setRunParam("addJob2"); pluginNode3.setRunParam("1");
pluginNodeConfig3.setFailedRetryCount(2); pluginNodeConfig3.setFailedRetryCount(2);
pluginNodeConfig3.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2)); pluginNodeConfig3.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig3.setNodeCron("0 0/7 * * * ? *"); pluginNodeConfig3.setNodeCron("0 0/7 * * * ? *");
...@@ -106,6 +108,8 @@ public class TestAddFlow { ...@@ -106,6 +108,8 @@ public class TestAddFlow {
PluginNodeConfig pluginNodeConfig4 = new PluginNodeConfig(); PluginNodeConfig pluginNodeConfig4 = new PluginNodeConfig();
pluginNode4.setName("中间节点3"); pluginNode4.setName("中间节点3");
pluginNode4.setDesc("中间节点3"); pluginNode4.setDesc("中间节点3");
//弱依赖
pluginNode4.setSuperSuccessRun("1");
pluginNode4.setType("node"); pluginNode4.setType("node");
pluginNode4.setAuthor("皇甫"); pluginNode4.setAuthor("皇甫");
pluginNode4.setJobType("JAVA"); pluginNode4.setJobType("JAVA");
......
...@@ -29,7 +29,7 @@ public class TestAddFlow2 { ...@@ -29,7 +29,7 @@ public class TestAddFlow2 {
.flowCron("0 0/1 * * * ? *") .flowCron("0 0/1 * * * ? *")
.flowTimeout(TimeUnit.MINUTES.toMillis(30)) .flowTimeout(TimeUnit.MINUTES.toMillis(30))
.priority("2") .priority("2")
.repeatCount(2) .repeatCount(1)
.scheduleFollow("1") .scheduleFollow("1")
.build(); .build();
......
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