Commit 3494022d by guominglei

修改补批查询的节点来源表

parent a4ef0522
...@@ -483,22 +483,22 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -483,22 +483,22 @@ public class ApiFlowServiceImpl implements ApiFlowService {
RunRecording runRecording = runRecordingMapper.findRunRecordingByFlowIdAndRunId(flow.getFlowId(), runInfo.getRunId()); RunRecording runRecording = runRecordingMapper.findRunRecordingByFlowIdAndRunId(flow.getFlowId(), runInfo.getRunId());
ValidationUtil.dataNotNull(runRecording, "查无此运行记录"); ValidationUtil.dataNotNull(runRecording, "查无此运行记录");
//TODO 这里应该查询的是版本表 与版本号挂钩 //TODO 这里应该查询的是版本表 与版本号挂钩
NodeVersion node = nodeVersionMapper.findNodeByNodeNameAndVersionNameAndFlowId(runInfo.getNodeName(),runRecording.getFlowVersionName(),runRecording.getFlowId()); // NodeVersion node = nodeVersionMapper.findNodeByNodeNameAndVersionNameAndFlowId(runInfo.getNodeName(),runRecording.getFlowVersionName(),runRecording.getFlowId());
//Node node = nodeMapper.getByNameAndFlow(runInfo.getNodeName(), flow.getFlowId()); // //Node node = nodeMapper.getByNameAndFlow(runInfo.getNodeName(), flow.getFlowId());
ValidationUtil.dataNotNull(node, runInfo.getNodeName() + "节点不存在"); // ValidationUtil.dataNotNull(node, runInfo.getNodeName() + "节点不存在");
//获取运行日志实例 //获取运行日志实例
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runInfo.getRunId(), node.getNodeId()); JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndFlowIdAndNodeName(runInfo.getRunId(), flow.getFlowId(), runInfo.getNodeName());
ValidationUtil.dataNotNull(jobTaskRunLog, "查无此运行记录"); ValidationUtil.dataNotNull(jobTaskRunLog, "查无此运行记录");
String userName = currentUserUtils.account(); String userName = currentUserUtils.account();
//校验上级是否成功 //校验上级是否成功
List<Integer> dependNodeList = nodeDependencyMapper.findDependIdByNodeId(node.getNodeId()); if (StringUtils.isNotEmpty(jobTaskRunLog.getNodeDepend())){
if (dependNodeList != null && dependNodeList.size() > 0){ List<String> dependNodeList = Arrays.asList(jobTaskRunLog.getNodeDepend().split(","));
dependNodeList.forEach(dependNodeId -> { dependNodeList.forEach(dependNodeId -> {
JobTaskRunLog dependJobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runInfo.getRunId(), dependNodeId); JobTaskRunLog dependJobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runInfo.getRunId(), Integer.valueOf(dependNodeId));
if (dependJobTaskRunLog != null) { if (dependJobTaskRunLog != null) {
ValidationUtil.isTrueValidation(!("1".equals(dependJobTaskRunLog.getRunCode()) || "3".equals(dependJobTaskRunLog.getRunCode())), "上级任务未运行成功!"); ValidationUtil.isTrueValidation(!("1".equals(dependJobTaskRunLog.getRunCode()) || "3".equals(dependJobTaskRunLog.getRunCode())), "上级任务未运行成功!");
} }
...@@ -525,11 +525,11 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -525,11 +525,11 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//判断重跑机制(单节点重跑,节点及下游重跑 //判断重跑机制(单节点重跑,节点及下游重跑
if (!"1".equals(runInfo.getRunState())) {//如果不是只重跑当前节点 if (!"1".equals(runInfo.getRunState())) {//如果不是只重跑当前节点
//查询依赖本节点的节点,并添加到集合中 //查询依赖本节点的节点,并添加到集合中
List<Integer> subNodeIdList = nodeDependencyMapper.findSubNodeList(node.getNodeId()); List<Integer> subNodeIdList = nodeVersionDependencyMapper.findSubNodeList(jobTaskRunLog.getNodeId());
Set<Integer> nodeIdSet = new HashSet<>(); Set<Integer> nodeIdSet = new HashSet<>();
Set<Integer> runNodeIdSet = new HashSet<>(); Set<Integer> runNodeIdSet = new HashSet<>();
runNodeIdSet.add(node.getNodeId()); runNodeIdSet.add(jobTaskRunLog.getNodeId());
nodeIdSet.add(node.getNodeId()); nodeIdSet.add(jobTaskRunLog.getNodeId());
//获取全部的需要重跑的节点 //获取全部的需要重跑的节点
if (subNodeIdList != null && subNodeIdList.size() > 0){ if (subNodeIdList != null && subNodeIdList.size() > 0){
Queue<Integer> queue = new LinkedList<>(); Queue<Integer> queue = new LinkedList<>();
...@@ -538,7 +538,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -538,7 +538,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
runNodeIdSet.add(nodeId); runNodeIdSet.add(nodeId);
}); });
while(!queue.isEmpty()){ while(!queue.isEmpty()){
List<Integer> innerSubNodeIdList = nodeDependencyMapper.findSubNodeList(queue.poll()); List<Integer> innerSubNodeIdList = nodeVersionDependencyMapper.findSubNodeList(queue.poll());
if (innerSubNodeIdList != null && innerSubNodeIdList.size() > 0){ if (innerSubNodeIdList != null && innerSubNodeIdList.size() > 0){
innerSubNodeIdList.forEach(nodeId->{ innerSubNodeIdList.forEach(nodeId->{
queue.offer(nodeId); queue.offer(nodeId);
...@@ -667,7 +667,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -667,7 +667,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
waitingTask.setRunParam(jobTaskRunLog.getRunParams()); waitingTask.setRunParam(jobTaskRunLog.getRunParams());
//查询当前节点的依赖节点 //查询当前节点的依赖节点
List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(jobTaskRunLog.getNodeId()); List<Integer> dependNodeIdList = nodeVersionDependencyMapper.findDependIdByNodeId(jobTaskRunLog.getNodeId());
if (dependNodeIdList != null && dependNodeIdList.size() > 0){ if (dependNodeIdList != null && dependNodeIdList.size() > 0){
List<Integer> realDependNodeIdList = new ArrayList<>(); List<Integer> realDependNodeIdList = new ArrayList<>();
for (Integer dependNodeId : dependNodeIdList) { for (Integer dependNodeId : dependNodeIdList) {
...@@ -681,7 +681,7 @@ public class ApiFlowServiceImpl implements ApiFlowService { ...@@ -681,7 +681,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
waitingTaskList.add(waitingTask); waitingTaskList.add(waitingTask);
} }
//查询依赖于当前节点的下级节点 //查询依赖于当前节点的下级节点
List<Integer> childNodeIdList = nodeDependencyMapper.findSubNodeList(childNodeId); List<Integer> childNodeIdList = nodeVersionDependencyMapper.findSubNodeList(childNodeId);
if (childNodeIdList != null && childNodeIdList.size() > 0){ if (childNodeIdList != null && childNodeIdList.size() > 0){
addDependNode(runId, triggerTime, waitingTaskList, childNodeIdList, userName, nodeIdSet, runNodeIdSet); addDependNode(runId, triggerTime, waitingTaskList, childNodeIdList, userName, nodeIdSet, runNodeIdSet);
} }
......
...@@ -156,4 +156,13 @@ public interface JobTaskRunLogMapper { ...@@ -156,4 +156,13 @@ public interface JobTaskRunLogMapper {
JobTaskRunLog findNewJavaTaskByJobName(Integer javaTaskId); JobTaskRunLog findNewJavaTaskByJobName(Integer javaTaskId);
int updateByFlowIdAndRunId(@Param("flowId")Integer flowId, @Param("runId")String runId); int updateByFlowIdAndRunId(@Param("flowId")Integer flowId, @Param("runId")String runId);
/**
* 根据runid和flowid,节点名称查询运行历史
* @param runId
* @param flowId
* @param nodeName
* @return
*/
JobTaskRunLog findByRunIdAndFlowIdAndNodeName(@Param("runId")String runId, @Param("flowId")Integer flowId, @Param("nodeName")String nodeName);
} }
\ No newline at end of file
...@@ -3,10 +3,16 @@ package com.byit.mapper; ...@@ -3,10 +3,16 @@ package com.byit.mapper;
import com.byit.model.NodeVersionDependencyKey; import com.byit.model.NodeVersionDependencyKey;
import org.apache.ibatis.annotations.Param; import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface NodeVersionDependencyMapper { public interface NodeVersionDependencyMapper {
int deleteById(NodeVersionDependencyKey key); int deleteById(NodeVersionDependencyKey key);
int insert(@Param("nodeVersionId") Integer nodeVersionId, @Param("dependencyId") Integer dependencyId); int insert(@Param("nodeVersionId") Integer nodeVersionId, @Param("dependencyId") Integer dependencyId);
int insertSelective(NodeVersionDependencyKey record); int insertSelective(NodeVersionDependencyKey record);
List<Integer> findSubNodeList(Integer nodeId);
List<Integer> findDependIdByNodeId(Integer nodeId);
} }
\ No newline at end of file
...@@ -233,6 +233,13 @@ ...@@ -233,6 +233,13 @@
) )
</select> </select>
<select id="findByRunIdAndFlowIdAndNodeName" resultMap="BaseResultMap">
select <include refid="Base_Column_List"/>
from job_task_run_log
where run_id = #{runId} and flow_id = #{flowId}
and node_name = #{nodeName}
</select>
<delete id="deleteById" parameterType="java.lang.Integer"> <delete id="deleteById" parameterType="java.lang.Integer">
delete from job_task_run_log delete from job_task_run_log
where log_id = #{logId,jdbcType=INTEGER} where log_id = #{logId,jdbcType=INTEGER}
......
...@@ -12,6 +12,16 @@ ...@@ -12,6 +12,16 @@
where node_version_id = #{nodeVersionId,jdbcType=INTEGER} where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
and dependency_id = #{dependencyId,jdbcType=INTEGER} and dependency_id = #{dependencyId,jdbcType=INTEGER}
</delete> </delete>
<select id="findSubNodeList" resultType="java.lang.Integer">
select node_version_id
from node_version_dependency
where dependency_id = #{nodeId,jdbcType=INTEGER}
</select>
<select id="findDependIdByNodeId" resultType="java.lang.Integer">
select dependency_id
from node_version_dependency
where node_version_id = #{nodeId,jdbcType=INTEGER}
</select>
<insert id="insert"> <insert id="insert">
<!-- generated @mbg.generated date: 2020-01-03 --> <!-- generated @mbg.generated date: 2020-01-03 -->
insert into node_version_dependency (node_version_id, dependency_id) insert into node_version_dependency (node_version_id, dependency_id)
......
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