设计讲解:Planner 已产出并行信息,但 Dispatcher 仍串行
一句话:要解决什么问题?
现在 planner 已经算出"哪些任务可以并行"(用group字段标记),但 dispatcher 还是一个一个串行跑——等于买了并行车道还在走单行道。加上没有重试、没有幂等、错误不分轻重,导致三个痛点:
- 多意图时 P95 线性叠加(串行跑 3 个 Agent = 3× 单次耗时)
- 瞬态错误(503、超时)和永久错误(403、Agent 不存在)一视同仁,浪费宝贵的重规划配额
- 重试时每次发新请求,下游无法去重,可能重复下单
核心改动 5 块,用一张运行时态图串起来
用户说:"帮我查食堂菜单,同时查会议室" │ ▼ ┌─── planner ───┐ │ step_1: 查食堂菜单 step_2: 查会议室 │ group=1 group=1 ← 同组=可并行 │ depends_on=[] depends_on=[] │ idempotency_key= │ "t1-r0-canteen-0" "t1-r0-meeting-1" │ max_group = 1 ← 写入 state └───────────────┘ │ ▼ ┌─ dispatch_step ─┐ │ current_group=1 │ │ 取 group==1 的全部 step │ call_with_retry(step_1) call_with_retry(step_2) ← asyncio.gather 并行 │ 某步失败? │ retryable=True? → 等待 → 重试(复用同一幂等键) │ retryable=False? → 直接返回 ok=False │ 两个 DispatchResult 汇入 dispatch_results │ completed_groups.append(1); current_group = 2 └──────────────────┘ │ ▼ current_group(2) > max_group(1) → 去 reflector ┌─ reflector ─┐ │ 全成功? → done → respond │ 有失败? → replan(差量:只重做失败步) └─────────────┘第 1 块:group 并行(最核心改动)
现状:一个一个跑,即使有两口灶空闲。
step=plan[step_index]result=awaitdispatch(step)step_index+=1# 条件边: step_index < len(plan) → 继续循环改后:组内并行,组间串行(组间有依赖)。
batch=[sforsinplanifs["group"]==current_group]results=awaitasyncio.gather(*[dispatch_with_retry(s)forsinbatch],return_exceptions=True)current_group+=1# 条件边: current_group <= max_group → 继续循环具体例子:
用户:"先查食堂菜单,再根据菜单查排队,同时查会议室" planner 产出: step_1: 查食堂菜单 group=1, depends_on=[] step_2: 查会议室 group=1, depends_on=[] ← 和 step_1 无依赖,同组并行 step_3: 查排队 group=2, depends_on=[step_1] ← 依赖 step_1,必须等 group 1 完成 执行: 第 1 轮:step_1 和 step_2 并行(asyncio.gather) 第 2 轮:step_3 执行(group 1 全完成后才到)安全细节:
return_exceptions=True:某步炸了不影响同组其他步max_parallel=5:防过载截断,超了 log warning 只跑前 5 个step_index字段废弃但不删(保序列化兼容),设为 -1
第 2 块:retryable 三路径解析
为什么需要:之前DispatchResult.error只是字符串"http_503",reflector 看到ok=False不知道能不能重试——503 可重试,403 重试一万次也白搭。
三条路按优先级走:
错误来了 │ ┌─ 路径1: error_detail.retryable 已有?─┐ │ (领域 Agent 在 A2A 响应里显式标注) │ │ 有 → 直接用 │ │ 无 ↓ │ ├─ 路径2: 从 error_detail.code 推导 ────┤ │ TIMEOUT / SERVER_ERROR → True │ │ FORBIDDEN / NOT_FOUND / NO_AGENT │ │ / PARAM_ERROR → False │ │ 无法推导 ↓ │ └─ 路径3: 保守兜底 → False ─────────────┘ (宁可不重试,也不浪费配额)错误码映射表:
| 场景 | code | retryable | 人的直觉 |
|---|---|---|---|
| 下游超时 | TIMEOUT | ✅ True | 等一会儿可能就好了 |
| 下游 503 | SERVER_ERROR | ✅ True | 服务暂时不可用,可能恢复 |
| 鉴权失败 403 | FORBIDDEN | ❌ False | 权限不够,重试没用 |
| Agent 不存在 | NO_AGENT | ❌ False | 不存在的东西再试也不存在 |
| A2A task.state=failed | SERVER_ERROR | ✅ True | 保守——可能下游 LLM 超时 |
| A2A task.state=rejected | PARAM_ERROR | ❌ False | 下游主动拒绝,参数有问题 |
第 3 块:步骤内重试 + 指数退避
现状:call()单次调用,失败就返回ok=False给 reflector 决定重规划整个 plan。
改后:call_with_retry包一层重试循环,不改变call()单次语义(保护现有测试)。
attempt 0: call() → ok=False + retryable=True → 等 1.0^0 = 1s attempt 1: call() → ok=False + retryable=True → 等 1.0^1 = 1s attempt 2: call() → ok=False + retryable=True → 重试耗尽,返回最后结果 如果 retry_backoff=2.0: attempt 0 → 等 1s attempt 1 → 等 2s attempt 2 → 等 4s关键区分(评审必问):重试和重规划是两个不同层级的重试
- 重试(
call_with_retry):同一个 step,同一个幂等键,微观重试,最多 2 次 - 重规划(reflector → planner):不同 plan,可能换策略,宏观重试,最多 3 次
第 4 块:幂等键(Idempotency Key)
为什么需要:用户点"订餐",请求发出去但超时了。重试——但下游已收到并处理了第一次请求。没有幂等键 → 下游执行两次 → 重复下单。
怎么工作:
planner 生成: idempotency_key = "{thread_id}-{replan_round}-{agent}-{idx}" 例如: "thread-abc-0-canteen-0" 透传链路: PlanStep.idempotency_key → _dispatch_step 读取 → dispatch_with_retry(idempotency_key=...) → call_with_retry(同一个 key 传给每次重试) → call(idempotency_key=...) → headers["Idempotency-Key"] = "thread-abc-0-canteen-0" → A2A 请求发出重试时复用同一个 key → 下游看到同一个 key → “这个处理过了,返回上次结果”。
Fail-Safe:如果idempotency_key为空(边界情况),不发空 header,而是发sup-{random-uuid}兜底,记 warning——宁可幂等保护失效,也不阻断调用。
第 5 块:配置化
现状:httpx.AsyncClient(timeout=120)硬编码 120 秒;没有重试;发现缓存 TTL 硬编码 10 秒。
改后(config/supervisor.yaml):
dispatch:a2a_timeout:"${DISPATCH_A2A_TIMEOUT:30}"# 单次 A2A 超时,从 120 降到 30a2a_retries:"${DISPATCH_A2A_RETRIES:2}"# 重试次数retry_backoff:"${DISPATCH_RETRY_BACKOFF:1.0}"# 退避基数discover_ttl:"${DISPATCH_DISCOVER_TTL:30}"# Agent 发现缓存 TTLmax_parallel:"${DISPATCH_MAX_PARALLEL:5}"# 单组最大并发idempotency_hours:24# 幂等键有效期(文档约束)注意a2a_timeout从 120 秒降到 30 秒——因为有了重试兜底,单次可以更激进地超时,整体反而更快(120s 卡死 vs 30s 超时 + 1s 等 + 30s 重试 = 最多 61s)。
数据模型变更一图览
DispatchError (新增) DispatchResult (增强) ┌─────────────────────┐ ┌─────────────────────────────┐ │ code: str │ │ agent, skill, group, ok │ ← 不变 │ retryable: bool │──────────────▶│ error: str │ ← 保留兼容 │ trace_id: str │ │ error_detail: DispatchError │ ← 新增 │ message: str │ │ retryable: bool │ ← 新增(快捷访问) └─────────────────────┘ │ idempotency_key: str|None │ ← 新增 │ step_id, round │ ← 已有 │ content, card_* │ ← 不变 └─────────────────────────────┘ PlanStep (增强) ┌─────────────────────────────────┐ │ agent, skill, message │ ← 不变 │ step_id, depends_on, group, slots│ ← 已有 │ idempotency_key: str │ ← 新增 └─────────────────────────────────┘图条件边变更(最关键的接线点)
# 旧:步进驱动g.add_conditional_edges("dispatch_step",lambdas:"dispatch_step"ifs["step_index"]<len(s["plan"])else"reflector")# 新:group 驱动g.add_conditional_edges("dispatch_step",lambdas:"dispatch_step"ifs["current_group"]<=s["max_group"]else"reflector")不做什么(评审时可能被问到的边界)
| 不做 | 原因 |
|---|---|
| Agent 白名单扩展 | 当前只开 canteen + transport,其他未部署,开了也会 404 |
| Reflector clarify 移除 / degraded 标记 | 独立设计,别混在一起 |
| call_parallel() 签名变更 | 那是给 agent 模式用的,workflow 走自己的 _dispatch_step 并行 |
| HITL 写操作确认 | 超出本次范围 |
| 补事务 undo | 超出 supervisor 职责 |
| 分批 gather | 99%+ 场景 group 内 ≤5 步,截断 + warning 够用 |
评审时重点关注的 5 个问题
return_exceptions=True但某 step 的 Exception 不是 DispatchResult——需在 gather 后把 Exception 包裹成DispatchResult(ok=False, error_detail=...)。- max_parallel 截断丢了 step——首版 log warning,但丢了 step 意味着
dispatch_results不完整,reflector 会判 incomplete → replan,是否可接受? - a2a_timeout 从 120→30——现有领域 Agent 的 P99 是否稳定在 30s 以内?如果不,大量超时 + 重试反而更慢。
- 幂等键格式
{thread_id}-{replan_round}-{agent}-{idx}——重规划时 replan_round 递增,同一 step 会换 key,下游幂等表是否有 TTL 导致旧 key 过期后的重复执行风险? - step_index 废弃但保留字段——checkpoint 里可能有旧格式 state,反序列化时 step_index 存在但不被使用,是否影响 replay?
需要我进D:\qed对照这份设计,看哪几块已经落地、哪几块还完全没做吗?这样评审前你能拿到"设计 vs 现状"的真实偏差。