# 工作流引擎 BSP → Dataflow 改造实施计划 > **文档目的**:将工作流执行模型从 BSP(层同步)改造为 Dataflow(数据流),使节点完成时立即触发后继节点,不再被同层慢节点阻塞。 > > **改造范围**:`backend/src/main/java/com/agent/management/engine/WorkflowLevelExecutor.java` 一个文件 > > **状态**:实施中 ## 1. 问题确认 ### 1.1 用户观察 > 一个节点,后继跟随多个并行节点时,多个并发节点目前是一起完成的。但应该是有快有慢。先完成的节点,又可以紧跟着触发它的后继节点,不必等其他并行节点一起。 ### 1.2 根因定位 `WorkflowLevelExecutor.java:162`: ```java // 等待当前层所有并行节点完成 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); ``` 当前执行模型是 **BSP(Bulk Synchronous Parallel,层同步)**: - 按 `dag.getLevels()` 逐层遍历 - 同层节点用 `nodeExecutor` 线程池并行 - **barrier**:整层必须全部完成才进入下一层 - 后继节点的调度时机由"所在层"决定,不由"前驱完成时刻"决定 ### 1.3 仓库历史 - 从未出现过 dataflow 模型 - commit `94a9f24` 仅把"层内串行"升级为"层内并行 + barrier" - 本次为全新实现 ## 2. 改造目标 ### 2.1 行为目标 | 模型 | 调度时机 | 总耗时(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)。 ### 2.2 必须保留的语义 | 编号 | 现有行为 | 保留方式 | |---|---|---| | 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) | 不动 | ## 3. 设计选型 ### 方案 A:CompletableFuture 回调链(**已选**) 每个节点用 `supplyAsync(executeOneNode, nodeExecutor)` 提交;`whenComplete` 中递减后继入度,归零则立即调度后继。 **优点**: - 真正满足"前驱完成即触发" - 改动聚焦于 `executeByLevel()` 内部 - `executeOneNode()` / 线程池 / `safeSend` 全部零改动 - 可通过 `Map` cancel 在飞 future ### 方案 B:Actor + 队列(未选) 过重,引入新并发原语。 ### 方案 C:移除 barrier 但保留 for-level(未选) 本质仍按层延迟,不满足需求。 ## 4. 实施阶段 ### 阶段 1:并发数据结构准备 **目标**:BSP 时代的 `HashSet` 升级为并发集合,新增运行时入度计数器;BSP 行为不变。 **改动点**: - `activeNodes` / `activeEdges` → `ConcurrentHashMap.newKeySet()` - 新增 `Map pendingInputs`:运行时入度 - 新增 `Map> inFlightFutures`:abort cancel 用 - 新增 `AtomicInteger sortOrderSeq`、`AtomicBoolean abortFlag`、`AtomicReference abortMsgRef` ### 阶段 2:核心调度循环重写 **目标**:删除 `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): - 条件节点 X 有 2 个分支,未选中分支的目标节点 T 入度已在阶段 1 算入 - X 完成时必须对未选中分支的所有 target 主动递减一次入度,否则 T 永远等不到 **主线程等待**: - `AtomicInteger remainingNodes`(初始 = 总节点数) - 每个节点完成(含 skipped 路径)递减,归零时 `complete(done)` - 主线程 `done.join()` ### 阶段 3:失败策略与 abort 取消 **skip 策略**:与成功一致,下游通过 envelope.status 感知失败 **abort 策略**: - `abortFlag.compareAndSet(false, true)` CAS 去重,仅首个失败节点推 `workflow_error` - 遍历 `inFlightFutures` 调用 `cancel(true)` - `scheduleNode` 与 `tryScheduleSuccessors` 入口检查 abortFlag,未启动节点不再调度 - 在飞节点完成后**仅清理 inFlightFutures**,不写 records、不推 SSE、不激活下游(避免与已发的 workflow_error 乱序) ### 阶段 4:sortOrder / finalOutputs / SSE 治理 - `sortOrder`:完成顺序,`AtomicInteger.getAndIncrement()` 在持锁区域 - `finalOutputs`:`clear() + putAll()` 在持锁区域原子执行 - `nodeRecords`:持锁 + 普通 ArrayList - `workflow_error`:阶段 3 CAS 已保证只触发一次 ### 阶段 5:单元测试 新增 `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 时在飞节点 | 在飞完成但不记录,下游不激活 | ### 阶段 6:回归验证 - 手动跑 3-5 个典型工作流,对比 finalOutputs / nodeRecords - 前端兼容性:检查 `WorkflowEditor.vue` / `RunHistory.vue` 是否依赖 sortOrder 全局顺序 - 故意失败场景验证 abort / skip - 16+ 节点并发压测,评估 `max(8, CPU*2)` 线程池容量 ## 5. 关键陷阱与对应阶段 | 陷阱 | 描述 | 处理阶段 | |---|---|---| | 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 | ## 6. 风险评估 ### 高风险(人工 review) - **H1 activeEdges 可见性**:worker 启动时必须看到刚 add 的 edge - **H2 条件分支入度补偿**:未选中分支目标节点入度必须主动递减,否则死锁 - **H3 abort 时在飞节点副作用**:完成后不能写 context / 推 SSE / 激活下游 - **H4 死锁风险**:done 加超时(30 分钟)兜底 ### 中风险(测试覆盖) - sortOrder 改完成顺序后前端展示 - finalOutputs 多 output 并发语义 - 线程池容量 ## 7. 不在本次范围内 1. ❌ 不引入 `executionModel = bsp|dataflow` 配置开关 2. ❌ 不修改 `DagResolver` 分层算法 3. ❌ 不修改 `NodeExecutor` 接口 4. ❌ 不引入新依赖(无 Akka / RxJava) 5. ❌ 不修改 `executeOneNode()` worker 逻辑 6. ❌ 不扩容线程池 7. ❌ 不改前端 ## 8. 回滚预案 - 改造集中单文件单方法 → `git revert` 原子回滚 - 无 DB schema 改动 - 建议阶段 2 与阶段 5 一起合并验证 ## 9. 复杂度估算 | 阶段 | 内容 | 估算(小时) | |---|---|---| | 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** |