Commit ba3f6ea8 by huangfusuper

修正查询实例状态接口(基于工作流名称查询最后时间的实例)

parent e0a6ed91
......@@ -15,11 +15,11 @@ import java.util.List;
*/
@Api(tags = "工作流查询api")
@RestController
@RequestMapping("api/find/flow/")
public class ApiFlowFindController {
@RequestMapping("api/find/runRecording/")
public class ApiRunRecordingFindController {
private final ApiFlowFindService apiFlowFindService;
public ApiFlowFindController(ApiFlowFindService apiFlowFindService) {
public ApiRunRecordingFindController(ApiFlowFindService apiFlowFindService) {
this.apiFlowFindService = apiFlowFindService;
}
......@@ -27,4 +27,9 @@ public class ApiFlowFindController {
public List<RunRecordingStatusDto> findRunRecordingStatusByRunId(@RequestBody SelectRecordingStatusCondition selectRecordingStatusCondition){
return apiFlowFindService.findRunRecordingStatusByRunId(selectRecordingStatusCondition);
}
@RequestMapping("findRunRecordingStatusByFlowNameAndWorkspaceName")
public List<RunRecordingStatusDto> findRunRecordingStatusByFlowNameAndWorkspaceName(@RequestBody SelectRecordingStatusCondition selectRecordingStatusCondition){
return apiFlowFindService.findRunRecordingStatusByFlowNameAndWorkspaceName(selectRecordingStatusCondition);
}
}
......@@ -2,6 +2,7 @@ package com.byit.service;
import com.byit.dto.recording.RunRecordingStatusDto;
import com.byit.dto.recording.SelectRecordingStatusCondition;
import org.springframework.web.bind.annotation.RequestBody;
import java.util.List;
......@@ -15,4 +16,11 @@ public interface ApiFlowFindService {
* @return 对应的实例状态信息
*/
List<RunRecordingStatusDto> findRunRecordingStatusByRunId(SelectRecordingStatusCondition selectRecordingStatusCondition);
/**
* 查询运行实例状态根据工作流名字和工作空间名字
* @param selectRecordingStatusCondition 查询条件
* @return 对应的实例状态信息
*/
List<RunRecordingStatusDto> findRunRecordingStatusByFlowNameAndWorkspaceName(SelectRecordingStatusCondition selectRecordingStatusCondition);
}
......@@ -5,14 +5,19 @@ import com.byit.dto.recording.RunRecordingStatusDto;
import com.byit.dto.recording.SelectRecordingStatusCondition;
import com.byit.enums.RunRecordingEnum;
import com.byit.mapper.RunRecordingMapper;
import com.byit.mapper.WorkspaceMapper;
import com.byit.model.RunRecording;
import com.byit.model.Workspace;
import com.byit.service.ApiFlowFindService;
import com.byit.utils.ValidationUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
......@@ -26,58 +31,79 @@ import java.util.stream.Collectors;
@Slf4j
public class ApiFlowFindServiceImpl implements ApiFlowFindService {
private final RunRecordingMapper runRecordingMapper;
private final WorkspaceMapper workspaceMapper;
public ApiFlowFindServiceImpl(RunRecordingMapper runRecordingMapper) {
public ApiFlowFindServiceImpl(RunRecordingMapper runRecordingMapper, WorkspaceMapper workspaceMapper) {
this.runRecordingMapper = runRecordingMapper;
this.workspaceMapper = workspaceMapper;
}
@Override
public List<RunRecordingStatusDto> findRunRecordingStatusByRunId(SelectRecordingStatusCondition selectRecordingStatusCondition) {
ValidationUtil.dataNotNull(selectRecordingStatusCondition,"查询条件不允许为空!");
ValidationUtil.dataNotNull(selectRecordingStatusCondition, "查询条件不允许为空!");
List<String> runIds = selectRecordingStatusCondition.getRunIds();
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(runIds), "运行标识不允许为空!");
//获取对应的工作流实例
List<RunRecording> allRunRecordingByRunIds = runRecordingMapper.findAllRunRecordingByRunIds(runIds);
log.debug("---------查询到有{}个实例-------",allRunRecordingByRunIds);
log.debug("---------查询到有{}个实例-------", allRunRecordingByRunIds);
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(allRunRecordingByRunIds), "没有找到运行标识对应的运行实例!");
return buildRunRecordingStatus(allRunRecordingByRunIds);
}
@Override
public List<RunRecordingStatusDto> findRunRecordingStatusByFlowNameAndWorkspaceName(SelectRecordingStatusCondition selectRecordingStatusCondition) {
ValidationUtil.dataNotNull(selectRecordingStatusCondition, "查询条件不允许为空!");
List<String> flowNameList = selectRecordingStatusCondition.getFlowNameList();
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(flowNameList), "要查询的工作流的名字不允许为空!");
ValidationUtil.dataNotBank(selectRecordingStatusCondition.getWorkspaceName(), "工作空间名称不允许为空!");
Workspace workspace = workspaceMapper.getByName(selectRecordingStatusCondition.getWorkspaceName());
ValidationUtil.dataNotNull(workspace, "工作空间不存在!");
List<RunRecording> runRecordingByFlowNamesAndWorkspace = runRecordingMapper.findRunRecordingByFlowNamesAndWorkspace(flowNameList, workspace.getWorkspaceId());
Collection<RunRecording> runRecordings = runRecordingByFlowNamesAndWorkspace.stream()
.collect(Collectors.toMap(RunRecording::getFlowName,
Function.identity(),
(c1, c2) -> c1.getTriggerTime() > c2.getTriggerTime() ? c1 : c2)).values();
log.debug("---------查询到有{}个实例-------", runRecordingByFlowNamesAndWorkspace);
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(runRecordingByFlowNamesAndWorkspace), "没有找到对应的运行实例!");
return buildRunRecordingStatus(new ArrayList<>(runRecordings));
}
/**
* 构建运行实例
*
* @param allRunRecordingByRunIds 所有的实例表
* @return 对应实例的状态
*/
private List<RunRecordingStatusDto> buildRunRecordingStatus(List<RunRecording> allRunRecordingByRunIds){
private List<RunRecordingStatusDto> buildRunRecordingStatus(List<RunRecording> allRunRecordingByRunIds) {
return allRunRecordingByRunIds.stream().map(runRecording -> {
RunRecordingStatusDto runRecordingStatusDto = new RunRecordingStatusDto();
runRecordingStatusDto.setFlowName(runRecording.getFlowName());
runRecordingStatusDto.setRunId(runRecording.getRunId());
runRecordingStatusDto.setTriggerTime(runRecording.getTriggerTime());
//工作流实例处于未开始的状态
if(RunRecordingEnum.FLOW_STATUS_NOT_RUN.getCode().equals(runRecording.getFlowStatus())){
if (RunRecordingEnum.FLOW_STATUS_NOT_RUN.getCode().equals(runRecording.getFlowStatus())) {
runRecordingStatusDto.setStatus(RunRecordingStatusDto.QUEUE_ING);
}else if(RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode().equals(runRecording.getFlowStatus())) {
} else if (RunRecordingEnum.FLOW_STATUS_RUN_ING.getCode().equals(runRecording.getFlowStatus())) {
//运行中
runRecordingStatusDto.setStatus(RunRecordingStatusDto.RUN_ING);
}else if(RunRecordingEnum.FLOW_STATUS_IS_STOP.getCode().equals(runRecording.getFlowStatus())){
} else if (RunRecordingEnum.FLOW_STATUS_IS_STOP.getCode().equals(runRecording.getFlowStatus())) {
//运行中
runRecordingStatusDto.setStatus(RunRecordingStatusDto.RUN_ING);
}else if(RunRecordingEnum.FLOW_STATUS_IS_END.getCode().equals(runRecording.getFlowStatus())){
} else if (RunRecordingEnum.FLOW_STATUS_IS_END.getCode().equals(runRecording.getFlowStatus())) {
//完结状态
//成功状态 正常成功
if(RunRecordingEnum.RUN_FLOW_SUCCESS.getCode().equals(runRecording.getFlowRunResult())){
if (RunRecordingEnum.RUN_FLOW_SUCCESS.getCode().equals(runRecording.getFlowRunResult())) {
runRecordingStatusDto.setStatus(RunRecordingStatusDto.SUCCESS);
}else if(RunRecordingEnum.RUN_FLOW_RE_SUCCESS.getCode().equals(runRecording.getFlowRunResult())){
} else if (RunRecordingEnum.RUN_FLOW_RE_SUCCESS.getCode().equals(runRecording.getFlowRunResult())) {
//补批成功
runRecordingStatusDto.setStatus(RunRecordingStatusDto.SUCCESS);
}else if(RunRecordingEnum.RUN_FLOW_FAILURE.getCode().equals(runRecording.getFlowRunResult())){
} else if (RunRecordingEnum.RUN_FLOW_FAILURE.getCode().equals(runRecording.getFlowRunResult())) {
//补批失败
runRecordingStatusDto.setStatus(RunRecordingStatusDto.FAILURE);
}else if(RunRecordingEnum.RUN_FLOW_RE_FAILURE.getCode().equals(runRecording.getFlowRunResult())){
} else if (RunRecordingEnum.RUN_FLOW_RE_FAILURE.getCode().equals(runRecording.getFlowRunResult())) {
//补批失败
runRecordingStatusDto.setStatus(RunRecordingStatusDto.FAILURE);
}else if(RunRecordingEnum.RUN_FLOW_KILL.getCode().equals(runRecording.getFlowRunResult())){
} else if (RunRecordingEnum.RUN_FLOW_KILL.getCode().equals(runRecording.getFlowRunResult())) {
//补批失败
runRecordingStatusDto.setStatus(RunRecordingStatusDto.FAILURE);
}
......
......@@ -6,6 +6,7 @@ import com.byit.dto.NodeRelyDto;
import com.byit.dto.executor.RunParamWrapped;
import com.byit.dto.specials.MakeUpReturn;
import com.byit.dto.specials.RepairFlow;
import com.byit.dto.specials.RepairTimeParam;
import com.byit.dto.specials.SpecialJobParam;
import com.byit.enums.NodeTypeEnum;
import com.byit.enums.PlaceholderEnum;
......@@ -95,7 +96,7 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
//获取公共参数
Map<String, String> publicParam = specialJobParam.getPublicParam();
//补批时间
List<String> repairTimeList = specialJobParam.getRepairTimeList();
List<RepairTimeParam> repairTimeList = specialJobParam.getRepairTimeList();
ValidationUtil.isTrueValidation(CollectionUtil.isEmpty(repairTimeList), "补批时间不允许为空");
//获取工作空间
Workspace workspace = workspaceService.getByName(workspaceName);
......@@ -265,14 +266,16 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
* @param repairTimeList 补时间
* @return 等待队列集合
*/
private List<WaitingRecord> buildWaitingRecordList(Flow flow, Integer nodeCount, String operator, List<String> repairTimeList) {
private List<WaitingRecord> buildWaitingRecordList(Flow flow, Integer nodeCount, String operator, List<RepairTimeParam> repairTimeList) {
//查看当前排队的工作流最大排队序号
Integer order = waitingRecordMapper.findOrderByFlowId(flow.getFlowId());
if (order == null) {
order = 0;
}
List<WaitingRecord> waitingRecords = new ArrayList<>(2);
for (String repairTime : repairTimeList) {
for (RepairTimeParam repairTimeParam : repairTimeList) {
String repeatTime = repairTimeParam.getRepairTime();
String timeTypeName = repairTimeParam.getTimeTypeName();
WaitingRecord waitingRecord = new WaitingRecord();
BeanUtils.copyProperties(flow, waitingRecord);
//设置排期
......@@ -282,10 +285,12 @@ public class ApiFlowOperatingServiceImpl implements ApiFlowOperatingService {
waitingRecord.setFlowNodeCount(nodeCount);
waitingRecord.setOperator(operator);
waitingRecord.setScheduleType(ScheduleTypeEnum.REPAIR.getCode());
waitingRecord.setRepeatTime(repairTime);
waitingRecord.setRepeatTime(repeatTime);
waitingRecord.setWaitOrder(++order);
//需要按着时间先后来设置时间
waitingRecord.setTriggerTime(System.currentTimeMillis());
long time = DateUtil.strFormatDate(repeatTime, timeTypeName).getTime();
waitingRecord.setTriggerTime(time);
waitingRecords.add(waitingRecord);
}
......
......@@ -33,6 +33,15 @@ public interface RunRecordingMapper {
*/
List<RunRecording> findAllRunRecordingByRunIds(@Param("runIds")List<String> runIds);
/**
* 查询统一工作空间下的工作流的最新的实例
* @param flowNameList 工作流名称
* @param workspaceId 工作空间名称
* @return 对应的最新的实例
*/
List<RunRecording> findRunRecordingByFlowNamesAndWorkspace(@Param("flowNameList")List<String> flowNameList,@Param("workspaceId")Integer workspaceId);
/**
* 查询该工作流中运行中的数据
* @param flowName
......
......@@ -66,6 +66,21 @@
</select>
<select id="findRunRecordingByFlowNamesAndWorkspace" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
from
run_recording
where
flow_name in (
<foreach collection="flowNameList" item="flowName" separator=",">
#{flowName}
</foreach>
) and workspace_id = #{workspaceId}
</select>
<select id="findAll" parameterType="com.byit.dto.FlowConditionDto" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
......
......@@ -34,6 +34,7 @@ public class RunRecordingStatusDto implements Serializable {
private String flowName;
private String runId;
private Long triggerTime;
/**
* 对应DDMP的质检状态: 1:排队中; 2:质检中; 3:质检完成; 4:质检失败
* 对应调度的运行状态: 1:未开始; 2:运行中/暂停中 4-1|4-3:成功 4-2|4-4|4-5:失败
......
......@@ -18,4 +18,13 @@ import java.util.List;
public class SelectRecordingStatusCondition implements Serializable {
private static final long serialVersionUID = -8489569647307394582L;
List<String> runIds;
/**
* 工作流名称
*/
private List<String> flowNameList;
/**
* 工作空间名称
*/
private String workspaceName;
}
......@@ -33,7 +33,7 @@ public class RepairFlow implements Serializable {
/**
* 补批的时间
*/
private List<String> repairTimes;
private List<RepairTimeParam> repairTimes;
/**
* 时间类型格式
*/
......
package com.byit.dto.specials;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* @date 2020年10月30日12:59:02
* @author huangfu
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class RepairTimeParam {
private String repairTime;
private String timeTypeName = "yyyyMMdd";
}
......@@ -52,7 +52,7 @@ public class SpecialJobParam implements Serializable {
/**
* 时间参数
*/
private List<String> repairTimeList;
private List<RepairTimeParam> repairTimeList;
/**
* 时间类型格式
*/
......
......@@ -45,7 +45,10 @@ public class JobUtils {
/**
* 根据运行标识查询工作流的状态
*/
private static final String FIND_RUNRECORDING_STATUS_BY_RUNID = "/api/find/flow/findRunRecordingStatusByRunId";
private static final String FIND_RUNRECORDING_STATUS_BY_RUNID = "/api/find/runRecording/findRunRecordingStatusByRunId";/**
* 根据工作流名称和工作空间名称查询工作流的状态
*/
private static final String FIND_RUNRECORDING_STATUS_BY_FLOWNAME_AND_WORKSPACENAME = "/api/find/runRecording/findRunRecordingStatusByFlowNameAndWorkspaceName";
private static final String KILL_NODE_RUN_ING = "/api/node/killNode";
......@@ -292,6 +295,19 @@ public class JobUtils {
return JSON.parseObject(response, ResponseResult.class);
}
/**
* 查询工作流实例的状态
*
* @return
*/
public static ResponseResult findRunRecordingStatusByFlowNameAndWorkspaceName(SelectRecordingStatusCondition selectRecordingStatusCondition) {
log.debug("-------------查询工作流实例的状态{}----------------------",selectRecordingStatusCondition);
String response = createHttpRequest(FIND_RUNRECORDING_STATUS_BY_FLOWNAME_AND_WORKSPACENAME, JSON.toJSONString(selectRecordingStatusCondition, WriteClassName));
log.debug("-------------查询工作流实例的状态:{}----------------------",response);
return JSON.parseObject(response, ResponseResult.class);
}
public static ResponseResult findNoeEnd(RunRecordingStatus runRecordingStatus) {
log.debug("-------------查询工作流实例是否全部完结{}----------------------",runRecordingStatus);
......
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