From 3170e867162fab60f849dab67b0bc768647a68a2 Mon Sep 17 00:00:00 2001 From: stream Date: Tue, 24 Mar 2026 10:59:50 +0800 Subject: [PATCH] =?UTF-8?q?refactor(flow):=20=E4=BC=98=E5=8C=96=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E6=8E=A7=E5=88=B6=E5=BC=82=E6=AD=A5=E6=89=A7=E8=A1=8C?= =?UTF-8?q?=E5=92=8C=E8=8A=82=E7=82=B9=E6=9F=A5=E6=89=BE=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 将终端停止操作改为异步执行以提升性能 - 替换子图开始节点查找方法为入度为零节点查找 - 扩展子图结束节点判断条件增强流程控制准确性 - 移除异常捕获中的空处理块简化代码结构 --- .../service/impl/EdgeHlcServiceImpl.java | 4 ++++ .../test/flow/control/FlowControlService.java | 7 ++++++- .../FlowGetCurrentObjNodeHandler.java | 21 ++++++++++--------- .../dispatcher/FlowLoopNodeHandler.java | 4 +++- .../flow/runtime/engine/FlowItemExecutor.java | 8 ++++--- .../runtime/engine/FlowTaskScheduler.java | 2 +- .../engine/support/FlowNodeParamPreparer.java | 5 +++-- .../message/TaskNodeExecuteContext.java | 1 + .../message/TaskNodeExecuteMessage.java | 2 ++ .../impl/TiVehicleFunctionServiceImpl.java | 2 +- 10 files changed, 37 insertions(+), 19 deletions(-) diff --git a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java index 4b28786..ff9dbe7 100644 --- a/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java +++ b/cmvr-iot-api/cmvr-iot-edge/cmvr-iot-grpc-client/src/main/java/com/cmvr/edge/client/service/impl/EdgeHlcServiceImpl.java @@ -3,6 +3,7 @@ package com.cmvr.edge.client.service.impl; import cmvr.api.HlcCommand; import cmvr.api.HlcServiceGrpc; import com.alibaba.fastjson2.JSON; +import com.cmvr.common.exception.GlobalException; import com.cmvr.edge.client.manage.GrpcServiceManager; import com.cmvr.edge.client.model.hlc.EdgeTouchVO; import com.cmvr.edge.client.service.EdgeHlcService; @@ -85,6 +86,9 @@ public class EdgeHlcServiceImpl implements EdgeHlcService { @Override public String touch(EdgeTouchVO edgeTouchVO) { + if (edgeTouchVO.getX() == edgeTouchVO.getY() && edgeTouchVO.getX() == 0) { + throw new GlobalException("触控坐标不能为0"); + } HlcServiceGrpc.HlcServiceBlockingStub stub = grpcServiceManager.getGrpcClient(edgeTouchVO.getTerminalId(), HlcServiceGrpc.HlcServiceBlockingStub.class); HlcCommand.Touch.Request request = HlcCommand.Touch.Request.newBuilder() .setHeader(EdgeCommonUtil.buildRequest(edgeTouchVO.getDeviceId())) diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java index c9af14c..5e35a9c 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java @@ -78,7 +78,10 @@ public class FlowControlService { if (ctx.isPaused()) { throw new GlobalException("任务已处于暂停状态"); } - edgeSystemService.stopAll(ctx.getTerminalId()); + // 异步停止终端 + executor.execute(() -> { + edgeSystemService.stopAll(ctx.getTerminalId()); + }); ctx.setPaused(true); ctx.setTerminalStatus(TerminalStatusEnum.RUNNING); ctx.setStatus(TaskStatusEnum.PAUSED); @@ -144,6 +147,7 @@ public class FlowControlService { TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage); resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留中断时的循环次数 + resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留中断时的循环次数 resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留中断时的路径 flowItemExecutor.executeNode( @@ -169,6 +173,7 @@ public class FlowControlService { TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage); resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留循环次数 + resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留循环次数 resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留路径 flowItemExecutor.executeNode( diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java index bee6c8e..09ff6f8 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowGetCurrentObjNodeHandler.java @@ -8,6 +8,7 @@ import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; +import java.util.Collections; import java.util.List; @Slf4j @@ -18,17 +19,17 @@ public class FlowGetCurrentObjNodeHandler implements FlowNodeTypeHandler { @Override public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) { try { - JSONObject inputParams = message.getInputParams(); - JSONArray jsonArray = inputParams.getJSONArray("array"); - - List iterations = message.getIterations(); - int last = CollUtil.getLast(iterations) - 1; - +// JSONObject inputParams = message.getInputParams(); +// JSONArray jsonArray = inputParams.getJSONArray("array"); +// +// List iterations = message.getIterations(); +// int last = CollUtil.getLast(iterations) - 1; + int index = message.getLoopNum() - 1; JSONObject output = new JSONObject(); - output.put("index", last); - - if (CollUtil.isNotEmpty(jsonArray) && last >= 0) { - output.put("object", jsonArray.get(last)); + output.put("index", index); + JSONArray objects = (JSONArray)message.getLoopArray(); + if (CollUtil.isNotEmpty(objects) && index >= 0) { + output.put("object", objects.get(index)); } return TaskNodeExecuteResult.success(output); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java index 7b11a01..a9e5bb5 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/dispatcher/FlowLoopNodeHandler.java @@ -37,6 +37,7 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler { String nodeId = message.getNodeId(); FlowNodeWrapper nodeWrapper = graph.getNode(nodeId); int loopCount = inputParams.getIntValue("loopNum"); // 获取循环次数 + Object loopArray = inputParams.get("loopArray"); // 获取循环次数 // 如果没有设置循环次数或者循环次数为 0,则跳过处理 if (loopCount <= 0) { @@ -65,8 +66,9 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler { // 克隆 message,并明确设置 loopNum TaskNodeExecuteMessage subMessage = new TaskNodeExecuteMessage(); BeanUtil.copyProperties(message, subMessage); -// subMessage.setLoopNum(i); // 当前 loop 第 i 次 + subMessage.setLoopNum(i); // 当前 loop 第 i 次 subMessage.setIterations(newIterations); // 完整路径 + subMessage.setLoopArray(loopArray); flowItemExecutor.executeSubGraph(subGraph, subMessage, latch::countDown, newIterations); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java index 1eac761..45557ed 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowItemExecutor.java @@ -73,7 +73,7 @@ public class FlowItemExecutor { List startInputDefs = subGraph.getStartInputParamsDefs(); // 子图的 start 节点 - FlowNodeWrapper subStartNode = subGraph.findSubStartNodeId(); + FlowNodeWrapper subStartNode = subGraph.getInDegreeZeroNode(); String subStartNodeId = subStartNode.getNodeId(); NodeExecutor executor = node -> executeNode( @@ -84,7 +84,7 @@ public class FlowItemExecutor { // 执行子图起始节点 executeNode(subGraph, subStartNode, subStartNodeId, startInputDefs, rootMessage, onFinished, iterations); - + log.info("当前执行-子图"); // 启动子图调度 flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()), executor, onFinished, iterations); @@ -131,7 +131,7 @@ public class FlowItemExecutor { } // END 节点 or 子图结束 - if (node.getNodeType() == NodeTypeEnum.END || node.getNodeType() == NodeTypeEnum.SUB_END) { + if (node.getNodeType() == NodeTypeEnum.END || node.getNodeType() == NodeTypeEnum.SUB_END || graph.isSubGraphEndNode(nodeId)) { try { log.info("{} 节点执行,释放线程", node.getNodeType()); onFinished.run(); @@ -172,6 +172,8 @@ public class FlowItemExecutor { // loopNum 不在这里计算,而是由 FlowLoopNodeHandler / resume 显式写入 // 如果 rootMessage 里已经带了 loopNum,就沿用它 message.setLoopNum(rootMessage.getLoopNum()); + + message.setLoopArray(rootMessage.getLoopArray()); return message; } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java index b291a53..74095fc 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskScheduler.java @@ -58,7 +58,7 @@ public class FlowTaskScheduler { Runnable onFinished, List iterations) { String instId = context.getInstId(); - FlowNodeWrapper startNode = graph.findSubStartNodeId(); + FlowNodeWrapper startNode = graph.getInDegreeZeroNode(); List nextNodes = graph.getNextNodes(startNode.getNodeId()); if (nextNodes.isEmpty()) { diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java index 78a49b7..cebde64 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowNodeParamPreparer.java @@ -84,11 +84,12 @@ public class FlowNodeParamPreparer { // 4 LOOP 节点:动态计算 loopCount if (node.getNodeType() == NodeTypeEnum.LOOP) { Object loopNumVal = input.get("loopNum"); + input.put("loopArray", loopNumVal); int loopCount = 0; // 根据迭代路径找到当前层的集合对象 - Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0); - +// Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0); + Object target = loopNumVal; if (target instanceof JSONArray) { loopCount = ((JSONArray) target).size(); } else if (target instanceof Collection) { diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java index 5c2f57c..e60597a 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteContext.java @@ -22,5 +22,6 @@ public class TaskNodeExecuteContext { private TaskNodeExecuteMessage rootMessage; private Runnable onFinished; private int loopIteration; + private Object loopArray; private List iterations = new ArrayList<>(); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java index b887a37..33462f3 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/message/TaskNodeExecuteMessage.java @@ -55,6 +55,8 @@ public class TaskNodeExecuteMessage { */ private List iterations = new ArrayList<>(); + private Object loopArray; + /** * 节点名称 */ diff --git a/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java b/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java index 5681d03..fcfb7e0 100644 --- a/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java +++ b/cmvr-iot-ti/src/main/java/com/cmvr/ti/service/impl/TiVehicleFunctionServiceImpl.java @@ -40,7 +40,7 @@ public class TiVehicleFunctionServiceImpl extends ServiceImpl