Commit 8bb1ef64 by guominglei

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

parents cb802f9a f7fa9088
*.class
*.iml
compiler.xml
encodings.xml
Maven_*.xml
misc.xml
modules.xml
Project_Default.xml
vcs.xml
workspace.xml
.idea
target
\ No newline at end of file
package com.byit.conf;
import com.byit.thread.EmailScanHelper;
import com.byit.thread.JobScheduleHelper;
import com.byit.thread.LogScanHelper;
import com.byit.thread.RunRecordingScanHelper;
import com.byit.thread.*;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.beans.factory.InitializingBean;
......@@ -23,13 +20,15 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
private final LogScanHelper logScanHelper;
private final RunRecordingScanHelper runRecordingScanHelper;
private final EmailScanHelper emailScanHelper;
private final FlowScanHelper flowScanHelper;
@Autowired
public MythJobScheduler(JobScheduleHelper jobScheduleHelper, LogScanHelper logScanHelper, RunRecordingScanHelper runRecordingScanHelper, EmailScanHelper emailScanHelper) {
public MythJobScheduler(JobScheduleHelper jobScheduleHelper, LogScanHelper logScanHelper, RunRecordingScanHelper runRecordingScanHelper, EmailScanHelper emailScanHelper, FlowScanHelper flowScanHelper) {
this.jobScheduleHelper = jobScheduleHelper;
this.logScanHelper = logScanHelper;
this.runRecordingScanHelper = runRecordingScanHelper;
this.emailScanHelper = emailScanHelper;
this.flowScanHelper = flowScanHelper;
}
/**
......@@ -42,6 +41,7 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
this.logScanHelper.doStop();
this.runRecordingScanHelper.doStop();
this.emailScanHelper.doStop();
this.flowScanHelper.doStop();
}
/**
......@@ -55,5 +55,6 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
this.logScanHelper.start();
this.runRecordingScanHelper.start();
this.emailScanHelper.start();
this.flowScanHelper.start();
}
}
......@@ -3,7 +3,16 @@ package com.byit.mapper;
import com.byit.model.Flow;
import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface FlowMapper {
/**
* 查询半个小时内即将要执行的工作流
* @param triggerNextTime
* @return
*/
List<Flow> findHalfAnHourFlow(@Param("triggerNextTime") Long triggerNextTime);
int deleteById(Integer flowId);
int insertSelective(Flow record);
......
......@@ -40,6 +40,13 @@ public interface JobTaskMapper {
int saveJobTask(JobTask jobTask);
/**
* 批量保存
* @param jobTasks
* @return
*/
int saveJobTasks(@Param("jobTasks") List<JobTask> jobTasks);
/**
* 修改任务表
* @param jobTask
* @return
......
......@@ -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
......
......@@ -5,7 +5,18 @@ import org.apache.ibatis.annotations.Param;
import java.util.List;
/**
* @author HUANGFU
*/
public interface NodeMapper {
/**
* 跟怒工作流ID和版本名称查询所有的节点
* @param flowId
* @param versionName
* @return
*/
List<Node> findNodeByFlowIdAndVersionName(@Param("flowId") Integer flowId,@Param("versionName") String versionName);
int deleteById(Integer nodeId);
int insertSelective(Node record);
......
......@@ -3,13 +3,20 @@ package com.byit.model;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import java.io.Serializable;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
*
*/
@ApiModel
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class JobTask implements Serializable {
/**
*/
......@@ -119,7 +126,7 @@ public class JobTask implements Serializable {
private String routingStrategy;
/**
* 运行标识
* 运行标识-----
*/
@ApiModelProperty("运行标识")
private String runId;
......
......@@ -52,7 +52,7 @@ public class RunRecording implements Serializable {
/**
* 执行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill
*/
@ApiModelProperty("执行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill")
@ApiModelProperty("执行结果 0未执行 1 成功 2 失败 3 补批成功 4 补批失败 5.kill")
private String flowRunResult;
/**
......
......@@ -3,6 +3,9 @@ package com.byit.service;
import com.byit.model.Flow;
import com.byit.model.vo.FlowVo;
import com.byit.model.vo.NodeVo;
import org.apache.ibatis.annotations.Param;
import java.util.List;
/**
* @description: 工作流业务逻辑接口
......@@ -10,6 +13,13 @@ import com.byit.model.vo.NodeVo;
* @create: 2019-12-23 17:33
*/
public interface FlowService {
/**
* 查询半个小时内即将要执行的工作流
* @param triggerNextTime
* @return
*/
List<Flow> findHalfAnHourFlow(Long triggerNextTime);
/**
* 保存工作流信息
......@@ -30,6 +40,8 @@ public interface FlowService {
*/
void updateFlow(FlowVo flowVo);
int updateByIdSelective(Flow record);
/**
* 真实删除当前表工作流信息
* @param flowId
......
......@@ -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 com.byit.model.JobTask;
import org.apache.ibatis.annotations.Param;
import java.util.List;
......@@ -25,6 +26,13 @@ public interface JobTaskService {
void addMythJobTask(JobTask jobTask);
/**
* 批量保存
* @param jobTasks
* @return
*/
int saveJobTasks(List<JobTask> jobTasks);
/**
* 删除已经被调度的任务根据ID
* @param id
*/
......
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
package com.byit.service;
import com.byit.model.Node;
import org.apache.ibatis.annotations.Param;
import java.util.List;
/**
* @description: 任务节点逻辑处理接口
* @author: gml
* @create: 2019-12-24 10:38
*/
public interface NodeService {
/**
* 跟怒工作流ID和版本名称查询所有的节点
* @param flowId
* @param versionName
* @return
*/
List<Node> findNodeByFlowIdAndVersionName(Integer flowId,String versionName);
}
package com.byit.service;
import com.byit.model.Flow;
import com.byit.model.Node;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;
import java.util.List;
/**
*封装事务
* @author huangfu
*/
public interface RunNodeServer {
/**
* 保存执行记录task
* @param flow
* @param nodes
*/
void saveRunRecAndTask(Flow flow, List<Node> nodes);
}
......@@ -44,6 +44,11 @@ public class FlowServiceImpl implements FlowService {
private NodeDependencyMapper nodeDependencyMapper;
@Override
public List<Flow> findHalfAnHourFlow(Long triggerNextTime) {
return flowMapper.findHalfAnHourFlow(triggerNextTime);
}
@Override
public Flow saveJobFlow(FlowVo flowVo) {
ValidationUtil.dataNotNull(flowVo.getFlowId(), "工作流Id不允许为空!");
......@@ -197,6 +202,11 @@ public class FlowServiceImpl implements FlowService {
}
@Override
public int updateByIdSelective(Flow record) {
return flowMapper.updateByIdSelective(record);
}
@Override
public void deleteFlow(Integer flowId) {
}
......
......@@ -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);
}
......
......@@ -46,6 +46,11 @@ public class JobTaskServiceImpl implements JobTaskService {
jobTaskMapper.saveJobTask(jobTask);
}
@Override
public int saveJobTasks(List<JobTask> jobTasks) {
return jobTaskMapper.saveJobTasks(jobTasks);
}
/**
* 根据id删除一个任务
* @param id 任务的id
......
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);
}
}
package com.byit.service.impl;
import com.byit.mapper.NodeMapper;
import com.byit.model.Node;
import com.byit.service.NodeService;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.util.List;
/**
* @description: 任务节点逻辑处理实现类
......@@ -17,4 +19,8 @@ public class NodeServiceImpl implements NodeService {
@Resource
private NodeMapper nodeMapper;
@Override
public List<Node> findNodeByFlowIdAndVersionName(Integer flowId, String versionName) {
return nodeMapper.findNodeByFlowIdAndVersionName(flowId,versionName);
}
}
package com.byit.service.impl;
import com.byit.model.Flow;
import com.byit.model.JobTask;
import com.byit.model.Node;
import com.byit.model.RunRecording;
import com.byit.service.JobTaskService;
import com.byit.service.RunNodeServer;
import com.byit.service.RunRecordingService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
/**
* @author huangfu
*/
@Service
@Transactional(propagation = Propagation.REQUIRED,rollbackFor = Exception.class)
@Slf4j
public class RunNodeServiceImpl implements RunNodeServer {
private final RunRecordingService runRecordingService;
private final JobTaskService jobTaskService;
public RunNodeServiceImpl(RunRecordingService runRecordingService, JobTaskService jobTaskService) {
this.runRecordingService = runRecordingService;
this.jobTaskService = jobTaskService;
}
@Override
public void saveRunRecAndTask(Flow flow, List<Node> nodes) {
log.info("---------saveRunRecAndTask start------【保存工作流:{}和节点:{}】-----------------------",flow,nodes);
String runId = UUID.randomUUID().toString().replace("-","");
log.info("-------------【开始保存运行记录runId为:{}】------------------",runId);
RunRecording build = new RunRecording();
BeanUtils.copyProperties(flow,build);
build.setRunId(runId);
build.setDispatchIp("127.0.0.1");
build.setFlowVersionName(flow.getVersionName());
build.setTriggerTime(flow.getTriggerNextTime());
runRecordingService.saveRunRecording(build);
boolean isScheduleFollow = "1".equals(flow.getScheduleFollow());
log.info("-------------【开始保存节点信息,任务是否为跟随工作流,{}】---------------",isScheduleFollow);
List<JobTask> jobTasks = new ArrayList<>(32);
nodes.forEach(node -> {
JobTask jobTask = new JobTask();
BeanUtils.copyProperties(node,jobTask);
if(isScheduleFollow){
jobTask.setTriggerTime(flow.getTriggerNextTime());
}
jobTask.setRunId(runId);
jobTasks.add(jobTask);
});
jobTaskService.saveJobTasks(jobTasks);
log.info("-------saveRunRecAndTask end-----------【运行结束】-----------------");
}
}
......@@ -46,6 +46,8 @@ public class EmailScanHelper {
dateAligned(5000);
log.info("--------------【com.byit.thread.EmailScanHelper.startScanNotSentEmailFlow】init success---------------");
while (!emailThreadIsStop){
//是否需要睡眠
boolean isSleep = false;
Connection conn = null;
Boolean connAutoCommit = null;
PreparedStatement preparedStatement = null;
......@@ -67,7 +69,7 @@ public class EmailScanHelper {
emailAlarmService.sendEmail(emailAlarmVo);
});
}else{
dateAligned(20000);
isSleep = true;
}
}catch (Exception e){
e.printStackTrace();
......@@ -111,6 +113,10 @@ public class EmailScanHelper {
}
}
}
if(isSleep){
dateAligned(20000);
}
}
});
emailThread.setDaemon(true);
......
package com.byit.thread;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.model.Flow;
import com.byit.model.JobTask;
import com.byit.model.Node;
import com.byit.model.RunRecording;
import com.byit.service.*;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Component;
import javax.sql.DataSource;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
/**
* 扫描任务流表的线程
* 目的:扫描即将要执行的任务流表 半个小时
* @author huangfu
*/
@Component
@Slf4j
public class FlowScanHelper {
private final Long PRE_TEST_TIME = System.currentTimeMillis()+TimeUnit.HOURS.toMillis(30);
private final FlowService flowService;
private final DataSource dataSource;
private final NodeService nodeService;
private final JobTaskService jobTaskService;
private final RunRecordingService runRecordingService;
private final RunNodeServer runNodeServer;
private Thread flowThread;
private volatile boolean flowThreadIsStop = false;
public FlowScanHelper(FlowService flowService, DataSource dataSource, NodeService nodeService, JobTaskService jobTaskService, RunRecordingService runRecordingService, RunNodeServer runNodeServer) {
this.flowService = flowService;
this.dataSource = dataSource;
this.nodeService = nodeService;
this.jobTaskService = jobTaskService;
this.runRecordingService = runRecordingService;
this.runNodeServer = runNodeServer;
}
public void start(){
flowThreadStart();
}
/**
* 运行扫描工作流的线程
*/
private void flowThreadStart(){
flowThread = new Thread(() ->{
dateAligned(5000);
log.info("--------------【com.byit.thread.FlowScanHelper#flowThreadStart】init success---------------");
while (!flowThreadIsStop){
Connection conn = null;
Boolean connAutoCommit = null;
PreparedStatement preparedStatement = null;
//是否需要睡眠
boolean isSleep = false;
try {
/**
* 添加行锁
*/
conn = dataSource.getConnection();
connAutoCommit = conn.getAutoCommit();
conn.setAutoCommit(false);
preparedStatement = conn.prepareStatement("SELECT * FROM JOB_LOCK WHERE LOCK_NAME = 'flow_lock' FOR UPDATE ");
preparedStatement.execute();
List<Flow> halfAnHourFlow = flowService.findHalfAnHourFlow(PRE_TEST_TIME);
if (CollectionUtil.isNotEmpty(halfAnHourFlow)) {
halfAnHourFlow.forEach(flow -> {
log.debug("-----------------【工作流{}的执行次数大于0,放行】-------------------------",flow.getFlowName());
String versionName = flow.getVersionName();
Integer flowId = flow.getFlowId();
List<Node> nodeByFlowIdAndVersionName = nodeService.findNodeByFlowIdAndVersionName(flowId, versionName);
runNodeServer.saveRunRecAndTask(flow,nodeByFlowIdAndVersionName);
flow.setRepeatCount(flow.getRepeatCount()-1);
flowService.updateByIdSelective(flow);
});
}else{
isSleep = true;
}
}catch (Exception e){
e.printStackTrace();
}finally {
//释放资源
if(conn != null){
try {
conn.commit();
} catch (SQLException e) {
if(!flowThreadIsStop){
log.error("--------------------【提交行锁出错】---------------------");
}
}
}
try {
if(conn != null){
conn.setAutoCommit(connAutoCommit);
}
} catch (SQLException e) {
if(!flowThreadIsStop){
log.error("--------------------【恢复自动提交出错】---------------------");
}
}
if(preparedStatement != null){
try {
preparedStatement.close();
} catch (SQLException e) {
if(!flowThreadIsStop){
log.error("--------------------【关闭执行器出错】---------------------");
}
}
}
try {
conn.close();
} catch (SQLException e) {
if(!flowThreadIsStop){
log.error("--------------------【关闭链接出错】---------------------");
}
}
}
if(isSleep){
try {
log.info("------------------【未扫描到要执行的工作流】-------------------");
TimeUnit.HOURS.sleep(10);
} catch (InterruptedException e) {
log.warn("----------------------【扫描工作流的线程被关闭了】----------------------------");
}
}
}
});
flowThread.setDaemon(true);
flowThread.setName("myth-job#【FlowScanHelper】#flowThreadStart");
flowThread.start();
}
public void doStop(){
this.flowThreadIsStop = true;
try {
TimeUnit.SECONDS.sleep(1);
} catch (InterruptedException e) {
e.printStackTrace( );
}
if (flowThread.getState() != Thread.State.TERMINATED) {
flowThread.interrupt();
try {
flowThread.join();
} catch (InterruptedException e) {
e.printStackTrace( );
}
}
log.warn("---------------【工作扫描线程被注销】-----------------------");
}
/**
* 对齐时钟。整秒运行
*/
private void dateAligned(long waitTime){
try {
TimeUnit.MILLISECONDS.sleep(waitTime - System.currentTimeMillis()%1000);
} catch (InterruptedException e) {
log.warn("----------------【线程被中断】-----------------------");
}
}
}
......@@ -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){
log.debug("任务:{}", jobTask);
JobTaskSchedule jobTaskSchedule = new JobTaskSchedule();
BeanUtils.copyProperties(jobTask,jobTaskSchedule);
jobTaskSchedules.add(jobTaskSchedule);
//虚节点的状态
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,15 +224,17 @@ public class JobScheduleHelper{
//设置名字
jobInfoThread.setName("myth-job,admin JobScheduleHelper#jobInfoThread");
jobInfoThread.start();
}
/**
* 开始操作任务排期表:
* 1. 扫描五秒要执行的数据
* 2.加载时,创建任务日志,传入调度时间
* 3.加载到任务调度轮盘
* 4.删除数据
*/
private void scheduleThreadStart(){
/**
* 开始操作任务排期表:
* 1. 扫描五秒要执行的数据
* 2.加载时,创建任务日志,传入调度时间
* 3.加载到任务调度轮盘
* 4.删除数据
*/
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 {
......
......@@ -2,6 +2,7 @@ package com.byit.util;
import com.byit.enums.DagCheckEnum;
import com.byit.job.dto.plugin.PluginBaseNode;
import lombok.extern.slf4j.Slf4j;
import java.util.*;
......@@ -10,6 +11,7 @@ import java.util.*;
* @author: gml
* @create: 2019-12-30 16:34
*/
@Slf4j
public class ApiFlowDagCheck {
//节点个数
......@@ -131,10 +133,10 @@ public class ApiFlowDagCheck {
}
}
if(number != nodeNum){
System.out.println("最后存在入度为1的结点,这个有向图是有回路的。");
log.debug("最后存在入度为1的结点,这个有向图是有回路的。");
return DagCheckEnum.LOOP;
}else{
System.out.println("这个有向图不存在回路,拓扑序列为:" + temp.toString());
log.debug("这个有向图不存在回路,拓扑序列为:{}", temp.toString());
return DagCheckEnum.PASS;
}
}
......
......@@ -33,6 +33,13 @@
author, add_time, start_up, principal, version_name, repeat_count, remaining_count,
schedule_follow, is_update
</sql>
<select id="findHalfAnHourFlow" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
from flow where trigger_next_time <![CDATA[ <= ]]> #{triggerNextTime,jdbcType=BIGINT} and remaining_count <![CDATA[ <> ]]> 0
</select>
<select id="getById" parameterType="java.lang.Integer" resultMap="BaseResultMap">
<!-- generated @mbg.generated date: 2019-12-31 -->
select
......
......@@ -58,7 +58,7 @@
</select>
<!--根据id查询-->
<select id="findJobTaskById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
select
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
......@@ -242,6 +242,29 @@
</if>
</trim>
</insert>
<insert id="saveJobTasks" parameterType="com.byit.model.JobTask">
insert into job_task (
id, node_id, block_strategy, plugin_token, failed_retry_count, flow_id, gateway_token,
job_type, handler_name, node_desc, node_name, map_flow_id, node_timeout, is_virtual,
plugin_urls, priority, failed_retry_interval, routing_strategy, run_id, run_param,
run_source_desc, script_urls, source_principal, trigger_time, trigger_status, version_name,
run_command,run_source
) values
<foreach collection="jobTasks" item="jobTask" separator =",">
(
#{jobTask.id,jdbcType=INTEGER}, #{jobTask.nodeId,jdbcType=INTEGER},#{jobTask.blockStrategy,jdbcType=VARCHAR},#{jobTask.pluginToken,jdbcType=VARCHAR},
#{jobTask.failedRetryCount,jdbcType=INTEGER},#{jobTask.flowId,jdbcType=INTEGER}, #{jobTask.gatewayToken,jdbcType=VARCHAR},#{jobTask.jobType,jdbcType=VARCHAR},
#{jobTask.handlerName,jdbcType=VARCHAR},#{jobTask.nodeDesc,jdbcType=VARCHAR},#{jobTask.nodeName,jdbcType=VARCHAR},#{jobTask.mapFlowId,jdbcType=INTEGER},
#{jobTask.nodeTimeout,jdbcType=BIGINT}, #{jobTask.isVirtual,jdbcType=CHAR},#{jobTask.pluginUrls,jdbcType=VARCHAR},#{jobTask.priority,jdbcType=CHAR},
#{jobTask.failedRetryInterval,jdbcType=BIGINT},#{jobTask.routingStrategy,jdbcType=VARCHAR},#{jobTask.runId,jdbcType=VARCHAR},
#{jobTask.runParam,jdbcType=VARCHAR},#{jobTask.runSourceDesc,jdbcType=VARCHAR},#{jobTask.scriptUrls,jdbcType=VARCHAR},
#{jobTask.sourcePrincipal,jdbcType=VARCHAR},#{jobTask.triggerTime,jdbcType=BIGINT},#{jobTask.triggerStatus,jdbcType=CHAR},
#{jobTask.versionName,jdbcType=VARCHAR},#{jobTask.runCommand,jdbcType=VARCHAR},#{jobTask.runSource,jdbcType=LONGVARCHAR}
)
</foreach>
</insert>
<update id="updateJobTask" parameterType="com.byit.model.JobTask">
<!-- generated @mbg.generated date: 2019-12-25 -->
update job_task
......
......@@ -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">
......
......@@ -47,10 +47,23 @@
routing_strategy, run_param, run_source_desc, script_urls, source_principal, source_update_time,
trigger_next_time, author, add_time, version_name, on_fork, run_command
</sql>
<sql id="Blob_Column_List">
<!-- generated @mbg.generated date: 2019-12-31 -->
run_source
</sql>
<select id="findNodeByFlowIdAndVersionName" resultMap="ResultMapWithBLOBs">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from node
where flow_id = #{flowId,jdbcType=INTEGER} AND version_name = #{versionName,jdbcType=VARCHAR}
</select>
<select id="getById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
<!-- generated @mbg.generated date: 2019-12-31 -->
select
......
......@@ -23,7 +23,7 @@ public class PluginBaseNode {
private String desc;
/**
* 类型 节点:node, 内嵌工作流:innerFlow
* 类型 节点:node, 工作流:flow
*/
private String type;
......@@ -33,6 +33,7 @@ public class PluginBaseNode {
private String author;
/**
* 上游
* 依赖的任务节点名称集合,只能是节点的名称,如果不是内嵌工作流不允许设置依赖
*/
private List<String> dependNodeNameList;
......
package com.byit.job.dto.plugin;
import lombok.Data;
import lombok.*;
import java.io.Serializable;
import java.util.List;
......@@ -11,6 +11,9 @@ import java.util.List;
* @create: 2019-12-30 10:44
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
@EqualsAndHashCode(callSuper = true)
public class PluginFlow extends PluginBaseNode implements Serializable {
/**
......
package com.byit.job.dto.plugin;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* @description: 插件端工作流的配置
......@@ -8,6 +11,9 @@ import lombok.Data;
* @create: 2019-12-30 11:16
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public class PluginFlowConfig {
/**
......
......@@ -14,7 +14,7 @@ import com.byit.job.vo.ReturnResult;
public class DemoJob extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String s) throws Exception {
System.out.println("--------------任务就这样运行了----------------"+s );
System.out.println("--------------任务就这样运行了 addJob----------------"+s );
return ReturnResult.SUCCESS;
}
}
......@@ -11,11 +11,10 @@ import com.byit.job.vo.ReturnResult;
* @date: 2019/11/20 12:40
**/
@JobHandler("END")
public class DemoJob1 extends BaseJobHandler {
public class EndNode extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String s) throws Exception {
Thread.sleep(10000);
System.out.println("--------------DemoJob1----------------"+s );
System.out.println("--------------END----------------"+s );
return ReturnResult.SUCCESS;
}
}
package com.byit.job;
import com.byit.annotations.JobHandler;
import com.byit.job.handler.BaseJobHandler;
import com.byit.job.vo.ReturnResult;
@JobHandler("start")
public class StartNode extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String param) throws Exception {
System.out.println("--------【start】----我摊牌了,我是亿万富翁------------");
return ReturnResult.SUCCESS;
}
}
package com.byit.job;
import com.byit.job.dto.plugin.*;
import com.byit.rpc.remoting.invoker.route.LoadBalance;
import com.byit.utils.JobUtils;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.TimeUnit;
public class TestAddFlow {
public static void main(String[] args) {
PluginPackage pluginPackage = new PluginPackage();
pluginPackage.setWorkspaceName("test");
pluginPackage.setFlow(createFlow());
JobUtils.publish(pluginPackage);
}
public static PluginFlow createFlow(){
PluginFlow pluginFlow = new PluginFlow();
PluginFlowConfig build = PluginFlowConfig.builder().alarmEmail("huangfukexin@byitgroup.com,guominglei@byitgroup.com")
.alarmlAction("1")
.execType("1")
.flowCron("0 0/5 * * * ? *")
.flowTimeout(TimeUnit.MINUTES.toMillis(30))
.priority("2")
.repeatCount(1)
.scheduleFollow("1")
.build();
pluginFlow.setName("测试任务流");
pluginFlow.setDesc("这是一个测试的任务流");
pluginFlow.setConfig(build);
pluginFlow.setPrincipal("皇甫科星");
pluginFlow.setRePublish(false);
pluginFlow.setAuthor("huangfusuper");
pluginFlow.setNodeList(createNodes());
return pluginFlow;
}
/**
* 创建节点
* @return
*/
public static List<PluginBaseNode> createNodes(){
PluginNode pluginNode1 = new PluginNode();
PluginNodeConfig pluginNodeConfig1 = new PluginNodeConfig();
pluginNode1.setName("start");
pluginNode1.setDesc("我是开始节点,打死你");
pluginNode1.setType("node");
pluginNode1.setAuthor("郭郭");
pluginNode1.setJobType("JAVA");
pluginNode1.setHandlerName("start");
pluginNodeConfig1.setFailedRetryCount(2);
pluginNodeConfig1.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig1.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig1.setNodeTimeout(-1L);
pluginNodeConfig1.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig1.setPriority("2");
pluginNodeConfig1.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode1.setConfig(pluginNodeConfig1);
PluginNode pluginNode2 = new PluginNode();
PluginNodeConfig pluginNodeConfig2 = new PluginNodeConfig();
pluginNode2.setName("中间节点1");
pluginNode2.setDesc("中间节点1");
pluginNode2.setType("node");
pluginNode2.setAuthor("皇甫");
pluginNode2.setJobType("JAVA");
pluginNode2.setHandlerName("addJob");
pluginNode2.setRunParam("addJob1");
pluginNodeConfig2.setFailedRetryCount(2);
pluginNodeConfig2.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig2.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig2.setNodeTimeout(-1L);
pluginNodeConfig2.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig2.setPriority("2");
pluginNodeConfig2.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode2.setConfig(pluginNodeConfig2);
pluginNode2.setDependNodeNameList(Collections.singletonList("start"));
PluginNode pluginNode3 = new PluginNode();
PluginNodeConfig pluginNodeConfig3 = new PluginNodeConfig();
pluginNode3.setName("中间节点2");
pluginNode3.setDesc("中间节点2");
pluginNode3.setType("node");
pluginNode3.setAuthor("皇甫");
pluginNode3.setJobType("JAVA");
pluginNode3.setHandlerName("addJob");
pluginNode3.setRunParam("addJob2");
pluginNodeConfig3.setFailedRetryCount(2);
pluginNodeConfig3.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig3.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig3.setNodeTimeout(-1L);
pluginNodeConfig3.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig3.setPriority("2");
pluginNodeConfig3.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode3.setConfig(pluginNodeConfig3);
pluginNode3.setDependNodeNameList(Collections.singletonList("中间节点1"));
PluginNode pluginNode4 = new PluginNode();
PluginNodeConfig pluginNodeConfig4 = new PluginNodeConfig();
pluginNode4.setName("中间节点3");
pluginNode4.setDesc("中间节点3");
pluginNode4.setType("node");
pluginNode4.setAuthor("皇甫");
pluginNode4.setJobType("JAVA");
pluginNode4.setHandlerName("addJob");
pluginNode4.setRunParam("addJob3");
pluginNodeConfig4.setFailedRetryCount(2);
pluginNodeConfig4.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig4.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig4.setNodeTimeout(-1L);
pluginNodeConfig4.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig4.setPriority("2");
pluginNodeConfig4.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode4.setConfig(pluginNodeConfig4);
pluginNode4.setDependNodeNameList(Collections.singletonList("中间节点2"));
PluginNode pluginNode5 = new PluginNode();
PluginNodeConfig pluginNodeConfig5 = new PluginNodeConfig();
pluginNode5.setName("end");
pluginNode5.setDesc("结束节点");
pluginNode5.setType("node");
pluginNode5.setAuthor("皇甫");
pluginNode5.setJobType("JAVA");
pluginNode5.setHandlerName("END");
pluginNode5.setRunParam("END");
pluginNodeConfig5.setFailedRetryCount(2);
pluginNodeConfig5.setFailedRetryInterval(TimeUnit.MINUTES.toSeconds(2));
pluginNodeConfig5.setNodeCron("0 0/7 * * * ? *");
pluginNodeConfig5.setNodeTimeout(-1L);
pluginNodeConfig5.setPluginUrls("http://127.0.0.1:8888");
pluginNodeConfig5.setPriority("2");
pluginNodeConfig5.setRoutingStrategy(LoadBalance.ROUND.name());
pluginNode5.setConfig(pluginNodeConfig5);
pluginNode5.setDependNodeNameList(Collections.singletonList("中间节点3"));
return Arrays.asList(pluginNode5, pluginNode4, pluginNode3, pluginNode2, pluginNode1);
}
}
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