Commit cd5f78bd by huangfusuper

执行节点监控

parent b297e9c6
......@@ -107,9 +107,8 @@ public class ApiFlowController {
@PostMapping("reStartSchedule")
@ApiOperation("重新开始某次调度")
public ResponseResult reStartSchedule(@RequestBody StopFlowParam stopFlowParam) {
public void reStartSchedule(@RequestBody StopFlowParam stopFlowParam) {
apiFlowService.reStartSchedule(stopFlowParam);
return ResponseResult.ok("SUCCESS");
}
@PostMapping("reRunJob")
......
......@@ -15,8 +15,13 @@ import java.util.List;
*/
@Repository
public interface JobTaskRunLogMapper {
/**
* 查询立即运行节点未完成的
* @return 所有运行中的
*/
List<JobTaskRunLog> findTempRunNodeList();
public List<JobTaskRunLogWithBLOBs> findAllByType(@Param("type") Integer type,@Param("jobName") String jobName);
List<JobTaskRunLogWithBLOBs> findAllByType(@Param("type") Integer type,@Param("jobName") String jobName);
List<JobTaskRunLog> findImmediatelyNode(String nodeName);
/**
......
......@@ -14,6 +14,12 @@ import java.util.List;
**/
public interface JobTaskRunLogService {
/**
* 查询立即运行节点未完成的
* @return 所有运行中的
*/
List<JobTaskRunLog> findTempRunNodeList();
List<JobTaskRunLogWithBLOBs> findAllByType(Integer type,String jobName);
List<JobTaskRunLog> findImmediatelyNode(String nodeName);
......
......@@ -41,6 +41,11 @@ public class JobTaskRunLogServiceImpl implements JobTaskRunLogService {
}
@Override
public List<JobTaskRunLog> findTempRunNodeList() {
return jobTaskRunLogMapper.findTempRunNodeList();
}
@Override
public List<JobTaskRunLogWithBLOBs> findAllByType(Integer type,String jobName) {
return jobTaskRunLogMapper.findAllByType(type,jobName);
......
package com.byit.thread.helper;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.model.JobTaskRunLog;
import com.byit.service.JobTaskRunLogService;
import com.byit.thread.BaseThreadRunHelper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import java.util.List;
/**
* @author huangfu
*/
@Slf4j
@Component
public class UnfinishedNodeMonitoringRunHelper extends BaseThreadRunHelper {
/**
* 心跳
*/
private static final String PENG = "PENG";
private final StringRedisTemplate stringRedisTemplate;
private final JobTaskRunLogService jobTaskRunLogService;
public UnfinishedNodeMonitoringRunHelper(StringRedisTemplate stringRedisTemplate, JobTaskRunLogService jobTaskRunLogService) {
this.stringRedisTemplate = stringRedisTemplate;
this.jobTaskRunLogService = jobTaskRunLogService;
}
@Override
public Long start() {
List<JobTaskRunLog> tempRunNodeList = jobTaskRunLogService.findTempRunNodeList();
if (CollectionUtil.isNotEmpty(tempRunNodeList)) {
tempRunNodeList.forEach(nodeLog -> stringRedisTemplate.opsForList().rightPush(nodeLog.getRunId(), PENG));
}
return 5000L;
}
@Override
public String getLockName() {
return "UnfinishedNodeMonitoringRunHelper";
}
}
......@@ -49,6 +49,10 @@
run_msg, trigger_msg
</sql>
<select id="findTempRunNodeList" resultMap="BaseResultMap">
select <include refid="Base_Column_List" /> from job_task_run_log where schedule_type = 4 and run_code = '0'
</select>
<!--根据runId 和 上级节点的集合 查询所有的上级节点-->
<select id="findJobTaskRunLogNotEndNodeByRunCodeCount" resultMap="BaseResultMap">
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