Commit 2abfbe8b by guo_minglei@163.com

修改重跑工作流

parent 7091be9f
...@@ -7,7 +7,7 @@ import com.byit.dto.plugin.*; ...@@ -7,7 +7,7 @@ import com.byit.dto.plugin.*;
import com.byit.enums.DagCheckEnum; import com.byit.enums.DagCheckEnum;
import com.byit.enums.FlowPropertyEnum; import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodePropertyEnum; import com.byit.enums.NodePropertyEnum;
import com.byit.enums.ScheduleEnum; import com.byit.enums.ScheduleTypeEnum;
import com.byit.enums.plugin.PluginNodeTypeEnum; import com.byit.enums.plugin.PluginNodeTypeEnum;
import com.byit.job.utils.CronExpression; import com.byit.job.utils.CronExpression;
import com.byit.job.utils.CurrentUserUtils; import com.byit.job.utils.CurrentUserUtils;
...@@ -485,7 +485,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -485,7 +485,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
jobTask.setRunId(reRunId); jobTask.setRunId(reRunId);
jobTask.setReRunId(jobTaskRunLog.getRunId()); jobTask.setReRunId(jobTaskRunLog.getRunId());
//设置为重跑 //设置为重跑
jobTask.setScheduleType(ScheduleEnum.REPEAT.getCode()); jobTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
jobTask.setOperator(userName); jobTask.setOperator(userName);
//查询当前节点的依赖节点 //查询当前节点的依赖节点
...@@ -540,22 +540,24 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -540,22 +540,24 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//不是内嵌工作流 //不是内嵌工作流
Long triggerTime = System.currentTimeMillis(); Long triggerTime = System.currentTimeMillis();
String reRunId = UUID.randomUUID().toString().replace("-",""); String reRunId = UUID.randomUUID().toString().replace("-","");
List<JobTask> jobTaskList = new ArrayList<>(); List<JobTask> jobTaskList = new ArrayList<>();
List<Node> nodeList = nodeMapper.findOnforkByFlowId(flow.getFlowId()); List<JobTaskRunLogWithBLOBs> jobTaskRunLogList = jobTaskRunLogMapper.findJobTaskRunLogWithBLOBsByFlowIdAndRunId(flow.getFlowId(), runInfo.getRunId());
ValidationUtil.dataNotNull(nodeList, "该工作流没有在调度上的任务"); ValidationUtil.dataNotNull(jobTaskRunLogList, "该工作流没有在调度上的任务");
String userName = currentUserUtils.account(); String userName = currentUserUtils.account();
nodeList.forEach(node -> {
jobTaskRunLogList.forEach(jobTaskRunLog -> {
JobTask jobTask = new JobTask(); JobTask jobTask = new JobTask();
BeanUtils.copyProperties(node, jobTask); BeanUtils.copyProperties(jobTaskRunLog, jobTask);
jobTask.setTriggerTime(triggerTime); jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1"); jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId); jobTask.setRunId(reRunId);
jobTask.setReRunId(runInfo.getRunId()); jobTask.setReRunId(runInfo.getRunId());
//设置为重跑 //设置为重跑
jobTask.setScheduleType(ScheduleEnum.REPEAT.getCode()); jobTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
jobTask.setOperator(userName); jobTask.setOperator(userName);
//查询当前节点的依赖节点 //查询当前节点的依赖节点
List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(node.getNodeId()); List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(jobTaskRunLog.getNodeId());
if (dependNodeIdList != null && dependNodeIdList.size() > 0){ if (dependNodeIdList != null && dependNodeIdList.size() > 0){
jobTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList)); jobTask.setNodeDepend(Joiner.on(",").join(dependNodeIdList));
} }
...@@ -575,7 +577,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -575,7 +577,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
jobTask.setRunId(reRunId); jobTask.setRunId(reRunId);
jobTask.setReRunId(runId); jobTask.setReRunId(runId);
//设置为重跑 //设置为重跑
jobTask.setScheduleType(ScheduleEnum.REPEAT.getCode()); jobTask.setScheduleType(ScheduleTypeEnum.REPEAT.getCode());
jobTask.setOperator(userName); jobTask.setOperator(userName);
//查询当前节点的依赖节点 //查询当前节点的依赖节点
...@@ -934,7 +936,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -934,7 +936,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
waitingRecord.setRunId(runId); waitingRecord.setRunId(runId);
waitingRecord.setFlowNodeCount(nodeList.size()); waitingRecord.setFlowNodeCount(nodeList.size());
waitingRecord.setOperator(userName); waitingRecord.setOperator(userName);
waitingRecord.setScheduleType(ScheduleEnum.REPAIR.getCode()); waitingRecord.setScheduleType(ScheduleTypeEnum.REPAIR.getCode());
waitingRecord.setWaitOrder(++order); waitingRecord.setWaitOrder(++order);
//设置实例 //设置实例
BeanUtils.copyProperties(waitingRecord, runRecording); BeanUtils.copyProperties(waitingRecord, runRecording);
...@@ -972,7 +974,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -972,7 +974,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
waitingTask.setNodeDepend(Joiner.on(",").join(dependNodeNameList)); waitingTask.setNodeDepend(Joiner.on(",").join(dependNodeNameList));
} }
waitingTask.setOperator(usernName); waitingTask.setOperator(usernName);
waitingTask.setScheduleType(ScheduleEnum.REPAIR.getCode()); waitingTask.setScheduleType(ScheduleTypeEnum.REPAIR.getCode());
waitingTaskList.add(waitingTask); waitingTaskList.add(waitingTask);
}); });
......
...@@ -4,7 +4,7 @@ import com.alibaba.fastjson.JSON; ...@@ -4,7 +4,7 @@ import com.alibaba.fastjson.JSON;
import com.byit.dto.plugin.RunLog; import com.byit.dto.plugin.RunLog;
import com.byit.dto.plugin.RunNode; import com.byit.dto.plugin.RunNode;
import com.byit.enums.NodePropertyEnum; import com.byit.enums.NodePropertyEnum;
import com.byit.enums.ScheduleEnum; import com.byit.enums.ScheduleTypeEnum;
import com.byit.job.WorkRoulette; import com.byit.job.WorkRoulette;
import com.byit.job.utils.PlaceholderUtils; import com.byit.job.utils.PlaceholderUtils;
import com.byit.mapper.JobTaskRunLogMapper; import com.byit.mapper.JobTaskRunLogMapper;
...@@ -63,7 +63,7 @@ public class ApiNodeServiceImpl implements ApiNodeService { ...@@ -63,7 +63,7 @@ public class ApiNodeServiceImpl implements ApiNodeService {
schedule.setScriptUrls(runNode.getScriptUrl()); schedule.setScriptUrls(runNode.getScriptUrl());
schedule.setJobType(runNode.getJobType()); schedule.setJobType(runNode.getJobType());
schedule.setTriggerTime(triggerTime); schedule.setTriggerTime(triggerTime);
schedule.setScheduleType(ScheduleEnum.REAL.getCode()); schedule.setScheduleType(ScheduleTypeEnum.REAL.getCode());
schedule.setRunParam(PlaceholderUtils.formatParam(schedule.getRunParam())); schedule.setRunParam(PlaceholderUtils.formatParam(schedule.getRunParam()));
......
...@@ -5,7 +5,7 @@ package com.byit.enums; ...@@ -5,7 +5,7 @@ package com.byit.enums;
* @author: gml * @author: gml
* @create: 2020/3/9 * @create: 2020/3/9
*/ */
public enum ScheduleEnum { public enum ScheduleTypeEnum {
NORMAL(1, "正常跑批"), NORMAL(1, "正常跑批"),
REPEAT(2, "重跑"), REPEAT(2, "重跑"),
...@@ -16,7 +16,7 @@ public enum ScheduleEnum { ...@@ -16,7 +16,7 @@ public enum ScheduleEnum {
private Integer code; private Integer code;
private String msg; private String msg;
private ScheduleEnum(Integer code, String msg){ private ScheduleTypeEnum(Integer code, String msg){
this.code = code; this.code = code;
this.msg = msg; this.msg = msg;
} }
......
...@@ -150,6 +150,11 @@ public class RunRecording implements Serializable { ...@@ -150,6 +150,11 @@ public class RunRecording implements Serializable {
@ApiModelProperty("重跑和补批的操作人") @ApiModelProperty("重跑和补批的操作人")
private String operator; private String operator;
/**
* 重跑的运行标识
*/
@ApiModelProperty("重跑的运行标识")
private String reRunId;
/** /**
*/ */
......
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