Commit 2ed1569c by huangfusuper

重构补批工作流

parent 88676f99
package com.byit.api;
import com.byit.dto.specials.RepairFlow;
import com.byit.dto.specials.SpecialJobParam;
import com.byit.service.ApiFlowOperatingService;
import io.swagger.annotations.Api;
......@@ -34,4 +35,13 @@ public class ApiFlowOperatingController {
public void specialRunBatch(@RequestBody SpecialJobParam specialJobParam){
apiFlowOperatingService.specialRunBatch(specialJobParam);
}
/**
* 常规的补批
* @param repairFlow 常规的补批
*/
@PostMapping("routineSupplementBatch")
public void routineSupplementBatch(@RequestBody RepairFlow repairFlow) {
apiFlowOperatingService.routineSupplementBatch(repairFlow);
}
}
package com.byit.service;
import com.byit.dto.specials.RepairFlow;
import com.byit.dto.specials.SpecialJobParam;
/**
......@@ -16,4 +17,10 @@ public interface ApiFlowOperatingService {
* @param specialJobParam 补批参数
*/
void specialRunBatch(SpecialJobParam specialJobParam);
/**
* 常规工作流补批
* @param repairFlow 常规工作流补批
*/
void routineSupplementBatch(RepairFlow repairFlow);
}
......@@ -4,6 +4,7 @@ import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.fastjson.JSON;
import com.byit.dto.NodeRelyDto;
import com.byit.dto.executor.RunParamWrapped;
import com.byit.dto.specials.RepairFlow;
import com.byit.dto.specials.SpecialJobParam;
import com.byit.enums.NodeTypeEnum;
import com.byit.enums.PlaceholderEnum;
......@@ -24,6 +25,7 @@ import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.text.ParseException;
import java.text.SimpleDateFormat;
......@@ -37,10 +39,10 @@ import java.util.stream.Collectors;
*/
@Slf4j
@Service
@Transactional(rollbackFor = Exception.class)
public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
@Value("${server.port}")
private Integer serverPort;
......@@ -134,7 +136,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
List<WaitingRecord> waitingRecords = buildWaitingRecordList(flowByName, runId, nodeRelyDtoSet.size(), operator, repairTimeList);
waitingRecords.forEach(waitingRecordMapper::insertSelective);
List<RunRecording> runRecordings = buildRunRecordingList(waitingRecords);
List<WaitingTask> waitingTaskList = buildWaitingTaskResult(waitingRecords, nodeRelyDtoSet,specialJobParam.getPublicParam(),timeFormatName, nowTimeFormatName);
List<WaitingTask> waitingTaskList = buildWaitingTaskResult(waitingRecords, nodeRelyDtoSet, specialJobParam.getPublicParam(), timeFormatName, nowTimeFormatName);
runRecordings.forEach(runRecordingMapper::saveRunRecording);
waitingTaskList.forEach(waitingTaskMapper::insertSelective);
}
......@@ -142,6 +144,18 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
}
@Override
public void routineSupplementBatch(RepairFlow repairFlow) {
ValidationUtil.dataNotNull(repairFlow, "补批核心参数不能为空!");
if ("1".equals(repairFlow.getAllNodes())) {
log.info("-------补批工作流-------");
SpecialJobParam specialJobParam = new SpecialJobParam();
BeanUtils.copyProperties(repairFlow, specialJobParam);
specialJobParam.setStartNodeTaskName("start");
specialRunBatch(specialJobParam);
}
}
/**
* 构建单个等待队列的节点
*
......@@ -149,7 +163,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
* @param nodeRelyDtoSet 节点对象
* @return 返回一个等待队的全部节点
*/
private List<WaitingTask> buildWaitingTask(WaitingRecord waitingRecord, Set<NodeRelyDto> nodeRelyDtoSet,Map<String,String> publicParam, String timeFormat, String nowDateFormat) {
private List<WaitingTask> buildWaitingTask(WaitingRecord waitingRecord, Set<NodeRelyDto> nodeRelyDtoSet, Map<String, String> publicParam, String timeFormat, String nowDateFormat) {
return nodeRelyDtoSet.stream().map(nodeRelyDto -> {
Node node = nodeRelyDto.getNode();
String relyId = nodeRelyDto.getRelyId();
......@@ -174,17 +188,17 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
Date repairDate = sdf.parse(repeatTime);
ValidationUtil.isTrueValidation(!repairDate.before(new Date()), "只能补过去时间的批次!");
} catch (ParseException e) {
log.error("补批日期不符合规范,例:{}",timeFormat);
ValidationUtil.isTrueValidation(true, "补批日期不符合规范,例:"+timeFormat);
log.error("补批日期不符合规范,例:{}", timeFormat);
ValidationUtil.isTrueValidation(true, "补批日期不符合规范,例:" + timeFormat);
}
//补批只替换不是java的节点
if (!NodeTypeEnum.JAVA.getCode().equals(waitingTask.getJobType())) {
if (StringUtils.isNotEmpty(waitingTask.getRunParam())) {
String param = PlaceholderUtils.formatBizDateParam(runParamWrapped.getPrivateParam(), PlaceholderEnum.DATE_PLACEHOLDER.getName() ,repeatTime, 0);
String param = PlaceholderUtils.formatBizDateParam(runParamWrapped.getPrivateParam(), PlaceholderEnum.DATE_PLACEHOLDER.getName(), repeatTime, 0);
String dateFormat = DateUtil.dateFormat(new Date(), nowDateFormat);
param = PlaceholderUtils.formatBizDateParam(param, PlaceholderEnum.DATE_PLACEHOLDER.getName() ,dateFormat, 0);
param = PlaceholderUtils.formatBizDateParam(param, PlaceholderEnum.DATE_PLACEHOLDER.getName(), dateFormat, 0);
//设置替换完成后的参数
runParamWrapped.setPrivateParam(param);
}
......@@ -202,7 +216,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
* @return 全部的节点信息
*/
private List<WaitingTask> buildWaitingTaskResult(List<WaitingRecord> waitingRecords, Set<NodeRelyDto> nodeRelyDtoSet,
Map<String,String> publicPram, String timeFormat, String nowDateFormat) {
Map<String, String> publicPram, String timeFormat, String nowDateFormat) {
List<WaitingTask> waitingTaskListResult = new ArrayList<>(8);
waitingRecords.forEach(waitingRecord -> {
List<WaitingTask> waitingTaskList = buildWaitingTask(waitingRecord, nodeRelyDtoSet, publicPram, timeFormat, nowDateFormat);
......@@ -214,8 +228,8 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
/**
* 构建实例队列
*
* @param waitingRecords
* @return
* @param waitingRecords 等待实例
* @return 运行实例
*/
private List<RunRecording> buildRunRecordingList(List<WaitingRecord> waitingRecords) {
return waitingRecords.stream().map(waitingRecord -> {
......@@ -244,7 +258,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
if (order == null) {
order = 0;
}
List<WaitingRecord> waitingRecords = new ArrayList<WaitingRecord>(2);
List<WaitingRecord> waitingRecords = new ArrayList<>(2);
for (String repairTime : repairTimeList) {
WaitingRecord waitingRecord = new WaitingRecord();
BeanUtils.copyProperties(flow, waitingRecord);
......
......@@ -235,17 +235,18 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
public void updateRunRecordingAndTask(WaitingRecord waitingRecord) {
log.debug("---------开始查询等待工作流{}对应的数据-------------",waitingRecord);
RunRecording runRecording = runRecordingService.findAllByRunID(waitingRecord.getRunId());
runRecording.setStartTime(new Date());
//将运行实例改为以执行
runRecording.setFlowStatus(RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode());
runRecordingService.updateRunRecordingById(runRecording);
log.info("-------修改运行实例表成功,查询对应等待实例{},的等待节点-------",waitingRecord);
List<WaitingTask> allByWaitId = taskService.findAllByWaitId(waitingRecord.getWaitId());
if(waitingRecord.getFlowNodeCount() != allByWaitId.size()){
log.info("------发现工作流数量与对应的等待节点数量不对等,实例节点数量为{},查出的等待节点数量为{},本次补批或重跑操作跳过----------",waitingRecord.getFlowNodeCount(), allByWaitId.size());
return;
}
runRecording.setStartTime(new Date());
//将运行实例改为以执行
runRecording.setFlowStatus(RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode());
runRecordingService.updateRunRecordingById(runRecording);
log.info("-------查询等待节点成功,查询对应的等待节点成功,开始保存对应的等待节点{}-------",allByWaitId);
List<Integer> ids = new ArrayList<>(2);
List<JobTask> jobTasks = allByWaitId.stream().map(waitingTask -> {
......
......@@ -12,6 +12,12 @@ import java.util.List;
@Data
public class RepairFlow implements Serializable {
private static final long serialVersionUID = -1984385012809451217L;
/**
* 1. 将工作流补批
* 2. 只补批节点
*/
private String allNodes = "1";
/**
* 工作空间名称
*/
......
......@@ -132,7 +132,7 @@ public class JobUtils {
/**
* 补批工作流
*/
public static final String REQUEST_REPAIRFLOW = "/api/flow/repairFlow";
public static final String REQUEST_REPAIRFLOW = "/api/operating/node/routineSupplementBatch";
/**
* 特殊的补批接口
*/
......
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