Commit 8fefb13e by guominglei

插件端重跑节点开发

parent 758cc4dd
......@@ -34,21 +34,30 @@ public class ApiFlowController {
@PostMapping("delete")
@ApiOperation("删除工作流,若当前工作流被依赖则删除失败,只允许删除不被依赖的工作流,当前工作流依赖其他工作流不影响")
public String deleteFlow(String flowName, String workspaceName){
public String deleteFlow(String param){
JSONObject jsonObject = JSON.parseObject(param);
String flowName = jsonObject.getString("flowName");
String workspaceName = jsonObject.getString("workspaceName");
apiFlowService.deleteFlow(flowName, workspaceName);
return "SUCCESS";
}
@PostMapping("start")
@ApiOperation("开始工作流的调度,将工作流启用调度")
public String start(String flowName, String workspaceName) throws ParseException {
public String start(String param) throws ParseException {
JSONObject jsonObject = JSON.parseObject(param);
String flowName = jsonObject.getString("flowName");
String workspaceName = jsonObject.getString("workspaceName");
apiFlowService.start(flowName, workspaceName);
return "SUCCESS";
}
@PostMapping("repealSchedule")
@ApiOperation("撤销工作流调度,只可以撤销总工作流,若为内嵌工作流不允许撤销")
public String repealSchedule(String flowName, String workspaceName) throws ParseException {
public String repealSchedule(String param) throws ParseException {
JSONObject jsonObject = JSON.parseObject(param);
String flowName = jsonObject.getString("flowName");
String workspaceName = jsonObject.getString("workspaceName");
apiFlowService.repealSchedule(flowName, workspaceName);
return "SUCCESS";
}
......@@ -78,4 +87,11 @@ public class ApiFlowController {
return "SUCCESS";
}
@PostMapping("reRun")
@ApiOperation("重跑")
public String reRun(String param){
apiFlowService.reRun(param);
return "SUCCESS";
}
}
......@@ -32,4 +32,10 @@ public interface ApiFlowService {
* @return
*/
void validateFlow(Integer workspaceId, PluginFlow flow);
/**
* 重跑任务
* @param param
*/
void reRun(String param);
}
package com.byit.service.impl;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.byit.enums.DagCheckEnum;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodePropertyEnum;
......@@ -62,6 +63,9 @@ public class ApiFlowServiceImpl implements ApiFlowService {
@Resource
private JobTaskMapper jobTaskMapper;
@Resource
private JobTaskRunLogMapper jobTaskRunLogMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public void publishFlow(String param) throws Exception {
......@@ -395,6 +399,48 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
}
/**
* 重跑任务
* @param param
*/
@Override
public void reRun(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, "节点名称不允许为空!");
//获取重跑机制(运行当前节点,或运行当前节点及以下节点)
String runState = jsonObject.getString("runState");
ValidationUtil.dataNotBank(runState, "重跑机制不允许为空!");
//获取要重跑的运行记录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, "查无此运行记录");
//校验通过,开始设置重跑
//判断重跑机制(单节点重跑,节点及下游重跑)
}
/**
* 工作流生成新版本
......
......@@ -73,4 +73,11 @@ public interface JobTaskRunLogMapper {
*/
int updateJobTaskRunLog(JobTaskRunLog record);
/**
* 根据runid和nodeid获取运行日志
* @param runId
* @param nodeId
* @return
*/
JobTaskRunLog findByRunIdAndNodeId(@Param("runId")String runId, @Param("nodeId")Integer nodeId);
}
\ No newline at end of file
......@@ -56,6 +56,12 @@
</foreach>
</select>
<select id="findByRunIdAndNodeId" resultMap="BaseResultMap">
select <include refid="Base_Column_List"/>
from job_task_run_log
where node_id = #{nodeId} and run_id = #{runId}
</select>
<select id="findJobTaskRunLogWithBLOBsByFlowIdAndRunId" resultMap="ResultMapWithBLOBs">
select
<include refid="Base_Column_List" />
......
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