Commit 9b0199d1 by huangfusuper

job_task扫描线程对于失败重试次数的控制

parent 54abcb67
...@@ -76,7 +76,7 @@ public class JobTaskRunLog implements Serializable { ...@@ -76,7 +76,7 @@ public class JobTaskRunLog implements Serializable {
/** /**
* 运行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill * 运行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill
*/ */
@ApiModelProperty("运行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill") @ApiModelProperty("运行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill 6.上级节点执行失败")
private String runCode; private String runCode;
/** /**
......
...@@ -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("2"); jobTaskRunLog.setRunCode("6");
jobTaskRunLog.setRunMsg("上级节点执行失败"); jobTaskRunLog.setRunMsg("上级节点执行失败");
jobTaskRunLog.setRunParams(jobTask.getRunParam()); jobTaskRunLog.setRunParams(jobTask.getRunParam());
jobTaskRunLog.setRunCommand(jobTask.getRunCommand()); jobTaskRunLog.setRunCommand(jobTask.getRunCommand());
......
...@@ -173,22 +173,26 @@ public class JobScheduleHelper{ ...@@ -173,22 +173,26 @@ public class JobScheduleHelper{
if (CollectionUtil.isNotEmpty(jobTaskRunLogList)) { if (CollectionUtil.isNotEmpty(jobTaskRunLogList)) {
//判断父类节点是否已经全部完成,只需要判断依赖节点的数目和查询出来的日志数据是否相同 //判断父类节点是否已经全部完成,只需要判断依赖节点的数目和查询出来的日志数据是否相同
if(dependIdByNodeId.size() == jobTaskRunLogList.size()){ if(dependIdByNodeId.size() == jobTaskRunLogList.size()){
//该节点如果为弱引用 //过滤失败的节点
if ("1".equals(jobTask.getSuperSuccessRun())) { List<JobTaskRunLog> errorJobLog = jobTaskRunLogList.stream().
//执行代码 filter(jobTaskRunLog -> ("2".equals(jobTaskRunLog.getRunCode()) || "4".equals(jobTaskRunLog.getRunCode())))
runJobTask(jobTask,jobTaskSchedules); .collect(Collectors.toList());
}else{ //判断剩余执行次数是否为0
//过滤失败的节点 if (parentNodeErrorCount(errorJobLog)) {
List<JobTaskRunLog> errorJobLog = jobTaskRunLogList.stream(). log.debug("--------------【{}的上级节点的失败节点已经全部重试完毕】------------------",jobTask);
filter(jobTaskRunLog -> (!("1".equals(jobTaskRunLog.getRunCode()) || "3".equals(jobTaskRunLog.getRunCode())))) //该节点如果为弱引用
.collect(Collectors.toList()); if ("1".equals(jobTask.getSuperSuccessRun())) {
//如果有失败的节点 就把改节点置为失败
if(CollectionUtil.isNotEmpty(errorJobLog)){
//删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask);
}else{
//执行代码 //执行代码
runJobTask(jobTask,jobTaskSchedules); runJobTask(jobTask,jobTaskSchedules);
}else{
//如果有失败的节点 就把该节点置为失败
if(CollectionUtil.isNotEmpty(errorJobLog)){
//删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask);
}else{
//执行代码
runJobTask(jobTask,jobTaskSchedules);
}
} }
} }
} }
...@@ -455,15 +459,23 @@ public class JobScheduleHelper{ ...@@ -455,15 +459,23 @@ public class JobScheduleHelper{
return jobTaskRunLog.getLogId(); return jobTaskRunLog.getLogId();
} }
private boolean parentIsEnd(List<JobTaskRunLog> jobTaskRunLogList){ /**
boolean flag = true; * 判断失败节点的重试次数是不是为0
for (JobTaskRunLog jobTaskRunLog : jobTaskRunLogList) { * @param errorJobLog 上级节点的全部失败节点
if ("0".equals(jobTaskRunLog.getRunCode())) { * @return
flag = false; */
break; 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 flag; return true;
} }
private void runJobTask(JobTask jobTask,List<JobTaskSchedule> jobTaskSchedules){ private void runJobTask(JobTask jobTask,List<JobTaskSchedule> jobTaskSchedules){
......
...@@ -109,7 +109,7 @@ public class TestAddFlow { ...@@ -109,7 +109,7 @@ public class TestAddFlow {
pluginNode4.setName("中间节点3"); pluginNode4.setName("中间节点3");
pluginNode4.setDesc("中间节点3"); pluginNode4.setDesc("中间节点3");
//弱依赖 //弱依赖
pluginNode4.setSuperSuccessRun("1"); pluginNode4.setSuperSuccessRun("0");
pluginNode4.setType("node"); pluginNode4.setType("node");
pluginNode4.setAuthor("皇甫"); pluginNode4.setAuthor("皇甫");
pluginNode4.setJobType("JAVA"); pluginNode4.setJobType("JAVA");
......
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