Commit 93a3cfea by huangfusuper

开发等待工作流扫描任务

parent fc2bff87
package com.byit.mapper; package com.byit.mapper;
import com.byit.model.WaitingRecord; import com.byit.model.WaitingRecord;
import org.springframework.stereotype.Repository;
import java.util.List;
/**
* @author huangfu
*/
@Repository
public interface WaitingRecordMapper { public interface WaitingRecordMapper {
/**
* 查询全部的工作流,同时根据执行时间查询对应的工作流数据
* @param triggerTime 本次的执行时间
* @return 即将执行的数据
*/
List<WaitingRecord> findAllByTriggerTime(long triggerTime);
int deleteById(Integer waitId); int deleteById(Integer waitId);
int insertSelective(WaitingRecord record); int insertSelective(WaitingRecord record);
......
package com.byit.mapper; package com.byit.mapper;
import com.byit.model.WaitingTask; import com.byit.model.WaitingTask;
import org.springframework.stereotype.Repository;
import java.util.List;
@Repository
public interface WaitingTaskMapper { public interface WaitingTaskMapper {
/**
* 根据等待实例查询等待节点
* @param waitId
* @return
*/
List<WaitingTask> findAllByWaitId(Integer waitId);
int deleteById(Integer id); int deleteById(Integer id);
int insertSelective(WaitingTask record); int insertSelective(WaitingTask record);
......
...@@ -3,13 +3,18 @@ package com.byit.model; ...@@ -3,13 +3,18 @@ package com.byit.model;
import io.swagger.annotations.ApiModel; import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty; import io.swagger.annotations.ApiModelProperty;
import java.io.Serializable; import java.io.Serializable;
import lombok.Data;
import lombok.*;
/** /**
* *
*/ */
@ApiModel @ApiModel
@Data @Data
@AllArgsConstructor
@NoArgsConstructor
@ToString
@Builder
public class WaitingRecord implements Serializable { public class WaitingRecord implements Serializable {
/** /**
* 运行标识 * 运行标识
......
package com.byit.service;
import com.byit.model.WaitingRecord;
import java.util.List;
/**
* 等待工作流的一个排序
* @author huangfu
*/
public interface WaitingRecordService {
/**
* 筛选即将执行的数据
* @param triggerTime
* @return
*/
List<WaitingRecord> findAllByTriggerTime(long triggerTime);
/**
* 修改数据
* @param waitingRecord
*/
void updateById(WaitingRecord waitingRecord);
}
package com.byit.service;
import com.byit.model.WaitingTask;
import java.util.List;
/**
* @author huangfu
*/
public interface WaitingTaskService {
/**
* 根据等待实例查询等待节点
* @param waitId
* @return
*/
List<WaitingTask> findAllByWaitId(Integer waitId);
}
package com.byit.service.impl;
import com.byit.mapper.WaitingRecordMapper;
import com.byit.model.WaitingRecord;
import com.byit.service.WaitingRecordService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.List;
/**
* 查询即将执行的数据
* @author huangfu
*/
@Service
@Slf4j
@Transactional(rollbackFor = Exception.class)
public class WaitingRecordServiceImpl implements WaitingRecordService {
private final WaitingRecordMapper waitingRecordMapper;
@Autowired
public WaitingRecordServiceImpl(WaitingRecordMapper waitingRecordMapper) {
this.waitingRecordMapper = waitingRecordMapper;
}
/**
* 查询一定时间内即将执行的数据
* @param triggerTime
* @return
*/
@Override
public List<WaitingRecord> findAllByTriggerTime(long triggerTime) {
log.info("------查询对应的等待工作流--------");
return waitingRecordMapper.findAllByTriggerTime(triggerTime);
}
@Override
public void updateById(WaitingRecord waitingRecord) {
waitingRecordMapper.updateByIdSelective(waitingRecord);
}
}
package com.byit.service.impl;
import com.byit.mapper.WaitingTaskMapper;
import com.byit.model.WaitingTask;
import com.byit.service.WaitingTaskService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.List;
/**
* 等待任务节点的实现
* @author huangfu
*/
@Service
@Transactional(rollbackFor = Exception.class)
@Slf4j
public class WaitingTaskServiceImpl implements WaitingTaskService {
private final WaitingTaskMapper waitingTaskMapper;
public WaitingTaskServiceImpl(WaitingTaskMapper waitingTaskMapper) {
this.waitingTaskMapper = waitingTaskMapper;
}
@Override
public List<WaitingTask> findAllByWaitId(Integer waitId) {
return waitingTaskMapper.findAllByWaitId(waitId);
}
}
package com.byit.service.mapservice; package com.byit.service.mapservice;
import com.byit.model.JobTask; import com.byit.model.JobTask;
import com.byit.model.WaitingRecord;
import java.net.UnknownHostException; import java.net.UnknownHostException;
...@@ -16,4 +17,12 @@ public interface RunRecordingAndJobTaskService { ...@@ -16,4 +17,12 @@ public interface RunRecordingAndJobTaskService {
* @throws Exception * @throws Exception
*/ */
void saveRunRecordingAndTask(JobTask jobTask) throws Exception; void saveRunRecordingAndTask(JobTask jobTask) throws Exception;
/**
* 修改运行实例 保存任务节点
* @param waitingRecord
*/
void updateRunRecordingAndTask(WaitingRecord waitingRecord);
} }
...@@ -24,20 +24,26 @@ import java.util.stream.Collectors; ...@@ -24,20 +24,26 @@ import java.util.stream.Collectors;
*/ */
@Service @Service
@Slf4j @Slf4j
@Transactional @Transactional(rollbackFor = Exception.class)
public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTaskService { public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTaskService {
private final JobTaskRunLogService jobTaskRunLogService; private final JobTaskRunLogService jobTaskRunLogService;
private final NodeService nodeService; private final NodeService nodeService;
private final FlowService flowService; private final FlowService flowService;
private final RunRecordingService runRecordingService; private final RunRecordingService runRecordingService;
private final JobTaskService jobTaskService; private final JobTaskService jobTaskService;
private final WaitingRecordService waitingRecordService;
private final WaitingTaskService taskService;
public RunRecordingAndJobTaskServiceImpl(JobTaskRunLogService jobTaskRunLogService, NodeService nodeService, FlowService flowService, RunRecordingService runRecordingService, JobTaskService jobTaskService) { public RunRecordingAndJobTaskServiceImpl(JobTaskRunLogService jobTaskRunLogService, NodeService nodeService,
FlowService flowService, RunRecordingService runRecordingService,
JobTaskService jobTaskService, WaitingRecordService waitingRecordService, WaitingTaskService taskService) {
this.jobTaskRunLogService = jobTaskRunLogService; this.jobTaskRunLogService = jobTaskRunLogService;
this.nodeService = nodeService; this.nodeService = nodeService;
this.flowService = flowService; this.flowService = flowService;
this.runRecordingService = runRecordingService; this.runRecordingService = runRecordingService;
this.jobTaskService = jobTaskService; this.jobTaskService = jobTaskService;
this.waitingRecordService = waitingRecordService;
this.taskService = taskService;
} }
/** /**
* 保存节点日志 * 保存节点日志
...@@ -58,7 +64,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask ...@@ -58,7 +64,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
//查看是否跟随工作流 //查看是否跟随工作流
boolean equals = "1".equals(mainFlow.getScheduleFollow()); boolean equals = "1".equals(mainFlow.getScheduleFollow());
boolean virFlag = "1".equals(virFlow.getScheduleFollow()); boolean virFlag = "1".equals(virFlow.getScheduleFollow());
//在日志表里面创建一条记录 //在日志表里面创建一条记录 虚节点在任务处理室是不会被保存的 所以需要在拉取虚节点的时候保存该节点数据
log.info("-----------【虚节点保存日志服务】--------------"); log.info("-----------【虚节点保存日志服务】--------------");
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs(); JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
jobTaskRunLog.setRunId(jobTask.getRunId()); jobTaskRunLog.setRunId(jobTask.getRunId());
...@@ -114,6 +120,32 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask ...@@ -114,6 +120,32 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
return task; return task;
}).collect(Collectors.toList()); }).collect(Collectors.toList());
jobTaskService.saveJobTasks(jobTasks); jobTaskService.saveJobTasks(jobTasks);
log.info("----------------【虚节点保存成功,删除虚节点】--------------------");
jobTaskService.removeMythJobTaskById(jobTask.getId());
log.info("-----------saveRunRecordingAndTask end【虚节点保存服务】--------------"); log.info("-----------saveRunRecordingAndTask end【虚节点保存服务】--------------");
}
@Override
public void updateRunRecordingAndTask(WaitingRecord waitingRecord) {
log.info("---------开始查询等待工作流{}对应的数据-------------",waitingRecord);
RunRecording runRecording = runRecordingService.findAllByRunID(waitingRecord.getRunId());
//将运行实例改为以执行
runRecording.setFlowStatus(RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode());
runRecordingService.updateRunRecordingById(runRecording);
log.info("-------修改运行实例表成功,查询对应等待实例{},的等待节点-------",waitingRecord);
List<WaitingTask> allByWaitId = taskService.findAllByWaitId(waitingRecord.getWaitId());
log.info("-------查询等待节点成功,查询对应的等待节点成功,开始保存对应的等待节点{}-------",allByWaitId);
List<JobTask> jobTasks = allByWaitId.stream().map(waitingTask -> {
JobTask target = new JobTask();
BeanUtils.copyProperties(waitingTask, target);
return target;
}).collect(Collectors.toList());
jobTaskService.saveJobTasks(jobTasks);
//保存等待实例
waitingRecord.setWaitOrder(-1);
log.info("-------保存jobTask成功,开始修改等待实例{}-----------------",waitingRecord);
waitingRecordService.updateById(waitingRecord);
} }
} }
...@@ -26,6 +26,15 @@ ...@@ -26,6 +26,15 @@
priority, trigger_time, principal, is_inner, flow_node_count, schedule_type, `operator`, priority, trigger_time, principal, is_inner, flow_node_count, schedule_type, `operator`,
wait_order, run_id wait_order, run_id
</sql> </sql>
<select id="findAllByTriggerTime" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
from waiting_record
where trigger_time<![CDATA[ <= ]]> #{triggerNextTime,jdbcType=BIGINT}
and wait_order > 0
</select>
<select id="getById" parameterType="java.lang.Integer" resultMap="BaseResultMap"> <select id="getById" parameterType="java.lang.Integer" resultMap="BaseResultMap">
<!-- generated @mbg.generated date: 2020-03-12 --> <!-- generated @mbg.generated date: 2020-03-12 -->
select select
......
...@@ -61,7 +61,17 @@ ...@@ -61,7 +61,17 @@
from waiting_task from waiting_task
where id = #{id,jdbcType=INTEGER} where id = #{id,jdbcType=INTEGER}
</select> </select>
<delete id="deleteById" parameterType="java.lang.Integer">
<select id="findAllByWaitId" resultType="com.byit.model.WaitingTask">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from waiting_task
where wait_id=#{waitId,jdbcType=INTEGER}
</select>
<delete id="deleteById" parameterType="java.lang.Integer">
<!-- generated @mbg.generated date: 2020-03-12 --> <!-- generated @mbg.generated date: 2020-03-12 -->
delete from waiting_task delete from waiting_task
where id = #{id,jdbcType=INTEGER} where id = #{id,jdbcType=INTEGER}
......
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