Commit 35caefab by huangfusuper

Merge remote-tracking branch 'origin/developer' into developer

parents 149f9398 9e41a972
......@@ -82,10 +82,17 @@ public class ApiFlowController {
return "SUCCESS";
}
@PostMapping("reRun")
@ApiOperation("重跑")
public String reRun(String param){
apiFlowService.reRun(param);
@PostMapping("reRunJob")
@ApiOperation("重跑节点")
public String reRunJob(String param){
apiFlowService.reRunJob(param);
return "SUCCESS";
}
@PostMapping("reRunFlow")
@ApiOperation("重跑工作流")
public String reRunFlow(String param){
apiFlowService.reRunFlow(param);
return "SUCCESS";
}
......
......@@ -37,11 +37,17 @@ public interface ApiFlowService {
* 重跑任务
* @param param
*/
void reRun(String param);
void reRunJob(String param);
/**
* 手动置为成功
* @param param
*/
void madeSuccess(String param);
/**
* 重跑工作流
* @param param
*/
void reRunFlow(String param);
}
......@@ -416,7 +416,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
* @param param
*/
@Override
public void reRun(String param) {
public void reRunJob(String param) {
ValidationUtil.dataNotBank(param, "请求参数不允许为空!");
JSONObject jsonObject = JSON.parseObject(param);
//获取工作空间名称
......@@ -471,7 +471,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//校验通过,开始设置重跑
//判断重跑机制(单节点重跑,节点及下游重跑
if (!"1".equals(runState)) {//如果是只重跑当前节点
if (!"1".equals(runState)) {//如果是只重跑当前节点
//查询依赖本节点的节点,并添加到集合中
List<Integer> subNodeIdList = nodeDependencyMapper.findSubNodeList(node.getNodeId());
if (subNodeIdList != null && subNodeIdList.size() > 0){
......@@ -483,6 +483,58 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
@Override
public void reRunFlow(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, "工作流名称不允许为空!");
//获取要重跑的运行记录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 + "工作流不存在");
//获取运行日志实例
RunRecording runRecording = runRecordingMapper.findRunRecordingByFlowIdAndRunId(flow.getFlowId(), runId);
ValidationUtil.dataNotNull(runRecording, "没有找到对应的运行记录");
//判断是否是内嵌工作流
if (FlowPropertyEnum.IS_INNER.getCode().equals(flow.getIsInner())){//是内嵌工作流
Node node = nodeMapper.getByMapFlowId(flow.getFlowId());
Map<String, Object> requestMap = new HashMap<>(10);
requestMap.put("workspaceName", workspaceName);
requestMap.put("flowName", flowName);
requestMap.put("runId", runId);
requestMap.put("runState", "1");
requestMap.put("nodeName", node.getNodeName());
reRunJob(JSON.toJSONString(requestMap));
}else {
//不是内嵌工作流
Long triggerTime = System.currentTimeMillis();
String reRunId = UUID.randomUUID().toString().replace("-","");
List<JobTask> jobTaskList = new ArrayList<>();
List<Node> nodeList = nodeMapper.findOnforkByFlowId(flow.getFlowId());
ValidationUtil.dataNotNull(nodeList, "该工作流没有在调度上的任务");
nodeList.forEach(node -> {
JobTask jobTask = new JobTask();
BeanUtils.copyProperties(node, jobTask);
jobTask.setTriggerTime(triggerTime);
jobTask.setTriggerStatus("1");
jobTask.setRunId(reRunId);
jobTask.setReRunId(runId);
jobTaskList.add(jobTask);
});
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);
......
......@@ -78,4 +78,11 @@ public interface NodeMapper {
* @param flowId
*/
int deleteByFlowId(Integer flowId);
/**
* 根据内嵌工作流的id来查找结点
* @param mapFlowId
* @return
*/
Node getByMapFlowId(Integer mapFlowId);
}
\ No newline at end of file
......@@ -114,6 +114,15 @@
where flow_id = #{flowId,jdbcType=INTEGER}
</select>
<select id="getByMapFlowId" parameterType="integer" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from node
where map_flow_id = #{mapFlowId,jdbcType=INTEGER}
</select>
<delete id="deleteVirtualNode">
delete from node
where flow_id = #{flowId,jdbcType=INTEGER}
......
......@@ -37,10 +37,41 @@ public class JobUtils {
*/
private static final String REQUEST_WORKSPACE_ADD = "/api/workspace/add";
/**
* 开始工作流
*/
private static final String REQUEST_FLOW_START = "/api/flow/start";
/**
* 暂停工作流
*/
private static final String REQUEST_FLOW_STOP = "/api/flow/stopSchedule";
/**
* 删除工作流
*/
private static final String REQUEST_FLOW_DELETE = "/api/flow/deleteFlow";
/**
* 撤销工作流调度
*/
private static final String REQUEST_FLOW_REPEAL = "/api/flow/repealSchedule";
/**
* 重新开始撤销的工作流的调度
*/
private static final String REQUEST_FLOW_REREPEAL = "/api/flow/reStartSchedule";
/**
* 杀死任务
*/
private static final String REQUEST_FLOW_KILL_JOB = "/api/flow/killJob";
/**
* 杀死任务
*/
private static final String REQUEST_FLOW_KILL_FLOW = "/api/flow/killFlow";
/**
* 重跑节点
*/
private static final String REQUEST_FLOW_RERUNJOB = "/api/flow/reRunJob";
/**
* 手动置为成功
*/
private static final String REQUEST_FLOW_MAKESUCCESS = "/api/flow/madeSuccess";
private static final String SERVER_PORT = "8998";
/**
* 当前项目运行环境 jar file
......@@ -96,6 +127,24 @@ public class JobUtils {
* @param workspaceName
* @return
*/
public static String startFlow(String flowName, String workspaceName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_START;
Map<String,String> map = new HashMap<>(5);
map.put("flowName",flowName);
map.put("workspaceName",workspaceName);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
log.info("--------------------开始接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 暂停工作流
* @param flowName
* @param workspaceName
* @return
*/
public static String stopFlow(String flowName, String workspaceName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_STOP;
......@@ -105,7 +154,140 @@ public class JobUtils {
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
//String addRequestResult = HttpUtil.post(requestUrl, JSON.toJSONString(pluginPackage, WriteClassName))
log.info("--------------------暂停成功,结果为:{}------------------------",addRequestResult);
log.info("--------------------暂停接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 删除工作流
* @param flowName
* @param workspaceName
* @return
*/
public static String deleteFlow(String flowName, String workspaceName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_DELETE;
Map<String,String> map = new HashMap<>(5);
map.put("flowName",flowName);
map.put("workspaceName",workspaceName);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
log.info("--------------------删除接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 删除工作流
* @param flowName
* @param workspaceName
* @return
*/
public static String repealSchedule(String flowName, String workspaceName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_REPEAL;
Map<String,String> map = new HashMap<>(5);
map.put("flowName",flowName);
map.put("workspaceName",workspaceName);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
log.info("--------------------删除接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 重新开始某次调度
* @param runId
* @return
*/
public static String startSchedule(String runId){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_REREPEAL;
Map<String,Object> map = new HashMap<>(5);
map.put("runId", runId);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl, map);
log.info("--------------------重新开始调度接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 杀死任务
* @param runId
* @param flowName
* @param nodeName
* @return
*/
public static String killJob(String runId, String flowName, String nodeName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_KILL_JOB;
Map<String,String> map = new HashMap<>(5);
map.put("runId", runId);
map.put("flowName", flowName);
map.put("nodeName", nodeName);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
log.info("--------------------杀死任务接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 杀死工作流
* @param runId
* @return
*/
public static String killFlow(String runId){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_KILL_FLOW;
Map<String, Object> map = new HashMap<>(2);
map.put("runId", runId);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl, map);
log.info("--------------------杀死工作流接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 重跑节点
* @param runId
* @param runState
* @param workspaceName
* @param flowName
* @param nodeName
* @return
*/
public static String reRunJob(String runId, String runState, String workspaceName, String flowName, String nodeName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_RERUNJOB;
Map<String,String> map = new HashMap<>(5);
map.put("runId",runId);
map.put("runState",runState);
map.put("workspaceName",workspaceName);
map.put("flowName",flowName);
map.put("nodeName",nodeName);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
log.info("--------------------重跑节点接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
/**
* 手动置为成功
* @param runId
* @param workspaceName
* @param flowName
* @param nodeName
* @return
*/
public static String makeSuccess(String runId, String workspaceName, String flowName, String nodeName){
//请求的路径
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_FLOW_MAKESUCCESS;
Map<String,String> map = new HashMap<>(5);
map.put("runId",runId);
map.put("workspaceName",workspaceName);
map.put("flowName",flowName);
map.put("nodeName",nodeName);
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl,"param="+JSON.toJSONString(map,WriteClassName));
log.info("--------------------手动置为成功接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
......@@ -121,7 +303,7 @@ public class JobUtils {
String requestUrl = REQUEST_PREFIX + "127.0.0.1" + ":" + SERVER_PORT + REQUEST_WORKSPACE_ADD;
//发送请求 添加任务
String addRequestResult = HttpUtil.post(requestUrl, "workspaceName="+ workspaceName);
log.info("--------------------添加任务完成,添加结果为:{}------------------------",addRequestResult);
log.info("--------------------创建工作空间接口调用成功,结果为:{}------------------------",addRequestResult);
return addRequestResult;
}
......
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