From eff8fdd4b6cf79ac5f6a4606f74e53230bd1ac0c Mon Sep 17 00:00:00 2001 From: lixiaolong <702156524@qq.com> Date: Fri, 14 Aug 2026 12:09:02 +0800 Subject: [PATCH] =?UTF-8?q?fix(flow):=20=E4=BF=AE=E5=A4=8D=E6=B5=81?= =?UTF-8?q?=E7=A8=8B=E8=B0=83=E5=BA=A6=E5=99=A8=E7=BC=93=E5=AD=98=E6=B8=85?= =?UTF-8?q?=E7=90=86=E9=80=BB=E8=BE=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 修改 FlowSchedulerCache 中的分隔符从下划线改为冒号 - 在 FlowTaskEngine 中为每个编排项添加调度器去重逻辑 - 移除 TeTaskOrchestrationController 中不必要的分页初始化 - 添加 FlowSchedulerCacheTest 测试用例验证复制节点ID的调度功能 --- .../test/TeTaskOrchestrationController.java | 1 - .../flow/runtime/engine/FlowTaskEngine.java | 3 ++ .../engine/support/FlowSchedulerCache.java | 5 +-- .../support/FlowSchedulerCacheTest.java | 33 +++++++++++++++++++ 4 files changed, 39 insertions(+), 3 deletions(-) create mode 100644 cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java diff --git a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java index 087674d..be5dc9e 100644 --- a/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java +++ b/cmvr-iot-admin/src/main/java/com/cmvr/web/controller/test/TeTaskOrchestrationController.java @@ -33,7 +33,6 @@ class TeTaskOrchestrationController extends BaseController { @NotEmpty(message = "任务ID不能为空") @RequestParam("taskId") String taskId ) { - startPage(); List list = teTaskOrchestrationService.queryByTaskId(taskId); return getDataTable(list); } diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java index a10077f..8d72639 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/FlowTaskEngine.java @@ -41,6 +41,9 @@ public class FlowTaskEngine { Map itemOccurrences = new HashMap<>(); for (int i = 0; i < totalItemCount; i++) { + // Scheduler deduplication is scoped to one orchestration item. Different + // items may legitimately contain the same node IDs (for example, copied flows). + flowTaskScheduler.clear(instId); TeQueryTaskDetailDTO item = items.get(i); // 等待当前 item 完成后再执行下一个 CountDownLatch latch = new CountDownLatch(1); diff --git a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java index e2977bc..e86e999 100644 --- a/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java +++ b/cmvr-iot-test/src/main/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCache.java @@ -24,7 +24,8 @@ public class FlowSchedulerCache { } public void clearContextCache(String instId) { - scheduled.keySet().removeIf(k -> k.startsWith(instId + "_")); - completedPreMap.keySet().removeIf(k -> k.startsWith(instId + "_")); + String prefix = instId + ":"; + scheduled.keySet().removeIf(k -> k.startsWith(prefix)); + completedPreMap.keySet().removeIf(k -> k.startsWith(prefix)); } } diff --git a/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java new file mode 100644 index 0000000..dce68a3 --- /dev/null +++ b/cmvr-iot-test/src/test/java/com/cmvr/test/flow/runtime/engine/support/FlowSchedulerCacheTest.java @@ -0,0 +1,33 @@ +package com.cmvr.test.flow.runtime.engine.support; + +import org.junit.Test; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class FlowSchedulerCacheTest { + + @Test + public void clearAllowsCopiedNodeIdToBeScheduledByNextItem() { + FlowSchedulerCache cache = new FlowSchedulerCache(); + String key = "inst-1:copied-node"; + + assertTrue(cache.trySchedule(key)); + assertFalse(cache.trySchedule(key)); + + cache.clearContextCache("inst-1"); + + assertTrue(cache.trySchedule(key)); + } + + @Test + public void clearDoesNotAffectAnotherInstanceWithSimilarId() { + FlowSchedulerCache cache = new FlowSchedulerCache(); + String otherInstanceKey = "inst-10:node-1"; + assertTrue(cache.trySchedule(otherInstanceKey)); + + cache.clearContextCache("inst-1"); + + assertFalse(cache.trySchedule(otherInstanceKey)); + } +}