Commit d3d4ec0d by huangfusuper

工作流快速失败

parent 3b31a26d
...@@ -3,13 +3,16 @@ package com.byit.service.mapservice; ...@@ -3,13 +3,16 @@ package com.byit.service.mapservice;
import com.byit.model.JobTask; import com.byit.model.JobTask;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import java.net.UnknownHostException;
/** /**
* @author 任务表和日志表的映射业务类,保证是一个事务+原子性 * @author 任务表和日志表的映射业务类,保证是一个事务+原子性
*/ */
public interface TaskAndLogServer { public interface TaskAndLogServer {
/** /**
* 添加失败日志,删除任务表的任务 * 添加失败日志,删除任务表的任务
* @param jobTask * @param jobTask 任务节点
* @param isInner 虚节点
*/ */
void addRunLogAndRemoveTask(JobTask jobTask); void addRunLogAndRemoveTask(JobTask jobTask,boolean isInner) throws UnknownHostException;
} }
package com.byit.service.mapservice.impl; package com.byit.service.mapservice.impl;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.enums.EmailEnum;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.RunRecordingEnum;
import com.byit.job.utils.CronExpression;
import com.byit.job.utils.PlaceholderUtils; import com.byit.job.utils.PlaceholderUtils;
import com.byit.model.JobTask; import com.byit.model.*;
import com.byit.model.JobTaskRunLogWithBLOBs; import com.byit.service.*;
import com.byit.service.JobTaskRunLogService;
import com.byit.service.JobTaskService;
import com.byit.service.mapservice.TaskAndLogServer; import com.byit.service.mapservice.TaskAndLogServer;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
...@@ -12,7 +15,12 @@ import org.springframework.beans.BeanUtils; ...@@ -12,7 +15,12 @@ import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import java.net.InetAddress;
import java.net.UnknownHostException;
import java.text.ParseException;
import java.util.Date; import java.util.Date;
import java.util.List;
import java.util.stream.Collectors;
/** /**
* @author huangfu * @author huangfu
...@@ -23,14 +31,25 @@ import java.util.Date; ...@@ -23,14 +31,25 @@ import java.util.Date;
public class TaskAndLogServerImpl implements TaskAndLogServer { public class TaskAndLogServerImpl implements TaskAndLogServer {
private final JobTaskService jobTaskService; private final JobTaskService jobTaskService;
private final JobTaskRunLogService jobTaskRunLogService; private final JobTaskRunLogService jobTaskRunLogService;
private final FlowService flowService;
private final RunRecordingService runRecordingService;
private final NodeService nodeService;
/**
* 节点依赖查询操作
*/
private final NodeDependencyService nodeDependencyService;
public TaskAndLogServerImpl(JobTaskService jobTaskService, JobTaskRunLogService jobTaskRunLogService) { public TaskAndLogServerImpl(JobTaskService jobTaskService, JobTaskRunLogService jobTaskRunLogService, FlowService flowService, RunRecordingService runRecordingService, NodeService nodeService, NodeDependencyService nodeDependencyService) {
this.jobTaskService = jobTaskService; this.jobTaskService = jobTaskService;
this.jobTaskRunLogService = jobTaskRunLogService; this.jobTaskRunLogService = jobTaskRunLogService;
this.flowService = flowService;
this.runRecordingService = runRecordingService;
this.nodeService = nodeService;
this.nodeDependencyService = nodeDependencyService;
} }
@Override @Override
public void addRunLogAndRemoveTask(JobTask jobTask) { public void addRunLogAndRemoveTask(JobTask jobTask, boolean isInner) throws UnknownHostException {
Date thisDate = new Date(); Date thisDate = new Date();
//删除任务节点 //删除任务节点
jobTaskService.removeMythJobTaskById(jobTask.getId()); jobTaskService.removeMythJobTaskById(jobTask.getId());
...@@ -75,5 +94,67 @@ public class TaskAndLogServerImpl implements TaskAndLogServer { ...@@ -75,5 +94,67 @@ public class TaskAndLogServerImpl implements TaskAndLogServer {
//添加日志节点 //添加日志节点
jobTaskRunLogService.saveJobTaskRunLog(jobTaskRunLog); jobTaskRunLogService.saveJobTaskRunLog(jobTaskRunLog);
if (isInner) {
log.error("------该节点是虚节点,先将该节点对应的实例保存--------");
//查询该节点对应的工作流
Integer mapFlowId = jobTask.getMapFlowId();
Flow virFlow = flowService.findFlowById(mapFlowId);
//保存进运行记录表 快速失败
RunRecording runRecording = RunRecording.builder()
.runId(jobTask.getRunId())
.flowId(jobTask.getMapFlowId())
.flowName(jobTask.getFlowName())
.flowVersionName(jobTask.getVersionName())
.flowStatus("4")
.flowRunResult("2")
.flowTimeout(virFlow.getFlowTimeout())
.dispatchIp(InetAddress.getLocalHost().getHostAddress())
.alarmEmail(virFlow.getAlarmEmail())
.alarmlAction(virFlow.getAlarmlAction())
.priority(virFlow.getPriority())
.triggerTime(jobTask.getTriggerTime())
.principal(virFlow.getPrincipal())
.startTime(new Date())
.flowNodeCount(virFlow.getFlowNodeCount())
.isAlarm(EmailEnum.IS_ALARM_NO.getCode())
.isInner(FlowPropertyEnum.IS_INNER.getCode())
.failFast(RunRecordingEnum.FAIL_FAST_YES.getCode())
.workspaceId(virFlow.getWorkspaceId())
.build();
runRecordingService.saveRunRecording(runRecording);
log.error("在将该虚节点对应的节点拉取过来");
//获取所有的节点,开始将所有节点保存到任务表
List<Node> nodeByFlowIdAndVersionName = nodeService.findNodeByFlowIdAndVersionName(mapFlowId);
List<JobTask> jobTasks = nodeByFlowIdAndVersionName.stream()
.map(node -> {
JobTask task = new JobTask();
List<Integer> dependIdByNodeId = nodeDependencyService.findDependIdByNodeId(jobTask.getNodeId());
if(CollectionUtil.isNotEmpty(dependIdByNodeId)){
String parentIds = StringUtils.join(dependIdByNodeId, ",");
task.setNodeDepend(parentIds);
}
BeanUtils.copyProperties(node, task);
task.setTriggerTime(jobTask.getTriggerTime());
task.setRunId(jobTask.getRunId());
task.setTriggerStatus("1");
task.setFlowName(runRecording.getFlowName());
task.setOperator(jobTask.getOperator());
task.setScheduleType(jobTask.getScheduleType());
return task;
}).collect(Collectors.toList());
jobTaskService.saveJobTasks(jobTasks);
log.info("-----------saveRunRecordingAndTask end【虚节点保存服务】--------------");
if(StringUtils.isNotBlank(virFlow.getFlowCron())){
try {
virFlow.setTriggerNextTime(new CronExpression(virFlow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
} catch (ParseException e) {
e.printStackTrace();
}
flowService.updateByIdSelective(virFlow);
}
}
} }
} }
...@@ -21,6 +21,7 @@ import org.springframework.beans.BeanUtils; ...@@ -21,6 +21,7 @@ import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import javax.sql.DataSource; import javax.sql.DataSource;
import java.net.UnknownHostException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Date; import java.util.Date;
import java.util.List; import java.util.List;
...@@ -105,7 +106,11 @@ public class JobTaskThreadRunHelper extends BaseThreadRunHelper { ...@@ -105,7 +106,11 @@ public class JobTaskThreadRunHelper extends BaseThreadRunHelper {
} }
}catch (SuperiorNodeRunException se) { }catch (SuperiorNodeRunException se) {
//删除这个数据 并且添加到日志 //删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask); try {
taskAndLogServer.addRunLogAndRemoveTask(jobTask,true);
} catch (UnknownHostException e) {
e.printStackTrace();
}
} catch (Exception e) { } catch (Exception e) {
e.printStackTrace(); e.printStackTrace();
} }
...@@ -125,7 +130,11 @@ public class JobTaskThreadRunHelper extends BaseThreadRunHelper { ...@@ -125,7 +130,11 @@ public class JobTaskThreadRunHelper extends BaseThreadRunHelper {
} }
}catch (SuperiorNodeRunException se) { }catch (SuperiorNodeRunException se) {
//删除这个数据 并且添加到日志 //删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(jobTask); try {
taskAndLogServer.addRunLogAndRemoveTask(jobTask,false);
} catch (UnknownHostException e) {
e.printStackTrace();
}
} }
} }
} }
......
...@@ -23,6 +23,7 @@ import org.springframework.beans.BeanUtils; ...@@ -23,6 +23,7 @@ import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import javax.sql.DataSource; import javax.sql.DataSource;
import java.net.UnknownHostException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Date; import java.util.Date;
import java.util.List; import java.util.List;
...@@ -166,7 +167,11 @@ public class TaskThreadRunHelper extends BaseThreadRunHelper { ...@@ -166,7 +167,11 @@ public class TaskThreadRunHelper extends BaseThreadRunHelper {
//如果有失败的节点 就把该节点置为失败 //如果有失败的节点 就把该节点置为失败
if(CollectionUtil.isNotEmpty(errorJobLog)){ if(CollectionUtil.isNotEmpty(errorJobLog)){
//删除这个数据 并且添加到日志 //删除这个数据 并且添加到日志
taskAndLogServer.addRunLogAndRemoveTask(thisJobTask); try {
taskAndLogServer.addRunLogAndRemoveTask(thisJobTask,false);
} catch (UnknownHostException e) {
e.printStackTrace();
}
}else{ }else{
//执行代码 //执行代码
runJobTask(thisJobTask,jobTaskSchedules); runJobTask(thisJobTask,jobTaskSchedules);
......
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