Commit 680f9c0e by guominglei

插件端开启调度和撤销调度

parent e124cbe9
package com.byit.mapper;
import com.byit.model.JobTask;
import com.byit.model.JobTaskSchedule;
import org.apache.ibatis.annotations.Param;
import org.springframework.stereotype.Repository;
......@@ -55,9 +56,9 @@ public interface JobTaskMapper {
/**
* 根据id删除多个
* @param jobTasks
* @param jobTaskSchedules
*/
void deleteInId(@Param("jobTasks") List<JobTask> jobTasks);
void deleteInId(@Param("jobTaskSchedules") List<JobTaskSchedule> jobTaskSchedules);
/**
* 根据runid删除运行信息
......
......@@ -15,10 +15,11 @@ import java.util.List;
public interface JobTaskRunLogMapper {
/**
* 查询没有结束的节点
* @param nodIds
* @param nodeIds
* @param runId
* @return
*/
int findJobTaskRunLogNotEndNodeByRunCodeCount(@Param("nodeIds") List<Integer> nodeIds);
List<JobTaskRunLog> findJobTaskRunLogNotEndNodeByRunCodeCount(@Param("nodeIds") List<Integer> nodeIds, @Param("runId") String runId);
/**
* 查根据flowId和RunId查询一批节点
......
......@@ -2,7 +2,6 @@ package com.byit.service;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskRunLogWithBLOBs;
import org.apache.ibatis.annotations.Param;
import java.util.List;
......@@ -17,9 +16,10 @@ public interface JobTaskRunLogService {
/**
* 查询没有结束的节点
* @param nodIds
* @param runId
* @return
*/
int findJobTaskRunLogNotEndNodeByRunCodeCount(List<Integer> nodIds);
List<JobTaskRunLog> findJobTaskRunLogNotEndNodeByRunCodeCount(List<Integer> nodIds, String runId);
/**
* 查根据flowId和RunId查询一批节点
* @param flowId
......
package com.byit.service;
import com.byit.model.JobTask;
import org.apache.ibatis.annotations.Param;
import com.byit.model.JobTaskSchedule;
import java.util.List;
......@@ -42,5 +42,5 @@ public interface JobTaskService {
* 根据集合删除
* @param jobTasks
*/
void removeMythJobTaskInIds(List<JobTask> jobTasks);
void removeMythJobTaskInIds(List<JobTaskSchedule> jobTasks);
}
......@@ -29,8 +29,8 @@ public class JobTaskRunLogServiceImpl implements JobTaskRunLogService {
}
@Override
public int findJobTaskRunLogNotEndNodeByRunCodeCount(List<Integer> nodIds) {
return jobTaskRunLogMapper.findJobTaskRunLogNotEndNodeByRunCodeCount(nodIds);
public List<JobTaskRunLog> findJobTaskRunLogNotEndNodeByRunCodeCount(List<Integer> nodIds, String runId) {
return jobTaskRunLogMapper.findJobTaskRunLogNotEndNodeByRunCodeCount(nodIds, runId);
}
@Override
......
......@@ -2,6 +2,7 @@ package com.byit.service.impl;
import com.byit.mapper.JobTaskMapper;
import com.byit.model.JobTask;
import com.byit.model.JobTaskSchedule;
import com.byit.service.JobTaskService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
......@@ -62,10 +63,10 @@ public class JobTaskServiceImpl implements JobTaskService {
/**
* 删除多个任务 根据id
* @param jobTasks 任务实体
* @param jobTaskSchedules 任务实体
*/
@Override
public void removeMythJobTaskInIds(List<JobTask> jobTasks) {
jobTaskMapper.deleteInId(jobTasks);
public void removeMythJobTaskInIds(List<JobTaskSchedule> jobTaskSchedules) {
jobTaskMapper.deleteInId(jobTaskSchedules);
}
}
......@@ -3,6 +3,7 @@ package com.byit.thread;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.job.WorkRoulette;
import com.byit.model.JobTask;
import com.byit.model.JobTaskRunLog;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.model.JobTaskSchedule;
import com.byit.service.JobTaskRunLogService;
......@@ -137,14 +138,25 @@ public class JobScheduleHelper{
//在日志表里面创建一条记录
//根据 map_flow_id查询当前的版本的工作流 使用祝工作流的runId 保存到执行记录表和任务表
}else{
if("start".equals(jobTask.getNodeName()) ){
log.debug("任务:{}", jobTask);
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
}else {
List<Integer> dependIdByNodeId = nodeDependencyService.findDependIdByNodeId(jobTask.getNodeId());
int jobTaskRunLogNotEndNodeByRunCodeCount = 0;
if (CollectionUtil.isNotEmpty(dependIdByNodeId)){
jobTaskRunLogNotEndNodeByRunCodeCount = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(dependIdByNodeId);
List<JobTaskRunLog> jobTaskRunLogList = jobTaskRunLogService.findJobTaskRunLogNotEndNodeByRunCodeCount(dependIdByNodeId, jobTask.getRunId());
if (CollectionUtil.isEmpty(jobTaskRunLogList)){
return;
}
if("start".equals(jobTask.getNodeName()) || (jobTaskRunLogNotEndNodeByRunCodeCount==0)){
log.debug("任务:{}", jobTask);
jobTaskRunLogList.forEach(e ->{
if(!("1".equals(e.getRunCode()) || "3".equals(e.getRunCode()))){
return;
}
});
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
......@@ -155,7 +167,7 @@ public class JobScheduleHelper{
if(CollectionUtil.isNotEmpty(jobTaskSchedules)) {
jobTaskScheduleService.saveAllData(jobTaskSchedules);
jobTaskService.removeMythJobTaskInIds(jobTasks);
jobTaskService.removeMythJobTaskInIds(jobTaskSchedules);
}
}else{
......@@ -397,10 +409,10 @@ public class JobScheduleHelper{
private Integer saveLog(JobTaskSchedule mythJobTaskSchedule){
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
//jobTaskRunLog.setRunId(mythJobTaskSchedule.getRunId());
jobTaskRunLog.setRunId(mythJobTaskSchedule.getRunId());
jobTaskRunLog.setFlowId(mythJobTaskSchedule.getFlowId());
jobTaskRunLog.setFlowName("1");
jobTaskRunLog.setRunId("qwer-tyui-opas-dfgh");
jobTaskRunLog.setNodeId(mythJobTaskSchedule.getNodeId());
jobTaskRunLog.setNodeName(mythJobTaskSchedule.getNodeName());
jobTaskRunLog.setRunParams(mythJobTaskSchedule.getRunParam());
jobTaskRunLog.setFailedRemainingCount(mythJobTaskSchedule.getFailedRetryCount());
......
......@@ -365,11 +365,11 @@
</update>
<!--多个删除-->
<delete id="deleteInId" parameterType="com.byit.model.JobTask">
<delete id="deleteInId" parameterType="com.byit.model.JobTaskSchedule">
delete from job_task
where id in
<foreach collection="jobTasks" item="jobTask" open="(" separator="," close=")">
#{jobTask.id}
<foreach collection="jobTaskSchedules" item="jobTaskSchedule" open="(" separator="," close=")">
#{jobTaskSchedule.id}
</foreach>
</delete>
<!--根据id删除-->
......
......@@ -40,10 +40,11 @@
run_msg, trigger_msg
</sql>
<select id="findJobTaskRunLogNotEndNodeByRunCodeCount" resultType="java.lang.Integer">
select count(*)
<select id="findJobTaskRunLogNotEndNodeByRunCodeCount" resultMap="BaseResultMap">
select <include refid="Base_Column_List" />
from job_task_run_log
where (run_code = '0' or run_code = '5')
where
run_id = #{runId}
and node_id in
<foreach item="nodeId" collection="nodeIds" open="(" separator="," close=")">
#{nodeId}
......
......@@ -3,38 +3,38 @@
<mapper namespace="com.byit.mapper.JobTaskScheduleMapper">
<resultMap id="BaseResultMap" type="com.byit.model.JobTaskSchedule">
<id column="id" jdbcType="INTEGER" property="id" />
<result column="node_id" jdbcType="INTEGER" property="nodeId" />
<result column="block_strategy" jdbcType="VARCHAR" property="blockStrategy" />
<result column="plugin_token" jdbcType="VARCHAR" property="pluginToken" />
<result column="failed_retry_count" jdbcType="INTEGER" property="failedRetryCount" />
<result column="flow_id" jdbcType="INTEGER" property="flowId" />
<result column="gateway_token" jdbcType="VARCHAR" property="gatewayToken" />
<result column="job_type" jdbcType="VARCHAR" property="jobType" />
<result column="handler_name" jdbcType="VARCHAR" property="handlerName" />
<result column="node_desc" jdbcType="VARCHAR" property="nodeDesc" />
<result column="node_name" jdbcType="VARCHAR" property="nodeName" />
<result column="map_flow_id" jdbcType="INTEGER" property="mapFlowId" />
<result column="node_timeout" jdbcType="BIGINT" property="nodeTimeout" />
<result column="is_virtual" jdbcType="CHAR" property="isVirtual" />
<result column="plugin_urls" jdbcType="VARCHAR" property="pluginUrls" />
<result column="priority" jdbcType="CHAR" property="priority" />
<result column="failed_retry_interval" jdbcType="BIGINT" property="failedRetryInterval" />
<result column="routing_strategy" jdbcType="VARCHAR" property="routingStrategy" />
<result column="run_id" jdbcType="VARCHAR" property="runId" />
<result column="run_param" jdbcType="VARCHAR" property="runParam" />
<result column="run_source_desc" jdbcType="VARCHAR" property="runSourceDesc" />
<result column="script_urls" jdbcType="VARCHAR" property="scriptUrls" />
<result column="source_principal" jdbcType="VARCHAR" property="sourcePrincipal" />
<result column="trigger_time" jdbcType="BIGINT" property="triggerTime" />
<result column="trigger_status" jdbcType="CHAR" property="triggerStatus" />
<result column="version_name" jdbcType="VARCHAR" property="versionName" />
<result column="log_id" jdbcType="INTEGER" property="logId" />
<result column="run_command" jdbcType="VARCHAR" property="runCommand" />
<id column="id" jdbcType="INTEGER" property="id"/>
<result column="node_id" jdbcType="INTEGER" property="nodeId"/>
<result column="block_strategy" jdbcType="VARCHAR" property="blockStrategy"/>
<result column="plugin_token" jdbcType="VARCHAR" property="pluginToken"/>
<result column="failed_retry_count" jdbcType="INTEGER" property="failedRetryCount"/>
<result column="flow_id" jdbcType="INTEGER" property="flowId"/>
<result column="gateway_token" jdbcType="VARCHAR" property="gatewayToken"/>
<result column="job_type" jdbcType="VARCHAR" property="jobType"/>
<result column="handler_name" jdbcType="VARCHAR" property="handlerName"/>
<result column="node_desc" jdbcType="VARCHAR" property="nodeDesc"/>
<result column="node_name" jdbcType="VARCHAR" property="nodeName"/>
<result column="map_flow_id" jdbcType="INTEGER" property="mapFlowId"/>
<result column="node_timeout" jdbcType="BIGINT" property="nodeTimeout"/>
<result column="is_virtual" jdbcType="CHAR" property="isVirtual"/>
<result column="plugin_urls" jdbcType="VARCHAR" property="pluginUrls"/>
<result column="priority" jdbcType="CHAR" property="priority"/>
<result column="failed_retry_interval" jdbcType="BIGINT" property="failedRetryInterval"/>
<result column="routing_strategy" jdbcType="VARCHAR" property="routingStrategy"/>
<result column="run_id" jdbcType="VARCHAR" property="runId"/>
<result column="run_param" jdbcType="VARCHAR" property="runParam"/>
<result column="run_source_desc" jdbcType="VARCHAR" property="runSourceDesc"/>
<result column="script_urls" jdbcType="VARCHAR" property="scriptUrls"/>
<result column="source_principal" jdbcType="VARCHAR" property="sourcePrincipal"/>
<result column="trigger_time" jdbcType="BIGINT" property="triggerTime"/>
<result column="trigger_status" jdbcType="CHAR" property="triggerStatus"/>
<result column="version_name" jdbcType="VARCHAR" property="versionName"/>
<result column="log_id" jdbcType="INTEGER" property="logId"/>
<result column="run_command" jdbcType="VARCHAR" property="runCommand"/>
</resultMap>
<resultMap extends="BaseResultMap" id="ResultMapWithBLOBs" type="com.byit.model.JobTaskSchedule">
<result column="run_source" jdbcType="LONGVARCHAR" property="runSource" />
<result column="run_source" jdbcType="LONGVARCHAR" property="runSource"/>
</resultMap>
<sql id="Base_Column_List">
......@@ -51,18 +51,18 @@
<!--查询五秒内将要执行的数据-->
<select id="findJobTaskScheduleByTriggerTimeLessThanEqual" resultMap="ResultMapWithBLOBs">
select
<include refid="Base_Column_List" />
<include refid="Base_Column_List"/>
,
<include refid="Blob_Column_List" />
<include refid="Blob_Column_List"/>
from job_task_schedule
where trigger_time <![CDATA[ <= ]]> #{triggerTime,jdbcType=BIGINT}
</select>
<select id="findJobTaskScheduleById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
select
<include refid="Base_Column_List" />
<include refid="Base_Column_List"/>
,
<include refid="Blob_Column_List" />
<include refid="Blob_Column_List"/>
from job_task_schedule
where id = #{id,jdbcType=INTEGER}
</select>
......@@ -267,16 +267,26 @@
run_command, run_source
)
values
<foreach collection="jobTaskSchedules" item="jobTaskSchedule" separator =",">
(#{jobTaskSchedule.id,jdbcType=INTEGER}, #{jobTaskSchedule.nodeId,jdbcType=INTEGER}, #{jobTaskSchedule.blockStrategy,jdbcType=VARCHAR},
#{jobTaskSchedule.pluginToken,jdbcType=VARCHAR}, #{jobTaskSchedule.failedRetryCount,jdbcType=INTEGER}, #{jobTaskSchedule.flowId,jdbcType=INTEGER},
#{jobTaskSchedule.gatewayToken,jdbcType=VARCHAR}, #{jobTaskSchedule.jobType,jdbcType=VARCHAR}, #{jobTaskSchedule.handlerName,jdbcType=VARCHAR},
#{jobTaskSchedule.nodeDesc,jdbcType=VARCHAR}, #{jobTaskSchedule.nodeName,jdbcType=VARCHAR}, #{jobTaskSchedule.mapFlowId,jdbcType=INTEGER},
#{jobTaskSchedule.nodeTimeout,jdbcType=BIGINT}, #{jobTaskSchedule.isVirtual,jdbcType=CHAR}, #{jobTaskSchedule.pluginUrls,jdbcType=VARCHAR},
#{jobTaskSchedule.priority,jdbcType=CHAR}, #{jobTaskSchedule.failedRetryInterval,jdbcType=BIGINT}, #{jobTaskSchedule.routingStrategy,jdbcType=VARCHAR},
#{jobTaskSchedule.runId,jdbcType=VARCHAR}, #{jobTaskSchedule.runParam,jdbcType=VARCHAR}, #{jobTaskSchedule.runSourceDesc,jdbcType=VARCHAR},
#{jobTaskSchedule.scriptUrls,jdbcType=VARCHAR}, #{jobTaskSchedule.sourcePrincipal,jdbcType=VARCHAR}, #{jobTaskSchedule.triggerTime,jdbcType=BIGINT},
#{jobTaskSchedule.triggerStatus,jdbcType=CHAR}, #{jobTaskSchedule.versionName,jdbcType=VARCHAR}, #{jobTaskSchedule.logId,jdbcType=INTEGER},
<foreach collection="jobTaskSchedules" item="jobTaskSchedule" separator=",">
(
#{jobTaskSchedule.id,jdbcType=INTEGER}, #{jobTaskSchedule.nodeId,jdbcType=INTEGER},
#{jobTaskSchedule.blockStrategy,jdbcType=VARCHAR},
#{jobTaskSchedule.pluginToken,jdbcType=VARCHAR}, #{jobTaskSchedule.failedRetryCount,jdbcType=INTEGER},
#{jobTaskSchedule.flowId,jdbcType=INTEGER},
#{jobTaskSchedule.gatewayToken,jdbcType=VARCHAR}, #{jobTaskSchedule.jobType,jdbcType=VARCHAR},
#{jobTaskSchedule.handlerName,jdbcType=VARCHAR},
#{jobTaskSchedule.nodeDesc,jdbcType=VARCHAR}, #{jobTaskSchedule.nodeName,jdbcType=VARCHAR},
#{jobTaskSchedule.mapFlowId,jdbcType=INTEGER},
#{jobTaskSchedule.nodeTimeout,jdbcType=BIGINT}, #{jobTaskSchedule.isVirtual,jdbcType=CHAR},
#{jobTaskSchedule.pluginUrls,jdbcType=VARCHAR},
#{jobTaskSchedule.priority,jdbcType=CHAR}, #{jobTaskSchedule.failedRetryInterval,jdbcType=BIGINT},
#{jobTaskSchedule.routingStrategy,jdbcType=VARCHAR},
#{jobTaskSchedule.runId,jdbcType=VARCHAR}, #{jobTaskSchedule.runParam,jdbcType=VARCHAR},
#{jobTaskSchedule.runSourceDesc,jdbcType=VARCHAR},
#{jobTaskSchedule.scriptUrls,jdbcType=VARCHAR}, #{jobTaskSchedule.sourcePrincipal,jdbcType=VARCHAR},
#{jobTaskSchedule.triggerTime,jdbcType=BIGINT},
#{jobTaskSchedule.triggerStatus,jdbcType=CHAR}, #{jobTaskSchedule.versionName,jdbcType=VARCHAR},
#{jobTaskSchedule.logId,jdbcType=INTEGER},
#{jobTaskSchedule.runCommand,jdbcType=VARCHAR}, #{jobTaskSchedule.runSource,jdbcType=LONGVARCHAR}
)
</foreach>
......
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