Commit cf5acc9c by huangfusuper

【删除失败节点】调整删除失败流程

【修复BUG】修复由单例模式引发的数据错乱
parent 86aaa505
......@@ -4,6 +4,7 @@ import com.byit.mapper.JobTaskRunLogMapper;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.service.JobTaskRunLogService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
......@@ -17,6 +18,7 @@ import java.util.List;
* @author: huangfu
* @date: 2019/12/20 19:43
**/
@Slf4j
@Service
@Transactional(rollbackFor = Exception.class)
public class JobTaskRunLogServiceImpl implements JobTaskRunLogService {
......@@ -56,6 +58,7 @@ public class JobTaskRunLogServiceImpl implements JobTaskRunLogService {
@Override
public int saveJobTaskRunLog(JobTaskRunLogWithBLOBs jobTaskRunLog) {
log.info("--------------saveJobTaskRunLog start【{}】------------------------",jobTaskRunLog);
return jobTaskRunLogMapper.saveJobTaskRunLog(jobTaskRunLog);
}
......
......@@ -55,6 +55,7 @@ public class RunNodeServiceImpl implements RunNodeServer {
jobTask.setTriggerTime(flow.getTriggerNextTime());
}
jobTask.setRunId(runId);
jobTask.setFlowName(flow.getFlowName());
jobTasks.add(jobTask);
});
jobTaskService.saveJobTasks(jobTasks);
......
......@@ -33,7 +33,7 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
jobTaskService.removeMythJobTaskById(jobTask.getId());
//添加任务日志
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
jobTaskRunLog.setReRunId(jobTask.getRunId());
jobTaskRunLog.setRunId(jobTask.getRunId());
jobTaskRunLog.setNodeId(jobTask.getNodeId());
jobTaskRunLog.setNodeName(jobTask.getNodeName());
jobTaskRunLog.setJobType(jobTask.getJobType());
......@@ -48,13 +48,20 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
jobTaskRunLog.setRunMsg("上级节点执行失败");
jobTaskRunLog.setRunParams(jobTask.getRunParam());
jobTaskRunLog.setRunCommand(jobTask.getRunCommand());
jobTaskRunLog.setRunType(jobTask.getJobType());
jobTaskRunLog.setRunType("2");
jobTaskRunLog.setTriggerCode("2");
jobTaskRunLog.setTriggerMsg("未执行调度");
jobTaskRunLog.setTriggerTime(thisDate);
jobTaskRunLog.setStartTime(thisDate);
jobTaskRunLog.setEndTime(thisDate);
jobTaskRunLog.setAlertEnd("1");
if("end".equals(jobTask.getNodeName())){
//未完成告警
jobTaskRunLog.setAlertEnd("0");
}else{
//已完成告警
jobTaskRunLog.setAlertEnd("1");
}
//添加日志节点
jobTaskRunLogService.saveJobTaskRunLog(jobTaskRunLog);
}
......
......@@ -79,6 +79,15 @@ public class JavaBeanJobTask implements TimerTask {
jobTaskRunLog.setTriggerTime(new Date());
jobTaskRunLog.setTriggerCode(dispatchResponseDto.getCode());
jobTaskRunLog.setTriggerMsg(dispatchResponseDto.getMsg());
if(!("1".equals(dispatchResponseDto.getCode()))){
Date thisTime = new Date();
jobTaskRunLog.setStartTime(thisTime);
jobTaskRunLog.setEndTime(thisTime);
jobTaskRunLog.setRunCode("2");
jobTaskRunLog.setAlertEnd("1");
jobTaskRunLog.setRunMsg(dispatchResponseDto.getMsg());
}
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLog);
}
}
......@@ -155,8 +155,13 @@ public class JobScheduleHelper{
//记录状态码
String runCode = jobTaskRunLog.getRunCode();
log.info("-------------------【查询到有失败的节点,运行状态为:{}】-------------------",runCode);
//删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask,runCode);
/**
* 判断父类节点是否全部都执行完毕了
*/
if(parentIsEnd(jobTaskRunLogList)){
//删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask,runCode);
}
}
return;
......@@ -428,4 +433,15 @@ 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;
}
}
return flag;
}
}
package com.byit.thread;
import cn.hutool.core.date.DateUtil;
import com.byit.job.dto.JobRunResultDto;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
......@@ -28,7 +29,9 @@ public class LogCallbackThread implements Runnable {
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
jobTaskRunLog.setLogId(jobRunResultDto.getLogId());
jobTaskRunLog.setStartTime(jobRunResultDto.getStartTime());
System.out.println(DateUtil.format(jobRunResultDto.getStartTime(),"yyyy-MM-dd HH:mm:ss"));
jobTaskRunLog.setEndTime(jobRunResultDto.getEndTime());
System.out.println(DateUtil.format(jobRunResultDto.getEndTime(),"yyyy-MM-dd HH:mm:ss"));
jobTaskRunLog.setRunCode(jobRunResultDto.getReturnResult().getCode());
jobTaskRunLog.setRunMsg(jobRunResultDto.getReturnResult().getMsg());
jobTaskRunLog.setAlertEnd("0");
......
......@@ -180,7 +180,7 @@ public class LogScanHelper {
String runId = endNode.getRunId( );
Integer flowId = endNode.getFlowId( );
//如果运行结果为null 那么就是调度都没成功 那么就取调度的值
String code = StringUtils.isNotBlank(endNode.getRunCode())?endNode.getRunCode() : endNode.getTriggerCode();
String code = endNode.getRunCode();
runRecordingService.updateRunRecordingByFlowIdAndRunId(RunRecording.builder().runId(runId).flowId(flowId).flowRunResult(code).flowStatus(RunRecordingEnum.FLOW_STATUS_IS_END.getCode()).build());
endNode.setAlertEnd("1");
jobTaskRunLogService.updateJobTaskRunLog(endNode);
......
......@@ -232,6 +232,7 @@ public class RunRecordingScanHelper {
if (CollectionUtil.isNotEmpty(jobTaskRunLogs)) {
jobTaskRunLogs.forEach(jobTaskRunLog -> {
long timeConsuming = TimeUnit.MILLISECONDS.toMinutes(jobTaskRunLog.getEndTime().getTime() - jobTaskRunLog.getStartTime().getTime());
log.info("######################{}#####################3",timeConsuming);
stringBuilder.append("<tr align='center'>")
.append(String.format("<td>%s</td>", jobTaskRunLog.getNodeName()))
.append(String.format("<td>%s</td>", DateUtil.format(jobTaskRunLog.getStartTime(), DATE_FORMAT)))
......
......@@ -64,7 +64,7 @@
<select id="findJobTaskRunLogEndOrFailureNode" resultMap="BaseResultMap">
select <include refid="Base_Column_List" /> from job_task_run_log
where (node_name='END' or run_code = '2' or run_code = '4' or trigger_code = '2') and alert_end = '0'
where node_name='END' and alert_end = '0' and run_code != '0'
</select>
<select id="findNotEndVirtualNode" resultMap="BaseResultMap">
......
......@@ -45,18 +45,19 @@ public class RunJobThread implements Runnable {
log.info("---------服务器端:{},花费时间:{}-------------", jobRunResultDto,jobRunResultDto.getStartTime().getTime()-jobRunResultDto.getEndTime().getTime());
}
private ReturnResult<String> runJob(String jobHandlerName,String param){
ReturnResult<String> returnResult= ReturnResult.SUCCESS;
private ReturnResult runJob(String jobHandlerName, String param){
Class<? extends IJobHandler> jobClass = JobUtils.jobCache.get(jobHandlerName);
try {
IJobHandler iJobHandler = jobClass.newInstance( );
returnResult = iJobHandler.execute(param);
IJobHandler iJobHandler = jobClass.newInstance();
return iJobHandler.execute(param);
} catch (Exception e) {
e.printStackTrace( );
ReturnResult returnResult = new ReturnResult();
returnResult.setMsg(e.getMessage());
returnResult.setCode(JobResultEnum.FAIL.getRes());
return returnResult;
}
return returnResult;
}
}
......@@ -33,8 +33,8 @@ public class TestAddFlow {
.build();
pluginFlow.setName("测试时间执行的工作流");
pluginFlow.setDesc("这是一个测试的任务流");
pluginFlow.setName("自动化测试原子弹");
pluginFlow.setDesc("自动化测试原子弹");
pluginFlow.setConfig(build);
pluginFlow.setPrincipal("皇甫科星");
pluginFlow.setRePublish(false);
......@@ -73,7 +73,7 @@ public class TestAddFlow {
pluginNode2.setAuthor("皇甫");
pluginNode2.setJobType("JAVA");
pluginNode2.setHandlerName("addJob");
pluginNode2.setRunParam("1");
pluginNode2.setRunParam("add1");
pluginNodeConfig2.setFailedRetryCount(2);
pluginNodeConfig2.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig2.setNodeCron("0 0/7 * * * ? *");
......@@ -111,7 +111,7 @@ public class TestAddFlow {
pluginNode4.setAuthor("皇甫");
pluginNode4.setJobType("JAVA");
pluginNode4.setHandlerName("addJob");
pluginNode4.setRunParam("addJob3");
pluginNode4.setRunParam("add3success");
pluginNodeConfig4.setFailedRetryCount(2);
pluginNodeConfig4.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig4.setNodeCron("0 0/7 * * * ? *");
......
......@@ -2,6 +2,7 @@ package com.byit.job;
import org.apache.commons.exec.CommandLine;
import org.apache.commons.exec.DefaultExecutor;
import org.apache.commons.exec.ExecuteWatchdog;
import org.apache.commons.exec.PumpStreamHandler;
import java.io.*;
......@@ -10,18 +11,23 @@ import java.nio.charset.StandardCharsets;
public class TestPy {
public static void main(String[] args) throws IOException {
public static void main(String[] args) throws IOException, InterruptedException {
ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
PumpStreamHandler pumpStreamHandler = new PumpStreamHandler(outputStream,outputStream,null);
ExecuteWatchdog watchdog = new ExecuteWatchdog(Integer.MAX_VALUE);
DefaultExecutor defaultExecutor = new DefaultExecutor();
defaultExecutor.setWatchdog(watchdog);
defaultExecutor.setExitValues(null);
defaultExecutor.setStreamHandler(pumpStreamHandler);
CommandLine commandline = new CommandLine("python");
commandline.addArgument("D:\\2020project\\byit-myth-job\\demo-client\\byit-demo-client\\src\\main\\resources\\test.py");
commandline.addArgument("D:\\2020project\\byit-myth-job\\demo-client\\byit-demo-client\\test.py");
int execute = defaultExecutor.execute(commandline);
byte[] bytes = outputStream.toByteArray();
Thread.sleep(3000);
watchdog.destroyProcess();
System.out.println(execute);
System.out.println(new String(bytes,0,bytes.length, StandardCharsets.UTF_8));
}
}
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print("--------------------------------------------")
print(100/10)
print(5/0)
\ No newline at end of file
import time
while True :
print("--------------------------------------------")
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