Commit b0d45cea by guominglei

重跑时根据版本表信息查询依赖关系

parent 14b49e56
......@@ -483,7 +483,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
RunRecording runRecording = runRecordingMapper.findRunRecordingByFlowIdAndRunId(flow.getFlowId(), runInfo.getRunId());
ValidationUtil.dataNotNull(runRecording, "查无此运行记录");
//TODO 这里应该查询的是版本表 与版本号挂钩
// NodeVersion node = nodeVersionMapper.findNodeByNodeNameAndVersionNameAndFlowId(runInfo.getNodeName(),runRecording.getFlowVersionName(),runRecording.getFlowId());
NodeVersion nodeVersion = nodeVersionMapper.findNodeByNodeNameAndVersionNameAndFlowId(runInfo.getNodeName(),runRecording.getFlowVersionName(),runRecording.getFlowId());
// //Node node = nodeMapper.getByNameAndFlow(runInfo.getNodeName(), flow.getFlowId());
// ValidationUtil.dataNotNull(node, runInfo.getNodeName() + "节点不存在");
......@@ -525,28 +525,28 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//判断重跑机制(单节点重跑,节点及下游重跑
if (!"1".equals(runInfo.getRunState())) {//如果不是只重跑当前节点
//查询依赖本节点的节点,并添加到集合中
List<Integer> subNodeIdList = nodeVersionDependencyMapper.findSubNodeList(jobTaskRunLog.getNodeId());
List<NodeVersion> subNodeList = nodeVersionMapper.findSubNodeList(nodeVersion.getNodeVersionId());
Set<Integer> nodeIdSet = new HashSet<>();
Set<Integer> runNodeIdSet = new HashSet<>();
runNodeIdSet.add(jobTaskRunLog.getNodeId());
nodeIdSet.add(jobTaskRunLog.getNodeId());
//获取全部的需要重跑的节点
if (subNodeIdList != null && subNodeIdList.size() > 0){
if (subNodeList != null && subNodeList.size() > 0){
Queue<Integer> queue = new LinkedList<>();
subNodeIdList.forEach(nodeId -> {
queue.offer(nodeId);
runNodeIdSet.add(nodeId);
subNodeList.forEach(subNode -> {
queue.offer(subNode.getNodeVersionId());
runNodeIdSet.add(subNode.getNodeId());
});
while(!queue.isEmpty()){
List<Integer> innerSubNodeIdList = nodeVersionDependencyMapper.findSubNodeList(queue.poll());
if (innerSubNodeIdList != null && innerSubNodeIdList.size() > 0){
innerSubNodeIdList.forEach(nodeId->{
queue.offer(nodeId);
runNodeIdSet.add(nodeId);
List<NodeVersion> innerSubNodeList = nodeVersionMapper.findSubNodeList(queue.poll());
if (innerSubNodeList != null && innerSubNodeList.size() > 0){
innerSubNodeList.forEach(innerSubNode->{
queue.offer(innerSubNode.getNodeVersionId());
runNodeIdSet.add(innerSubNode.getNodeId());
});
}
}
addDependNode(runInfo.getRunId(), triggerTime, waitingTaskList, subNodeIdList, userName, nodeIdSet, runNodeIdSet);
addDependNode(runInfo.getRunId(), triggerTime, waitingTaskList, subNodeList, userName, nodeIdSet, runNodeIdSet);
}
}
......@@ -672,9 +672,9 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
}
private void addDependNode(String runId, Long triggerTime, List<WaitingTask> waitingTaskList, List<Integer> subNodeIdList, String userName, Set<Integer> nodeIdSet, Set<Integer> runNodeIdSet) {
subNodeIdList.forEach(childNodeId -> {
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, childNodeId);
private void addDependNode(String runId, Long triggerTime, List<WaitingTask> waitingTaskList, List<NodeVersion> subNodeList, String userName, Set<Integer> nodeIdSet, Set<Integer> runNodeIdSet) {
subNodeList.forEach(childNode -> {
JobTaskRunLog jobTaskRunLog = jobTaskRunLogMapper.findByRunIdAndNodeId(runId, childNode.getNodeId());
WaitingTask waitingTask = new WaitingTask();
BeanUtils.copyProperties(jobTaskRunLog, waitingTask);
waitingTask.setTriggerTime(triggerTime);
......@@ -685,23 +685,23 @@ public class ApiFlowServiceImpl implements ApiFlowService {
waitingTask.setRunParam(jobTaskRunLog.getRunParams());
//查询当前节点的依赖节点
List<Integer> dependNodeIdList = nodeVersionDependencyMapper.findDependIdByNodeId(jobTaskRunLog.getNodeId());
if (dependNodeIdList != null && dependNodeIdList.size() > 0){
if (StringUtils.isNotEmpty(jobTaskRunLog.getNodeDepend())){
List<String> dependNodeIdList = Arrays.asList(jobTaskRunLog.getNodeDepend().split(","));
List<Integer> realDependNodeIdList = new ArrayList<>();
for (Integer dependNodeId : dependNodeIdList) {
if (!runNodeIdSet.add(dependNodeId)){
realDependNodeIdList.add(dependNodeId);
for (String dependNodeId : dependNodeIdList) {
if (!runNodeIdSet.add(Integer.valueOf(dependNodeId))){
realDependNodeIdList.add(Integer.valueOf(dependNodeId));
}
}
waitingTask.setNodeDepend(Joiner.on(",").join(realDependNodeIdList));
}
if (nodeIdSet.add(childNodeId)){
if (nodeIdSet.add(childNode.getNodeId())){
waitingTaskList.add(waitingTask);
}
//查询依赖于当前节点的下级节点
List<Integer> childNodeIdList = nodeVersionDependencyMapper.findSubNodeList(childNodeId);
if (childNodeIdList != null && childNodeIdList.size() > 0){
addDependNode(runId, triggerTime, waitingTaskList, childNodeIdList, userName, nodeIdSet, runNodeIdSet);
List<NodeVersion> childNodeList = nodeVersionMapper.findSubNodeList(childNode.getNodeId());
if (childNodeList != null && childNodeList.size() > 0){
addDependNode(runId, triggerTime, waitingTaskList, childNodeList, userName, nodeIdSet, runNodeIdSet);
}
});
}
......
......@@ -43,4 +43,6 @@ public interface NodeVersionMapper {
* @return
*/
int removeByFlowId(Integer flowId);
List<NodeVersion> findSubNodeList(Integer nodeId);
}
\ No newline at end of file
......@@ -88,8 +88,15 @@
where flow_id = #{flowId} and version_name = #{versionName} and node_name = #{nodeName}
</select>
<select id="findSubNodeList" resultMap="BaseResultMap">
select node.node_version_id, node.node_id
from node_version node
left join node_version_dependency dependency on node.node_version_id = dependency.node_version_id
where dependency.dependency_id = #{nodeId,jdbcType=INTEGER}
</select>
<delete id="deleteById" parameterType="java.lang.Integer">
<delete id="deleteById" parameterType="java.lang.Integer">
<!-- generated @mbg.generated date: 2019-12-31 -->
delete from node_version
where node_version_id = #{nodeVersionId,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