fix(flow): 修复流程调度器缓存清理逻辑

- 修改 FlowSchedulerCache 中的分隔符从下划线改为冒号
- 在 FlowTaskEngine 中为每个编排项添加调度器去重逻辑
- 移除 TeTaskOrchestrationController 中不必要的分页初始化
- 添加 FlowSchedulerCacheTest 测试用例验证复制节点ID的调度功能
This commit is contained in:
lixiaolong 2026-08-14 12:09:02 +08:00
parent 54220ccaac
commit eff8fdd4b6
4 changed files with 39 additions and 3 deletions

View File

@ -33,7 +33,6 @@ class TeTaskOrchestrationController extends BaseController {
@NotEmpty(message = "任务ID不能为空") @NotEmpty(message = "任务ID不能为空")
@RequestParam("taskId") String taskId @RequestParam("taskId") String taskId
) { ) {
startPage();
List<TeQueryTaskOrchestraItemVO> list = teTaskOrchestrationService.queryByTaskId(taskId); List<TeQueryTaskOrchestraItemVO> list = teTaskOrchestrationService.queryByTaskId(taskId);
return getDataTable(list); return getDataTable(list);
} }

View File

@ -41,6 +41,9 @@ public class FlowTaskEngine {
Map<String, Integer> itemOccurrences = new HashMap<>(); Map<String, Integer> itemOccurrences = new HashMap<>();
for (int i = 0; i < totalItemCount; i++) { 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); TeQueryTaskDetailDTO item = items.get(i);
// 等待当前 item 完成后再执行下一个 // 等待当前 item 完成后再执行下一个
CountDownLatch latch = new CountDownLatch(1); CountDownLatch latch = new CountDownLatch(1);

View File

@ -24,7 +24,8 @@ public class FlowSchedulerCache {
} }
public void clearContextCache(String instId) { public void clearContextCache(String instId) {
scheduled.keySet().removeIf(k -> k.startsWith(instId + "_")); String prefix = instId + ":";
completedPreMap.keySet().removeIf(k -> k.startsWith(instId + "_")); scheduled.keySet().removeIf(k -> k.startsWith(prefix));
completedPreMap.keySet().removeIf(k -> k.startsWith(prefix));
} }
} }

View File

@ -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));
}
}