KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
06. 并行任务怎样汇聚而不丢失因果关系? — keel 龙骨
事件诊断需要同时检查指标和最近部署。两个任务只依赖同一份事件摘要,彼此不依赖,因此可以 fan-out 并行执行,完成后再 fan-in 汇聚。
事件诊断需要同时检查指标和最近部署。两个任务只依赖同一份事件摘要,彼此不依赖,因此可以 fan-out 并行执行,完成后再 fan-in 汇聚。
先证明任务真的独立
如果 Change Agent 必须先读取 Metrics Agent 的结论,它们就不是并行任务。可以用依赖图判断:
flowchart LR
I[Incident Context] --> M[Metrics Task]
I --> D[Deployment Task]
M --> G[Gather]
D --> G
G --> R[Review Task]
只有同一层且没有相互依赖的节点可以安全并行。把顺序依赖强行并行,会让 Agent 使用不完整或旧数据。
并行不是“同时 print”
一个完整并行协议需要:
- 为每个子任务分配唯一
task_id; - 记录它属于哪个
run_id和父任务; - 限制最大并行数;
- 独立记录开始、完成、失败和取消;
- 按 task ID 汇聚,而不是按返回顺序猜测;
- 处理重复、缺失和迟到结果;
- 根据 required/optional 策略决定最终状态。
模型输出完成顺序不稳定,下面的代码是错误的:
# 错误:假设第一个返回的一定是 metrics。
metrics_result, deployment_result = completion_order
正确做法是显式关联:
results_by_task[result.task_id] = result
Gather 不只是字符串拼接
Gather 需要先做完整性检查:
期望任务:metrics(required), deployment(required), owner(optional)
收到结果:metrics(succeeded), deployment(failed), owner(succeeded)
若策略是 all_required,运行不能声称完整成功;若策略是 best_effort,可以生成降级报告,但必须明确缺失 deployment 证据。
常见汇聚策略:
| 策略 | 行为 | 适用情况 |
|---|---|---|
| all_required | 任一必需任务失败,汇聚失败 | 合规检查、必须覆盖全部来源 |
| best_effort | 使用成功结果并标记缺失 | 辅助分析、允许降级 |
| quorum | 达到指定数量或权重才继续 | 多次独立判断 |
| first_valid | 第一个通过验收的结果即可 | 多供应商 fallback |
最大并行数是稳定性边界
同时创建 100 个 Agent 不一定更快。本地模型可能串行加载,外部 API 可能限流,工具服务也可能被打满。
待执行任务 -> bounded queue -> N 个 worker slot -> results
max_parallelism、每 Agent max_concurrency 和服务限流共同形成背压。上游必须在容量不足时排队或拒绝,不能无限创建任务。
取消和迟到结果
用户取消 TeamRun 后,运行时应:
- 标记 run 为 cancelling/cancelled;
- 停止分派尚未开始的任务;
- 向可取消的子任务传播 cancellation token;
- 接收但隔离已经无法阻止的迟到结果;
- 禁止迟到结果把 cancelled 改回 completed。
线程、HTTP 请求和外部任务的取消能力不同。取消是协作协议,不是保证底层计算瞬间停止。
运行并行实验
python courses/advanced/multi-agent-collaboration/course/project/examples/03_parallel_team.py
示例故意让两个 Worker 使用不同延迟。观察结果是否仍按 task_id 关联,以及事件序号是否反映真实完成顺序。
然后让一个必需 Worker 抛出异常,分别使用 all_required 和 best_effort。最终状态应该不同,但成功任务的 Artifact 都应保留用于诊断。
并行中的重复执行
网络重试或 worker lease 过期可能让同一 task attempt 被执行两次。对于只读分析,可以使用唯一 attempt 接受第一个有效结果;涉及副作用时,必须结合 Tool Calling 的幂等键和业务唯一约束。
Agent 层的 task_id 不能替代工具层的 idempotency key,两者保护不同边界。
检查理解
- 为什么返回顺序不能代表任务类型?
best_effort为什么不等于忽略错误?- max parallelism 解决的是什么问题?
- 取消后为什么仍可能收到结果?
并行结果要进入同一次运行,下一步必须分清哪些数据可以共享、哪些只能引用。