Commit 22c83f51 by huangfusuper

【修改遗留】修改添加任务时,判断上级节点是否完成

parent 2bfb89a8
......@@ -14,6 +14,13 @@ import java.util.List;
@Repository
public interface JobTaskRunLogMapper {
/**
* 查询没有结束的节点
* @param nodIds
* @return
*/
int findJobTaskRunLogNotEndNodeByRunCodeCount(@Param("nodIds") List<Integer> nodIds);
/**
* 查根据flowId和RunId查询一批节点
* @param flowId
* @param runId
......
......@@ -2,6 +2,7 @@ package com.byit.service;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskRunLogWithBLOBs;
import org.apache.ibatis.annotations.Param;
import java.util.List;
......@@ -12,6 +13,13 @@ import java.util.List;
* @date: 2019/12/20 19:42
**/
public interface JobTaskRunLogService {
/**
* 查询没有结束的节点
* @param nodIds
* @return
*/
int findJobTaskRunLogNotEndNodeByRunCodeCount(List<Integer> nodIds);
/**
* 查根据flowId和RunId查询一批节点
* @param flowId
......
package com.byit.service;
import java.util.List;
/**
* @author huangfu
*/
public interface NodeDependencyService {
/**
* 根据节点id查询本节点依赖的节点
* @param nodeId
* @return
*/
List<Integer> findDependIdByNodeId(Integer nodeId);
}
\ No newline at end of file
......@@ -6,4 +6,5 @@ package com.byit.service;
* @create: 2019-12-24 10:38
*/
public interface NodeService {
}
......@@ -29,6 +29,11 @@ public class JobTaskRunLogServiceImpl implements JobTaskRunLogService {
}
@Override
public int findJobTaskRunLogNotEndNodeByRunCodeCount(List<Integer> nodIds) {
return jobTaskRunLogMapper.findJobTaskRunLogNotEndNodeByRunCodeCount(nodIds);
}
@Override
public List<JobTaskRunLogWithBLOBs> findJobTaskRunLogWithBLOBsByFlowIdAndRunId(Integer flowId, String runId) {
return jobTaskRunLogMapper.findJobTaskRunLogWithBLOBsByFlowIdAndRunId(flowId,runId);
}
......
package com.byit.service.impl;
import com.byit.mapper.NodeDependencyMapper;
import com.byit.service.NodeDependencyService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import javax.xml.ws.Action;
import java.util.List;
/**
* @author Administrator
*/
@Service
@Slf4j
public class NodeDependencyServiceImpl implements NodeDependencyService {
private final NodeDependencyMapper nodeDependencyMapper;
public NodeDependencyServiceImpl(NodeDependencyMapper nodeDependencyMapper) {
this.nodeDependencyMapper = nodeDependencyMapper;
}
@Override
public List<Integer> findDependIdByNodeId(Integer nodeId) {
return nodeDependencyMapper.findDependIdByNodeId(nodeId);
}
}
......@@ -5,8 +5,10 @@ import com.byit.job.WorkRoulette;
import com.byit.model.JobTask;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.model.JobTaskSchedule;
import com.byit.service.JobTaskRunLogService;
import com.byit.service.JobTaskScheduleService;
import com.byit.service.JobTaskService;
import com.byit.service.NodeDependencyService;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.task.JavaBeanJobTask;
import com.byit.util.SpringUtil;
......@@ -39,6 +41,8 @@ public class JobScheduleHelper{
private final JobTaskService jobTaskService;
private final JobTaskScheduleService jobTaskScheduleService;
private final NodeDependencyService nodeDependencyService;
private final JobTaskRunLogService jobTaskRunLogService;
/**
* 读取任务节点的预读
......@@ -67,15 +71,25 @@ public class JobScheduleHelper{
private volatile boolean scheduleThreadToStop = false;
@Autowired
public JobScheduleHelper(JobTaskScheduleService jobTaskScheduleService, JobTaskService jobTaskService) {
public JobScheduleHelper(JobTaskScheduleService jobTaskScheduleService, JobTaskService jobTaskService, NodeDependencyService nodeDependencyService, JobTaskRunLogService jobTaskRunLogService) {
this.jobTaskScheduleService = jobTaskScheduleService;
this.jobTaskService = jobTaskService;
this.nodeDependencyService = nodeDependencyService;
this.jobTaskRunLogService = jobTaskRunLogService;
}
/**
* 启动两条线程
*/
public void start(){
jobInfoThreadStart();
scheduleThreadStart();
}
/**
* 任务节点扫描
*/
public void jobInfoThreadStart(){
jobInfoThread = new Thread(()->{
try {
TimeUnit.MILLISECONDS.sleep(PRE_READ_MS - System.currentTimeMillis()%1000 );
......@@ -117,12 +131,23 @@ public class JobScheduleHelper{
* 大概思路,根据任务流id,从任务流执行回溯表查询该任务流的所有节点,查看上级节点是否已经执行成功
* //TODO 需要修改 判断父节点是否执行完毕 注意 父节点是一个集合
*/
if(true){
//虚节点的状态
if("0".equals(jobTask.getIsVirtual())){
//在日志表里面创建一条记录
//根据 map_flow_id查询当前的版本的工作流 使用祝工作流的runId 保存到执行记录表和任务表
}else{
List<Integer> dependIdByNodeId = nodeDependencyService.findDependIdByNodeId(jobTask.getNodeId());
int jobTaskRunLogNotEndNodeByRunCodeCount = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(dependIdByNodeId);
if("start".equals(jobTask.getNodeName()) || (jobTaskRunLogNotEndNodeByRunCodeCount==0)){
log.debug("任务:{}", jobTask);
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
}
}
});
if(CollectionUtil.isNotEmpty(jobTaskSchedules)) {
......@@ -199,7 +224,7 @@ public class JobScheduleHelper{
//设置名字
jobInfoThread.setName("myth-job,admin JobScheduleHelper#jobInfoThread");
jobInfoThread.start();
}
/**
* 开始操作任务排期表:
......@@ -208,6 +233,8 @@ public class JobScheduleHelper{
* 3.加载到任务调度轮盘
* 4.删除数据
*/
private void scheduleThreadStart(){
scheduleThread = new Thread(() ->{
try {
TimeUnit.MILLISECONDS.sleep(4000 - System.currentTimeMillis()%1000 );
......@@ -320,6 +347,7 @@ public class JobScheduleHelper{
scheduleThread.setDaemon(true);
scheduleThread.start();
}
public void doStop(){
this.jobInfoThreadToStop = true;
try {
......
......@@ -39,6 +39,15 @@
<sql id="Blob_Column_List">
run_msg, trigger_msg
</sql>
<select id="findJobTaskRunLogNotEndNodeByRunCodeCount" resultType="java.lang.Integer">
select count(id) from job_task_run_log where node_id in
<foreach item="nodeId" collection="nodIds" open="(" separator="," close=")">
#{nodeId}
</foreach>
and (run_code = '0' or run_code = '5')
</select>
<select id="findJobTaskRunLogWithBLOBsByFlowIdAndRunId" resultMap="ResultMapWithBLOBs">
select
<include refid="Base_Column_List" />
......@@ -56,7 +65,7 @@
<select id="findNotEndVirtualNode" resultMap="BaseResultMap">
select <include refid="Base_Column_List" /> from job_task_run_log
where is_virtual='0' and (run_code is null or run_code = '')
where is_virtual='0' and run_code = '0'
</select>
<select id="findJobTaskRunLogByLogId" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
......
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