Commit a396df40 by guominglei

Merge remote-tracking branch 'origin/developer' into developer

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