Commit 6fdadeb5 by guo_minglei@163.com

杀死进程开发

parent 6f0b1eff
......@@ -86,4 +86,21 @@ public interface JobTaskRunLogMapper {
* @return
*/
JobTaskRunLog findByRunIdAndNodeId(@Param("runId")String runId, @Param("nodeId")Integer nodeId);
/**
* 根据runid和flowname和nodename获取一条运行记录
* @param runId
* @param flowName
* @param nodeName
* @return
*/
JobTaskRunLog findByRunIdAndFlowAndNode(@Param("runId")String runId, @Param("flowName")String flowName, @Param("nodeName")String nodeName);
/**
* 根据runid和flowname获取运行记录
* @param runId
* @param flowName
* @return
*/
List<JobTaskRunLog> findByRunIdAndFlowName(@Param("runId")String runId, @Param("flowName")String flowName);
}
\ No newline at end of file
......@@ -96,4 +96,19 @@ public interface RunRecordingMapper {
*/
List<RunRecording> findStopRunCordByRunId(String runId);
/**
* 杀死整个runid
* @param runId
* @return
*/
int killFlow(String runId);
/**
* 杀死工作流内的某个内嵌工作流
* @param runId
* @param flowName
* @return
*/
int killInnerFlow(String runId, String flowName);
}
\ No newline at end of file
......@@ -74,9 +74,9 @@ public class JobTaskRunLog implements Serializable {
private String isVirtual;
/**
* 运行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill
* 运行结果 0运行中 1 成功 2 失败 3 补批成功 4 补批失败 5.kill 6.上级节点执行失败
*/
@ApiModelProperty("运行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill 6.上级节点执行失败")
@ApiModelProperty("运行结果 0运行中 1 成功 2 失败 3 补批成功 4 补批失败 5.kill 6.上级节点执行失败")
private String runCode;
/**
......
......@@ -83,4 +83,13 @@ public interface FlowService {
* @return
*/
Boolean includeFlow(FlowVo flowVo, NodeVo nodeVo);
/**
* 杀死一条任务
* @param runId 运行实例id
* @param flowName 工作流名称
* @param nodeName 任务名称
* @return
*/
Boolean killJob(String runId, String flowName, String nodeName) throws InterruptedException;
}
package com.byit.service.impl;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodePropertyEnum;
import com.byit.job.utils.CronExpression;
import com.byit.mapper.FlowMapper;
import com.byit.mapper.NodeDependencyMapper;
import com.byit.mapper.NodeMapper;
import com.byit.job.vo.ResponseResult;
import com.byit.mapper.*;
import com.byit.model.Flow;
import com.byit.model.JobTaskRunLog;
import com.byit.model.Node;
import com.byit.model.vo.FlowVo;
import com.byit.model.vo.NodeVo;
......@@ -20,9 +22,7 @@ import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.*;
/**
* @description: 工作流业务逻辑实现类
......@@ -42,6 +42,10 @@ public class FlowServiceImpl implements FlowService {
private NodeMapper nodeMapper;
@Resource
private NodeDependencyMapper nodeDependencyMapper;
@Resource
private JobTaskRunLogMapper jobTaskRunLogMapper;
@Resource
private RunRecordingMapper runRecordingMapper;
@Override
public Flow findFlowById(Integer id) {
......@@ -297,4 +301,59 @@ public class FlowServiceImpl implements FlowService {
flowMapper.updateByIdSelective(virtualFlow);
return null;
}
@Override
public Boolean killJob(String runId, String flowName, String nodeName) throws InterruptedException {
ValidationUtil.dataNotBank(runId, "运行实例id不允许为空!");
ValidationUtil.dataNotBank(flowName, "工作流名称不允许为空!");
ValidationUtil.dataNotBank(nodeName, "节点名称不允许为空!");
Boolean result = true;
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndFlowAndNode(runId, flowName, nodeName);
ValidationUtil.dataNotNull(jobTaskRunLog, "不存在【" + flowName + "】下节点【" + nodeName + "】的运行日志,或尚未开始调度!");
ValidationUtil.isTrueValidation(!"0".equals(jobTaskRunLog.getRunCode()), "该任务已经运行结束!");
if (FlowPropertyEnum.IS_INNER.getCode().equals(jobTaskRunLog.getIsVirtual())){
runRecordingMapper.killInnerFlow(runId, flowName);
List<JobTaskRunLog> jobTaskRunLogList = jobTaskRunLogMapper.findByRunIdAndFlowName(runId, flowName);
ValidationUtil.dataNotNull(jobTaskRunLogList, "该任务尚未开始调度");
for (JobTaskRunLog innerJobTaskRunLog : jobTaskRunLogList){
if (!"0".equals(innerJobTaskRunLog.getRunCode())){
if (StringUtils.isEmpty(innerJobTaskRunLog.getJobGroupIp())){
log.warn("已经调度成功但是还未返回具体的调用机器的ip地址");
Thread.sleep(3000);
innerJobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndFlowAndNode(runId, flowName, nodeName);
ValidationUtil.isTrueValidation(!"0".equals(innerJobTaskRunLog.getRunCode()), "该任务已经运行结束!");
ValidationUtil.dataNotBank(innerJobTaskRunLog.getJobGroupIp(), "尚未分配执行机请稍后再试");
}
if (! killJob(innerJobTaskRunLog.getLogId(), innerJobTaskRunLog.getJobGroupIp())){
result = false;
}
}
}
}
if (StringUtils.isEmpty(jobTaskRunLog.getJobGroupIp())){
log.warn("已经调度成功但是还未返回具体的调用机器的ip地址");
Thread.sleep(3000);
jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndFlowAndNode(runId, flowName, nodeName);
ValidationUtil.isTrueValidation(!"0".equals(jobTaskRunLog.getRunCode()), "该任务已经运行结束!");
ValidationUtil.dataNotBank(jobTaskRunLog.getJobGroupIp(), "尚未分配执行机请稍后再试");
}
result = killJob(jobTaskRunLog.getLogId(), jobTaskRunLog.getJobGroupIp());
return result;
}
private synchronized Boolean killJob(Integer logId, String exectUrl){
//具体访问的URL
String killUrl = exectUrl + "/processManager/killJob";
Map<String, Integer> requestMap = new HashMap<>(2);
requestMap.put("logId", logId);
String killResult = HttpUtil.post(killUrl, JSON.toJSONString(requestMap));
log.info("请求结果{}", killResult);
ResponseResult responseResult = JSON.parseObject(killResult, ResponseResult.class);
if ("SUCCESS".equals(responseResult.getResult())){
return true;
}
return false;
}
}
\ No newline at end of file
......@@ -70,6 +70,19 @@
where node_id = #{nodeId} and run_id = #{runId}
</select>
<select id="findByRunIdAndFlowAndNode" resultMap="BaseResultMap">
select <include refid="Base_Column_List"/>
from job_task_run_log
where node_id = #{nodeId} and flow_name = #{flowName}
and node_name = #{nodeName}
</select>
<select id="findByRunIdAndFlowName" resultMap="BaseResultMap">
select <include refid="Base_Column_List"/>
from job_task_run_log
where node_id = #{nodeId} and flow_name = #{flowName}
</select>
<select id="findJobTaskRunLogWithBLOBsByFlowIdAndRunId" resultMap="ResultMapWithBLOBs">
select
<include refid="Base_Column_List" />
......
......@@ -341,4 +341,14 @@
and flow_status = '3'
</update>
<update id="killFlow" parameterType="string">
update run_recording set flow_status = '4',flow_run_result = '5'
where run_id = #{runId,jdbcType=VARCHAR}
</update>
<update id="killInnerFlow" >
update run_recording set flow_status = '4',flow_run_result = '5'
where run_id = #{runId,jdbcType=VARCHAR} and flow_name = ${flowName,jdbcType=VARCHAR}
</update>
</mapper>
\ No newline at end of file
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