Commit 04a283e5 by huangfusuper

上级节点不满足的节点不执行立即失败,而是也会走声明周期函数

parent 80c055fc
...@@ -464,6 +464,7 @@ public class FlowServiceImpl implements FlowService { ...@@ -464,6 +464,7 @@ public class FlowServiceImpl implements FlowService {
runRecording.setFlowRunResult("5"); runRecording.setFlowRunResult("5");
runRecording.setEndTime(new Date()); runRecording.setEndTime(new Date());
runRecording.setIsAlarm(EmailEnum.IS_ALARM_NO.getCode()); runRecording.setIsAlarm(EmailEnum.IS_ALARM_NO.getCode());
runRecording.setFailFast("1");
runRecordingMapper.updateRunRecordingById(runRecording); runRecordingMapper.updateRunRecordingById(runRecording);
}); });
//杀死所有的任务 //杀死所有的任务
...@@ -494,13 +495,13 @@ public class FlowServiceImpl implements FlowService { ...@@ -494,13 +495,13 @@ public class FlowServiceImpl implements FlowService {
if (!killJob(jobTaskRunLog.getLogId(), jobTaskRunLog.getJobGroupIp(), jobTaskRunLog.getNodeName(), errorMsg)) { if (!killJob(jobTaskRunLog.getLogId(), jobTaskRunLog.getJobGroupIp(), jobTaskRunLog.getNodeName(), errorMsg)) {
result = false; result = false;
} }
jobTaskRunLog.setRunCode("2"); // jobTaskRunLog.setRunCode("2");
jobTaskRunLog.setRunMsg("节点被杀死!"); // jobTaskRunLog.setRunMsg("节点被杀死!");
jobTaskRunLog.setEndTime(new Date()); // jobTaskRunLog.setEndTime(new Date());
if (jobTaskRunLog.getStartTime() == null) { // if (jobTaskRunLog.getStartTime() == null) {
jobTaskRunLog.setEndTime(new Date()); // jobTaskRunLog.setEndTime(new Date());
} // }
jobTaskRunLogMapper.updateJobTaskRunLog(jobTaskRunLog); // jobTaskRunLogMapper.updateJobTaskRunLog(jobTaskRunLog);
} }
} }
} }
...@@ -516,8 +517,7 @@ public class FlowServiceImpl implements FlowService { ...@@ -516,8 +517,7 @@ public class FlowServiceImpl implements FlowService {
}); });
flow.setScanMark("1"); flow.setScanMark("1");
flowMapper.updateByIdSelective(flow); flowMapper.updateByIdSelective(flow);
result = true; return new KillDto(result, errorMsg.toString());
return new KillDto(result, "停止成功");
} }
private synchronized Boolean killJob(Integer logId, String exectUrl, String nodeName, StringBuffer errorMsg) { private synchronized Boolean killJob(Integer logId, String exectUrl, String nodeName, StringBuffer errorMsg) {
......
package com.byit.service.mapservice.impl; package com.byit.service.mapservice.impl;
import cn.hutool.core.bean.BeanUtil;
import cn.hutool.core.collection.CollectionUtil; import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.byit.dto.BeanStrategyPackage;
import com.byit.dto.executor.RunParamWrapped; import com.byit.dto.executor.RunParamWrapped;
import com.byit.enums.EmailEnum; import com.byit.enums.EmailEnum;
import com.byit.enums.FlowPropertyEnum; import com.byit.enums.FlowPropertyEnum;
...@@ -13,6 +15,9 @@ import com.byit.job.utils.PlaceholderUtils; ...@@ -13,6 +15,9 @@ import com.byit.job.utils.PlaceholderUtils;
import com.byit.model.*; import com.byit.model.*;
import com.byit.service.*; import com.byit.service.*;
import com.byit.service.mapservice.TaskAndLogServer; import com.byit.service.mapservice.TaskAndLogServer;
import com.byit.strategy.TaskRunTheLifeCycleCallback;
import com.byit.util.ClassSortUtil;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.BeanUtils; import org.springframework.beans.BeanUtils;
...@@ -24,6 +29,7 @@ import java.net.UnknownHostException; ...@@ -24,6 +29,7 @@ import java.net.UnknownHostException;
import java.text.ParseException; import java.text.ParseException;
import java.util.Date; import java.util.Date;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.stream.Collectors; import java.util.stream.Collectors;
/** /**
...@@ -64,28 +70,41 @@ public class TaskAndLogServerImpl implements TaskAndLogServer { ...@@ -64,28 +70,41 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
cutBackDay = 0; cutBackDay = 0;
} }
//执行快速失败时 将参数替换强行替换进去 替换的是私有参数 //执行快速失败时 将参数替换强行替换进去 替换的是私有参数
String runParam = jobTask.getRunParam(); // String runParam = jobTask.getRunParam();
if (StringUtils.isNotBlank(runParam)) { // if (StringUtils.isNotBlank(runParam)) {
RunParamWrapped runParamWrapped = JSON.parseObject(runParam, RunParamWrapped.class); // RunParamWrapped runParamWrapped = JSON.parseObject(runParam, RunParamWrapped.class);
String privateParam = runParamWrapped.getPrivateParam(); // String privateParam = runParamWrapped.getPrivateParam();
String newParam = PlaceholderUtils.formatBizDateParam(privateParam, PlaceholderEnum.DATE_PLACEHOLDER.getName(), cutBackDay); // String newParam = PlaceholderUtils.formatBizDateParam(privateParam, PlaceholderEnum.DATE_PLACEHOLDER.getName(), cutBackDay);
newParam = PlaceholderUtils.formatBizDateParam(newParam, PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName(), 0); // newParam = PlaceholderUtils.formatBizDateParam(newParam, PlaceholderEnum.NOW_DATE_PLACEHOLDER.getName(), 0);
runParamWrapped.setPrivateParam(newParam); // runParamWrapped.setPrivateParam(newParam);
jobTaskRunLog.setRunParams(JSON.toJSONString(runParamWrapped)); // jobTaskRunLog.setRunParams(JSON.toJSONString(runParamWrapped));
} // }
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask, jobTaskSchedule);
//获取该节点的全部声明周期函数
Map<String, TaskRunTheLifeCycleCallback> stringTaskRunTheLifeCycleCallbackMap = SpringUtil.getBeansOfType(TaskRunTheLifeCycleCallback.class);
//数据排序
List<BeanStrategyPackage<TaskRunTheLifeCycleCallback>> beanStrategyPackages = ClassSortUtil.objectSort(stringTaskRunTheLifeCycleCallbackMap);
beanStrategyPackages.forEach(beanStrategyPackage -> {
log.info("--------------脚本节点生命周期开始回调{}--------------", beanStrategyPackage.getBeanName());
TaskRunTheLifeCycleCallback beanStrategyPackageBean = beanStrategyPackage.getBean();
if (beanStrategyPackageBean.matchType(jobTaskSchedule)) {
beanStrategyPackageBean.postProcessAfterInitialization(jobTaskSchedule);
}
});
BeanUtils.copyProperties(jobTask, jobTaskRunLog); BeanUtils.copyProperties(jobTaskSchedule, jobTaskRunLog);
jobTaskRunLog.setRunId(jobTask.getRunId()); jobTaskRunLog.setRunId(jobTaskSchedule.getRunId());
jobTaskRunLog.setNodeId(jobTask.getNodeId()); jobTaskRunLog.setNodeId(jobTaskSchedule.getNodeId());
jobTaskRunLog.setNodeName(jobTask.getNodeName()); jobTaskRunLog.setNodeName(jobTaskSchedule.getNodeName());
jobTaskRunLog.setJobType(jobTask.getJobType()); jobTaskRunLog.setJobType(jobTaskSchedule.getJobType());
jobTaskRunLog.setFlowId(jobTask.getFlowId()); jobTaskRunLog.setFlowId(jobTaskSchedule.getFlowId());
jobTaskRunLog.setFailedRemainingCount(jobTask.getFailedRetryCount()); jobTaskRunLog.setFailedRemainingCount(jobTaskSchedule.getFailedRetryCount());
jobTaskRunLog.setVersionName(jobTask.getVersionName()); jobTaskRunLog.setVersionName(jobTaskSchedule.getVersionName());
jobTaskRunLog.setFlowName(jobTask.getFlowName()); jobTaskRunLog.setFlowName(jobTaskSchedule.getFlowName());
jobTaskRunLog.setHandlerName(jobTask.getHandlerName()); jobTaskRunLog.setHandlerName(jobTaskSchedule.getHandlerName());
jobTaskRunLog.setIsVirtual(jobTask.getIsVirtual()); jobTaskRunLog.setIsVirtual(jobTaskSchedule.getIsVirtual());
jobTaskRunLog.setMapFlowId(jobTask.getMapFlowId()); jobTaskRunLog.setMapFlowId(jobTaskSchedule.getMapFlowId());
if (taskError) { if (taskError) {
jobTaskRunLog.setRunCode("2"); jobTaskRunLog.setRunCode("2");
jobTaskRunLog.setRunMsg(StringUtils.join(errorMsg,",")); jobTaskRunLog.setRunMsg(StringUtils.join(errorMsg,","));
...@@ -94,21 +113,14 @@ public class TaskAndLogServerImpl implements TaskAndLogServer { ...@@ -94,21 +113,14 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
jobTaskRunLog.setRunMsg("上级节点执行失败"); jobTaskRunLog.setRunMsg("上级节点执行失败");
} }
jobTaskRunLog.setRunParams(jobTask.getRunParam()); jobTaskRunLog.setRunParams(jobTaskSchedule.getRunParam());
jobTaskRunLog.setRunCommand(jobTask.getRunCommand()); jobTaskRunLog.setRunCommand(jobTaskSchedule.getRunCommand());
jobTaskRunLog.setRunType("2"); jobTaskRunLog.setRunType("2");
jobTaskRunLog.setTriggerCode("2"); jobTaskRunLog.setTriggerCode("2");
jobTaskRunLog.setTriggerMsg("未执行调度"); jobTaskRunLog.setTriggerMsg("未执行调度");
jobTaskRunLog.setTriggerTime(thisDate); jobTaskRunLog.setTriggerTime(thisDate);
jobTaskRunLog.setStartTime(thisDate); jobTaskRunLog.setStartTime(thisDate);
jobTaskRunLog.setEndTime(thisDate); jobTaskRunLog.setEndTime(thisDate);
// if("end".equals(jobTask.getNodeName())){
// //未完成告警
// jobTaskRunLog.setAlertEnd("0");
// }else{
// //已完成告警
// jobTaskRunLog.setAlertEnd("1");
// }
jobTaskRunLog.setAlertEnd("0"); jobTaskRunLog.setAlertEnd("0");
//添加日志节点 //添加日志节点
......
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