workflow-dataflow-execution-plan.md 8.4 KB

工作流引擎 BSP → Dataflow 改造实施计划

文档目的:将工作流执行模型从 BSP(层同步)改造为 Dataflow(数据流),使节点完成时立即触发后继节点,不再被同层慢节点阻塞。

改造范围backend/src/main/java/com/agent/management/engine/WorkflowLevelExecutor.java 一个文件

状态:实施中

1. 问题确认

1.1 用户观察

一个节点,后继跟随多个并行节点时,多个并发节点目前是一起完成的。但应该是有快有慢。先完成的节点,又可以紧跟着触发它的后继节点,不必等其他并行节点一起。

1.2 根因定位

WorkflowLevelExecutor.java:162

// 等待当前层所有并行节点完成
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<nodeId, Future> cancel 在飞 future

方案 B:Actor + 队列(未选)

过重,引入新并发原语。

方案 C:移除 barrier 但保留 for-level(未选)

本质仍按层延迟,不满足需求。

4. 实施阶段

阶段 1:并发数据结构准备

目标:BSP 时代的 HashSet 升级为并发集合,新增运行时入度计数器;BSP 行为不变。

改动点

  • activeNodes / activeEdgesConcurrentHashMap.newKeySet()
  • 新增 Map<String, AtomicInteger> pendingInputs:运行时入度
  • 新增 Map<String, CompletableFuture<NodeExecutionResult>> inFlightFutures:abort cancel 用
  • 新增 AtomicInteger sortOrderSeqAtomicBoolean abortFlagAtomicReference<String> 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)
  • scheduleNodetryScheduleSuccessors 入口检查 abortFlag,未启动节点不再调度
  • 在飞节点完成后仅清理 inFlightFutures,不写 records、不推 SSE、不激活下游(避免与已发的 workflow_error 乱序)

阶段 4:sortOrder / finalOutputs / SSE 治理

  • sortOrder:完成顺序,AtomicInteger.getAndIncrement() 在持锁区域
  • finalOutputsclear() + 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