From 710f33d882116e0002ea9ed543d45a427d624776 Mon Sep 17 00:00:00 2001 From: lixiaolong <702156524@qq.com> Date: Tue, 14 Jul 2026 10:12:04 +0800 Subject: [PATCH] =?UTF-8?q?refactor(flow):=20=E9=87=8D=E6=9E=84=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E6=89=A7=E8=A1=8C=E4=BA=8B=E4=BB=B6=E5=92=8C=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E7=8A=B6=E6=80=81=E7=AE=A1=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 更新 FlowExecutionEvent 类的注释和字段描述,明确流程执行事件的用途和各字段含义 - 重命名 TaskInstHolder.markSuccess 为 markNodeSuccess 方法,并拆分任务完成逻辑 - 添加 markTaskCompleted 私有方法专门处理任务完成状态更新和事件发布 - 保留旧方法名作为兼容性支持,标记为 @Deprecated - 修复 FlowLoggingInterceptor 中的方法调用以使用新的 markNodeSuccess 方法 - 优化 FlowTaskRuntimeEntry 中的终端ID检查逻辑,移除冗余空值判断 - 改进 TaskThreadRegistry.interruptAll 方法,精确匹配实例ID并添加线程中断状态检查 - 在任务中断后释放所有 latch 避免死锁,提升系统稳定性 --- .../test/flow/context/TaskInstHolder.java | 82 ++++++++++++++----- .../test/flow/context/TaskThreadRegistry.java | 13 +-- .../runtime/engine/FlowTaskRuntimeEntry.java | 9 +- .../runtime/event/FlowExecutionEvent.java | 32 ++++---- .../interceptor/FlowLoggingInterceptor.java | 2 +- 5 files changed, 90 insertions(+), 48 deletions(-) diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java index 879252d..1edaae8 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskInstHolder.java @@ -49,30 +49,68 @@ public class TaskInstHolder { nodeInstService.logRunning(instId, taskId, itemId, nodeId, nodeType, operate, action, paramsIn, iterations); } + /** + * 标记节点执行成功(仅记录节点日志,不更新任务状态) + * + * @param instId 任务实例ID + * @param itemId 检测项ID + * @param nodeId 节点ID + * @param pendingItemCount 剩余待执行检测项数量 + * @param nodeType 节点类型 + * @param paramsOut 输出参数 + * @param message 日志消息 + * @param iterations 迭代路径 + */ + public void markNodeSuccess(String instId, String itemId, String nodeId, int pendingItemCount, String nodeType, + String paramsOut, String message, List iterations) { + // 仅记录节点执行日志 + nodeInstService.logSuccess(instId, itemId, nodeId, paramsOut, message, iterations); + + // 如果是最后一个检测项的END节点,标记任务完成 + if (pendingItemCount == 0 && NodeTypeEnum.END.getCode().equalsIgnoreCase(nodeType)) { + markTaskCompleted(instId, itemId); + } + } + + /** + * 标记任务执行完成(更新任务状态为SUCCESS,发布完成事件) + * + * @param instId 任务实例ID + * @param itemId 检测项ID + */ + private void markTaskCompleted(String instId, String itemId) { + // 先获取上下文(在注销前) + TaskContext ctx = taskContextManager.get(instId); + String taskId = ctx != null ? ctx.getTaskId() : null; + String moduleCode = ctx != null ? ctx.getModuleCode() : null; + + // 更新任务状态为SUCCESS + syncStatus(instId, TaskStatusEnum.SUCCESS); + + // 注销任务上下文(释放终端锁) + taskContextManager.unregister(instId); + + // 发布任务完成事件 + eventPublisher.publishEvent(FlowExecutionEvent.builder() + .eventType(FlowExecutionEvent.EventType.TASK_COMPLETED) + .instId(instId) + .taskId(taskId) + .moduleCode(moduleCode) + .itemId(itemId) + .status(TaskStatusEnum.SUCCESS) + .build()); + + log.info("任务执行完成: instId={}, taskId={}", instId, taskId); + } + + /** + * 兼容旧方法名(已废弃,请使用 markNodeSuccess) + * @deprecated 使用 markNodeSuccess 替代 + */ + @Deprecated public void markSuccess(String instId, String itemId, String nodeId, int pendingItemCount, String nodeType, String paramsOut, String message, List iterations) { - nodeInstService.logSuccess(instId, itemId, nodeId, paramsOut, message,iterations); - syncStatus(instId, TaskStatusEnum.SUCCESS); - // 如果没有下一个检测项要执行 且是end节点 注销任务上下文 - if (pendingItemCount == 0 && nodeType.equalsIgnoreCase(NodeTypeEnum.END.getCode())) { - // 先获取上下文(在注销前) - TaskContext ctx = taskContextManager.get(instId); - String taskId = ctx != null ? ctx.getTaskId() : null; - String moduleCode = ctx != null ? ctx.getModuleCode() : null; - - // 注销任务上下文 - taskContextManager.unregister(instId); - - // 发布任务完成事件 - eventPublisher.publishEvent(FlowExecutionEvent.builder() - .eventType(FlowExecutionEvent.EventType.TASK_COMPLETED) - .instId(instId) - .taskId(taskId) - .moduleCode(moduleCode) - .itemId(itemId) - .status(TaskStatusEnum.SUCCESS) - .build()); - } + markNodeSuccess(instId, itemId, nodeId, pendingItemCount, nodeType, paramsOut, message, iterations); } public void markFailed(String instId, String taskId, String itemId, diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java index b14e31f..2669d4d 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/context/TaskThreadRegistry.java @@ -73,18 +73,21 @@ public class TaskThreadRegistry { */ public void interruptAll(String instId) { for (Map.Entry entry : nodeContextMap.entrySet()) { - if (entry.getKey().startsWith(instId)) { + if (entry.getKey().startsWith(instId + "_")) { FlowNodeWrapper node = entry.getValue().getNode(); if (node.getNodeType().equals(NodeTypeEnum.LOOP)) { - log.info("LOOP节点不直接中断"); + log.info("LOOP节点不直接中断,等待子图完成"); continue; } Thread thread = entry.getValue().getThread(); - log.info("中断线程:{}", thread.getName()); - thread.interrupt(); + if (thread != null && !thread.isInterrupted()) { + log.info("中断线程:{}", thread.getName()); + thread.interrupt(); + } } } -// releaseAllLatches(instId); + // ✅ 释放所有 latch,避免死锁 + releaseAllLatches(instId); log.info("-------------------------------------"); log.info("pending:\n{}", getPendingNodes(instId).stream().map(x -> x.getNode().getNodeId()).collect(Collectors.toList())); log.info("-------------------------------------"); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java index eabe6d6..f731ede 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskRuntimeEntry.java @@ -61,9 +61,9 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { JSONObject runParams = taskExecuteTrailVO.getRunParams(); String terminalId = taskExecuteTrailVO.getTerminalId(); if (StrUtil.isEmpty(terminalId)) { -// throw new GlobalException("终端ID不能为空"); + throw new GlobalException("终端ID不能为空"); } - if (terminalId != null && taskInstHolder.isTerminalLocked(terminalId)) { + if (taskInstHolder.isTerminalLocked(terminalId)) { throw new GlobalException("当前终端正在执行其他任务,请稍后再试"); } @@ -124,12 +124,13 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService { @Transactional(rollbackFor = Exception.class) public String executeTaskInternal(String terminalId, String taskId, String moduleCode, RunModeEnum runMode, JSONObject runParams, Object originalVO) { + // 终端id可能为空 if (StrUtil.isEmpty(terminalId)) { // throw new GlobalException("终端ID不能为空"); } // 检查终端状态 - if (terminalId != null && taskInstHolder.isTerminalLocked(terminalId)) { -// throw new GlobalException("当前终端正在执行其他任务,请稍后再试"); + if (!StrUtil.isEmpty(terminalId) && taskInstHolder.isTerminalLocked(terminalId)) { + throw new GlobalException("当前终端正在执行其他任务,请稍后再试"); } // 查询任务详情(检测项信息) List details = taskOrchestrationService.queryTaskDetail(taskId); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java index 91df724..5e47fd1 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/event/FlowExecutionEvent.java @@ -7,9 +7,9 @@ import lombok.Data; import lombok.NoArgsConstructor; /** - * 流程执行事件 + * 流程执行事件。 *

- * 携带流程执行过程中的关键信息,供外部模块监听和处理 + * 用于描述流程引擎执行过程中的关键节点和任务状态变化,供业务模块监听并处理。 */ @Data @Builder @@ -18,62 +18,62 @@ import lombok.NoArgsConstructor; public class FlowExecutionEvent { /** - * 事件类型 + * 事件类型。 */ private EventType eventType; /** - * 任务实例 ID + * 流程任务实例ID。 */ private String instId; /** - * 任务 ID + * 流程任务ID。 */ private String taskId; /** - * 模块编码 + * 事件所属业务模块编码,用于监听器路由隔离。 */ private String moduleCode; /** - * 节点数量 + * 当前检测项内的节点总数。 */ private Integer nodeCount; /** - * 当前检测项 ID + * 当前检测项ID。 */ private String itemId; /** - * 节点 ID(节点级事件时有值) + * 当前节点ID,节点级事件时有值。 */ private String nodeId; /** - * 节点类型(节点级事件时有值) + * 当前节点类型,节点级事件时有值。 */ private String nodeType; /** - * 节点名称 + * 当前节点名称。 */ private String nodeName; /** - * 任务状态 + * 当前任务状态。 */ private TaskStatusEnum status; /** - * 错误信息(失败事件时有值) + * 错误信息,失败事件时有值。 */ private String errorMessage; /** - * 事件类型枚举 + * 流程执行事件类型。 */ public enum EventType { /** 节点执行完成 */ @@ -84,9 +84,9 @@ public class FlowExecutionEvent { TASK_COMPLETED, /** 任务执行失败 */ TASK_FAILED, - /** 任务被暂停 */ + /** 任务暂停 */ TASK_PAUSED, - /** 任务被终止 */ + /** 任务终止 */ TASK_STOPPED, /** 任务恢复执行 */ TASK_RESUMED, diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java index 4fc24b7..9113185 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/interceptor/FlowLoggingInterceptor.java @@ -48,7 +48,7 @@ public class FlowLoggingInterceptor extends AbstractFlowMsgPreInterceptor { return null; } - taskInstHolder.markSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), + taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(), result.getOutputParams().toJSONString(), StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations()); long end = System.currentTimeMillis();