Commit 9f793105 by guominglei

重跑节点和手动置为成功

parent a9b6a88a
...@@ -94,4 +94,11 @@ public class ApiFlowController { ...@@ -94,4 +94,11 @@ public class ApiFlowController {
return "SUCCESS"; return "SUCCESS";
} }
@PostMapping("madeSuccess")
@ApiOperation("手动置为成功")
public String madeSuccess(String param){
apiFlowService.madeSuccess(param);
return "SUCCESS";
}
} }
...@@ -38,4 +38,10 @@ public interface ApiFlowService { ...@@ -38,4 +38,10 @@ public interface ApiFlowService {
* @param param * @param param
*/ */
void reRun(String param); void reRun(String param);
/**
* 手动置为成功
* @param param
*/
void madeSuccess(String param);
} }
...@@ -435,6 +435,15 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -435,6 +435,15 @@ public class ApiFlowServiceImpl implements ApiFlowService {
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, node.getNodeId()); JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, node.getNodeId());
ValidationUtil.dataNotNull(jobTaskRunLog, "查无此运行记录"); ValidationUtil.dataNotNull(jobTaskRunLog, "查无此运行记录");
//校验上级是否成功
List<Integer> dependNodeList = nodeDependencyMapper.findDependIdByNodeId(node.getNodeId());
if (dependNodeList != null && dependNodeList.size() > 0){
dependNodeList.forEach(dependNodeId -> {
JobTaskRunLog dependJobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, dependNodeId);
ValidationUtil.isTrueValidation(!("1".equals(dependJobTaskRunLog.getRunCode()) || "3".equals(dependJobTaskRunLog.getRunCode())), "上级任务未运行成功!");
});
}
//获取当前时间 //获取当前时间
String reRunId = UUID.randomUUID().toString().replace("-",""); String reRunId = UUID.randomUUID().toString().replace("-","");
Long triggerTime = System.currentTimeMillis(); Long triggerTime = System.currentTimeMillis();
...@@ -443,17 +452,79 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -443,17 +452,79 @@ public class ApiFlowServiceImpl implements ApiFlowService {
BeanUtils.copyProperties(node, jobTask); BeanUtils.copyProperties(node, jobTask);
jobTask.setTriggerTime(triggerTime); jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1"); jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId);
jobTask.setReRunId(jobTaskRunLog.getRunId());
jobTaskList.add(jobTask);
//校验通过,开始设置重跑 //校验通过,开始设置重跑
//判断重跑机制(单节点重跑,节点及下游重跑 //判断重跑机制(单节点重跑,节点及下游重跑
if ("1".equals(runState)) {//如果是只重跑当前节点 if (!"1".equals(runState)) {//如果是只重跑当前节点
jobTaskList.add(jobTask); //查询依赖本节点的节点,并添加到集合中
List<Integer> subNodeIdList = nodeDependencyMapper.findSubNodeList(node.getNodeId());
if (subNodeIdList != null && subNodeIdList.size() > 0){
addDependNode(reRunId, runId, triggerTime, jobTaskList, subNodeIdList);
}
} }
jobTaskMapper.saveJobTasks(jobTaskList); jobTaskMapper.saveJobTasks(jobTaskList);
} }
private void addDependNode(String reRunId, String runId, Long triggerTime, List<JobTask> jobTaskList, List<Integer> subNodeIdList) {
subNodeIdList.forEach(childNodeId -> {
Node subNode = nodeMapper.getById(childNodeId);
JobTask jobTask = new JobTask();
BeanUtils.copyProperties(subNode, jobTask);
jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId);
jobTask.setReRunId(runId);
jobTaskList.add(jobTask);
List<Integer> childNodeIdList = nodeDependencyMapper.findSubNodeList(childNodeId);
if (childNodeIdList != null && childNodeIdList.size() > 0){
addDependNode(reRunId, runId, triggerTime, jobTaskList, childNodeIdList);
}
});
}
/**
* 重跑任务
* @param param
*/
@Override
public void madeSuccess(String param){
ValidationUtil.dataNotBank(param, "请求参数不允许为空!");
JSONObject jsonObject = JSON.parseObject(param);
//获取工作空间名称
String workspaceName = jsonObject.getString("workspaceName");
ValidationUtil.dataNotBank(workspaceName, "工作空间名称不允许为空!");
//获取工作流名称
String flowName = jsonObject.getString("flowName");
ValidationUtil.dataNotBank(flowName, "工作流名称不允许为空!");
//获取节点名称
String nodeName = jsonObject.getString("nodeName");
ValidationUtil.dataNotBank(nodeName, "节点名称不允许为空!");
//获取要重跑的运行记录id
String runId = jsonObject.getString("runId");
ValidationUtil.dataNotBank(runId, "运行实例id不允许为空!");
//开始校验
Workspace workspace = workspaceMapper.getByName(workspaceName);
ValidationUtil.dataNotNull(workspace, workspaceName + "工作空间不存在");
Flow flow = flowMapper.getByWorkSpaceAndName(workspace.getWorkspaceId(), flowName);
ValidationUtil.dataNotNull(flow, flowName + "工作流不存在");
Node node = nodeMapper.getByNameAndFlow(nodeName, flow.getFlowId());
ValidationUtil.dataNotNull(node, nodeName + "节点不存在");
//获取运行日志实例
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, node.getNodeId());
ValidationUtil.dataNotNull(jobTaskRunLog, "查无此运行记录");
jobTaskRunLog.setRunCode("1");
jobTaskRunLogMapper.updateJobTaskRunLog(jobTaskRunLog);
}
/** /**
* 工作流生成新版本 * 工作流生成新版本
......
...@@ -27,4 +27,11 @@ public interface NodeDependencyMapper { ...@@ -27,4 +27,11 @@ public interface NodeDependencyMapper {
List<Integer> findDependIdByNodeId(Integer nodeId); List<Integer> findDependIdByNodeId(Integer nodeId);
List<NodeDependencyKey> findByNodeId(Integer nodeId); List<NodeDependencyKey> findByNodeId(Integer nodeId);
/**
* 查询依赖于当前nodeid的节点id
* @param nodeId
* @return
*/
List<Integer> findSubNodeList(Integer nodeId);
} }
\ No newline at end of file
...@@ -19,6 +19,12 @@ ...@@ -19,6 +19,12 @@
where node_id = #{nodeId,jdbcType=INTEGER} where node_id = #{nodeId,jdbcType=INTEGER}
</select> </select>
<select id="findSubNodeList" parameterType="integer" resultType="integer">
select node_id
from node_dependency
where dependency_id = #{nodeId,jdbcType=INTEGER}
</select>
<delete id="deleteById" parameterType="com.byit.model.NodeDependencyKey"> <delete id="deleteById" parameterType="com.byit.model.NodeDependencyKey">
<!-- generated @mbg.generated date: 2019-12-31 --> <!-- generated @mbg.generated date: 2019-12-31 -->
delete from node_dependency delete from node_dependency
......
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