文档目的:将工作流执行模型从 BSP(层同步)改造为 Dataflow(数据流),使节点完成时立即触发后继节点,不再被同层慢节点阻塞。
改造范围:
backend/src/main/java/com/agent/management/engine/WorkflowLevelExecutor.java一个文件状态:实施中
一个节点,后继跟随多个并行节点时,多个并发节点目前是一起完成的。但应该是有快有慢。先完成的节点,又可以紧跟着触发它的后继节点,不必等其他并行节点一起。
WorkflowLevelExecutor.java:162:
// 等待当前层所有并行节点完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
当前执行模型是 BSP(Bulk Synchronous Parallel,层同步):
dag.getLevels() 逐层遍历nodeExecutor 线程池并行94a9f24 仅把"层内串行"升级为"层内并行 + barrier"| 模型 | 调度时机 | 总耗时(A→B+C→D,B 5s、C 1s、D 1s) |
|---|---|---|
| BSP(当前) | 按层 | A(1s) + max(B,C)(5s) + D(1s) = 7s |
| Dataflow(目标) | 前驱完成即触发 | A(1s) + max(B, C+D)(5s) = 6s |
更关键的差别:dataflow 下,C 完成的瞬间 D 立刻开始(不等 B)。
| 编号 | 现有行为 | 保留方式 |
|---|---|---|
| S1 | executeOneNode() worker 逻辑 |
完全不动 |
| S2 | 条件分支 sourceHandle 过滤 |
在后继激活中保留 |
| S3 | NodeWorkspaceBuilder.build(..., activeEdges) 可达前驱过滤 |
activeEdges 改并发集合后语义不变 |
| S4 | output 节点的 envelope 写 finalOutputs | 改为持锁写入 |
| S5 | sortOrder 单调递增 |
改为 AtomicInteger,按完成顺序 |
| S6 | failStrategy = abort / skip 语义 | abort 时设置标志 + 取消在飞 future;skip 正常调度下游 |
| S7 | skipped 节点的 SSE 推送与记录 | 改由"入度永远到不了 0"判定 |
| S8 | SSE 线程安全(safeSend synchronized) |
不动 |
每个节点用 supplyAsync(executeOneNode, nodeExecutor) 提交;whenComplete 中递减后继入度,归零则立即调度后继。
优点:
executeByLevel() 内部executeOneNode() / 线程池 / safeSend 全部零改动Map<nodeId, Future> cancel 在飞 future过重,引入新并发原语。
本质仍按层延迟,不满足需求。
目标:BSP 时代的 HashSet 升级为并发集合,新增运行时入度计数器;BSP 行为不变。
改动点:
activeNodes / activeEdges → ConcurrentHashMap.newKeySet()Map<String, AtomicInteger> pendingInputs:运行时入度Map<String, CompletableFuture<NodeExecutionResult>> inFlightFutures:abort cancel 用AtomicInteger sortOrderSeq、AtomicBoolean abortFlag、AtomicReference<String> abortMsgRef目标:删除 for (level) + allOf().join(),改为事件驱动。
新增方法:
scheduleNode(nodeId):检查 abortFlag → 提交线程池 → 挂 whenComplete 回调onNodeCompleted(nodeId, result, ex):处理结果 → 持锁写 records/output → 触发后继tryScheduleSuccessors(nodeId, result):递减后继入度,归零则调度关键顺序(保证可见性,陷阱 8):
activeEdges.add(edge) ← 1. 先加边
pendingInputs.decrement ← 2. 再递减入度(happens-before)
if (入度 == 0) schedule ← 3. 入度归零才调度
条件分支入度补偿(陷阱 3/10):
主线程等待:
AtomicInteger remainingNodes(初始 = 总节点数)complete(done)done.join()skip 策略:与成功一致,下游通过 envelope.status 感知失败
abort 策略:
abortFlag.compareAndSet(false, true) CAS 去重,仅首个失败节点推 workflow_errorinFlightFutures 调用 cancel(true)scheduleNode 与 tryScheduleSuccessors 入口检查 abortFlag,未启动节点不再调度sortOrder:完成顺序,AtomicInteger.getAndIncrement() 在持锁区域finalOutputs:clear() + putAll() 在持锁区域原子执行nodeRecords:持锁 + 普通 ArrayListworkflow_error:阶段 3 CAS 已保证只触发一次新增 WorkflowLevelExecutorTest.java,覆盖 10 个场景:
| # | 用例 | 验证点 |
|---|---|---|
| 1 | 菱形 DAG(A→B+C→D),B 慢 C 快 | D 在 C 完成后立即调度 |
| 2 | 纯串行 A→B→C | 与 BSP 行为一致 |
| 3 | 单层 8 节点并行 | sortOrder 唯一单调 |
| 4 | 中间节点 abort | 下游不调度,workflow_error 一次 |
| 5 | 中间节点 skip | 下游通过 envelope 感知,仍执行 |
| 6 | 条件分支未选中 | 未选分支目标节点入度补偿 |
| 7 | merge 节点多入度 | 等所有 active 前驱完成 |
| 8 | 多 output 并发 | finalOutputs 原子写入 |
| 9 | 100 节点 fan-out 压测 | sortOrder 0..99 全覆盖 |
| 10 | abort 时在飞节点 | 在飞完成但不记录,下游不激活 |
WorkflowEditor.vue / RunHistory.vue 是否依赖 sortOrder 全局顺序max(8, CPU*2) 线程池容量| 陷阱 | 描述 | 处理阶段 |
|---|---|---|
| 1 | activeNodes / activeEdges 改并发集合 | 阶段 1 |
| 2 | 新增运行时入度计数器 | 阶段 1 |
| 3 | 条件分支入度补偿 | 阶段 2 |
| 4 | inFlightFutures Map 用于 cancel | 阶段 1 + 3 |
| 5 | sortOrder 全局原子递增 | 阶段 1 + 4 |
| 6 | finalOutputs 并发写入 | 阶段 4 |
| 7 | workflow_error CAS 去重 | 阶段 3 |
| 8 | activeEdges 实时可见性 | 阶段 2 |
| 9 | 线程池容量评估 | 阶段 6 |
| 10 | 条件分支 vs 普通节点入度去重对齐 | 阶段 2 |
| 11 | abort 阻止已就绪下游调度 | 阶段 3 |
executionModel = bsp|dataflow 配置开关DagResolver 分层算法NodeExecutor 接口executeOneNode() worker 逻辑git revert 原子回滚| 阶段 | 内容 | 估算(小时) |
|---|---|---|
| 1 | 并发数据结构准备 | 1.5 - 2 |
| 2 | 核心调度循环重写 | 3 - 4 |
| 3 | 失败策略与 abort 取消 | 2.5 - 3 |
| 4 | sortOrder / finalOutputs / SSE 治理 | 1.5 - 2 |
| 5 | 单元测试(10 用例) | 3 - 4 |
| 6 | 回归验证 | 2 - 3 |
| 合计 | 13.5 - 18 |