ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

dispacher支持多任务编排(串行/并行执行/高可用)设计

dispacher支持多任务编排(串行/并行执行/高可用)设计 下面用「角色比喻 大白话」帮你梳理汇报时直接照着讲就行。一、先认识 4 个角色贯穿全文角色干什么类比Planner规划员把用户问题拆成多个子任务按依赖关系分组项目经理排任务Dispatcher派单员把子任务分发给各个子 Agent派活的组长子 Agent真正干活的出行 Agent、食堂 Agent…一线员工Reflector复盘员看执行结果决定继续 / 重规划 / 收工质检员用户问一句帮我导航到公司顺便看看食堂今天有啥→ Planner 拆成 2 个子任务 → Dispatcher 派给两个子 Agent → Reflector 汇总结果。二、现在有什么问题4 个痛点痛点 1能一起干的活却排队串行干 最影响性能现状Planner 已经把任务分好组group了但 Dispatcher无视分组一个一个排队派。后果3 件互不相干的事总耗时 三件事耗时之和。用户等得久P95 飙高。讽刺的是代码里其实已经写好了并行能力call_parallel用了asyncio.gather但从来没人调用它——是个死代码。痛点 2失败后不分青红皂白一律重规划 最浪费资源现状子 Agent 失败后Reflector 只看成功/失败这一个信号不分原因一律触发「重新规划整个计划」。问题网络抖动、服务临时挂了 →再试一次就好重规划整个计划是杀鸡用牛刀会议室已满、余额不足 →再试 100 次也是满的重规划纯属浪费配额和时间事实平台契约里子 Agent已经会返回retryable能否重试字段但我们的代码根本没去读它。大白话版这段在解决一个核心问题——子 Agent 调用失败了中枢 supervisor 该怎么反应是重试、是直接放弃、还是硬着头皮继续背景为什么需要这个supervisor 把任务派给子 Agentcanteen/transport 等。子 Agent 偶尔会失败服务挂了、会议室已满、超时。失败分两种性质处理方式必须不同可重试的失败比如网络抖、服务 5xx、超时——等会儿再试可能就成功了不可重试的失败比如会议室已被占用“参数非法”——再试 100 次也是同样结果纯浪费旧逻辑大概是一刀切失败就 replan重新规划重试但这样遇到会议室已满这种死局会白费重试配额遇到服务挂了又可能一次就放弃。改动 1子 Agent 返回失败时先搞清楚能不能重试sub_provider.py里解析子 Agent 的失败响应按优先级三条路径判断路径1最优先看 JSON-RPC 错误信封——子 Agent 明确回了{error: {data: {retryable: true/false}}}。子 Agent 自己最清楚能不能重试听它的。默认retryableFalse保守不轻易重试。路径2看任务状态——如果 task 状态是failed/error/rejected再试着从 metadata 里读retryable标注。路径3兜底看 HTTP 状态码——没前面那些信息就自己猜5xx/超时 → 可重试4xx → 业务错误不可重试。每条路径最终都给失败结果打上一个retryable标签能重试 / 不能重试。改动 2reflector 拿到所有结果后做差量决策reflector 节点收集这一轮所有子 Agent 的返回分成两堆不可重试的失败堆4xx、会议室已满等这种不重试直接判定done_degraded——任务算完成但降级给用户一个部分服务不可用的提示卡片把能拿到的部分结果拼上不再浪费 replan 配额。可重试的失败堆5xx、超时这种才走replan利用 replan 配额做差量补偿只重跑失败的跳过已成功的。改动 3降级卡片长什么样当 reflector 判done_degraded时doneTrue直接结束不 replan回复里追加一句部分服务暂时不可用把那些不可重试失败的原文比如会议室已满拼进回复让用户知道哪部分没办成不重复派发已经失败的内容一句话总结给失败分了类能重试的服务抖→ 重新规划再试不能重试的业务死局→ 直接降级交付不瞎重试。判定依据优先听子 Agent 自己说其次看任务状态最后看 HTTP 码兜底。这样既不浪费重试配额也不会对死局傻等。需要我确认这段代码在sub_provider.py/reflector.py里现在是不是已经落地还是只是方案痛点 3写操作重试可能重复执行 有业务风险场景用户说帮我订会议室第一次请求超时了系统重试 → 下游收到两次独立的请求→ 可能订了两个会议室。现状代码每次都发一个新的随机 UUID当幂等键形同虚设。契约要求A2A 请求头要带Idempotency-Key幂等键重试时必须复用同一个键下游才能认出来这是同一个请求别重复处理。幂等键类比就像银行转账的流水号。同一个流水号重复提交银行只处理一次。痛点 4参数写死在代码里改一次要发一次版现状HTTP 超时120 秒、重试次数、并发数、发现缓存 TTL —— 全部硬编码在代码里。后果线上想调超时时间必须改代码 → 提 PR → 评审 → 发版。运维没法根据现场情况快速调整。三、改什么4 件事一一对应上面 4 个痛点改动 1按组并行派活 → 解决慢现在任务1 → 任务2 → 任务3 总耗时 T1 T2 T3 改后┌ 任务1 ┐ ├ 任务2 ┤ 同时派出 总耗时 ≈ max(T1, T2, T3) └ 任务3 ┘规则组内并行、组间串行。同一个 group 的任务 → 互不依赖同时派用asyncio.gather不同 group 的任务 → 后面的可能依赖前面的结果必须按顺序来保护机制单任务失败不影响同组其他任务return_exceptionsTrue失败的那个转成失败结果继续走流程新增max_parallel默认 5控制并发上限比如一组 10 个任务就分 3 批55执行防止一下子打爆下游顺带激活 3 个字段current_group/completed_groups这两个字段早就定义了但从来没用过死字段现在真正启用step_index因为不再需要而废弃保留字段只为兼容旧数据。改动 2读懂能不能重试 → 解决瞎重规划让代码去读子 Agent 返回的retryable字段按三条路径判断判断路径场景结论① JSON-RPC 错误包里的retryable子 Agent 明确标注了以它说的为准② 任务状态是 failed/error/rejected状态里带了 metadata 标注读标注③ HTTP 状态码兜底子 Agent 没标的情况5xx / 超时 → 可重试4xx 业务错误 → 不可重试推导Reflector 拿到这个信号后的两种决策失败是不可重试的会议室已满 → 直接 done_degraded降级收工 → 给用户看部分服务暂时不可用 已成功的部分结果 → 不浪费重规划配额 ✅ 失败是可重试的网络抖动/5xx → 走 replan重新规划差量补偿✅收益不可重试的错误立刻收口不再无脑重规划省时间、省 LLM token。改动 3幂等键 → 解决重复执行生成Planner 在生成每个任务时就分配一个幂等键格式{会话ID}-{重规划轮次}-{子Agent名}-{序号} 例如thread-abc-0-canteen-agent-2透传链路一路传下去不丢Planner 生成 idempotency_key ↓ dispatch_step 从任务里读出来 ↓ dispatch() → SubAgentProvider.call() ↓ A2A 请求头 Idempotency-Key: thread-abc-0-canteen-agent-2重试复用重试时用的是同一个键下游 Agent 认键去重 →订会议室重试 3 次也只会订 1 个✅兜底万一幂等键生成失败仍然发一个sup-{随机UUID}保证调用不中断只记 warning不阻断。改动 4配置化 → 解决改参数要发版新增配置文件段dispatch配置项默认值作用a2a_timeout30s单次调用超时原来硬编码 120sa2a_retries1可重试失败的重试次数retry_backoff1.0s重试退避等 1s、2s…线性退避max_parallel5单组最大并发discover_ttl10s服务发现缓存有效期收益运维改配置即可调优不用改代码、不用发版a2a_timeout和max_parallel还支持Nacos 热更新改完立即生效不用重启。四、领导最关心的验收指标指标目标怎么测编排成功率≥90%端到端评测集含并行 差量重试场景端到端 P95≤3s多意图场景并行 vs 串行压测对比并行正确性组内并发、组间串行、结果不丢单元测试重试判断准确可重试→重规划不可重试→降级单元测试写操作幂等带幂等键重试复用同一键单元测试断言请求头重试次数受控可重试配1次 → 最多调2次不可重试 → 不重试单元测试配置生效改 yaml → 行为变化单元测试五、领导一定会问出问题怎么办风险兜底方案并行后某个任务挂了该任务转失败结果同组其他任务不受影响流程继续子 Agent 没标retryable默认不可重试保守策略不浪费配额幂等键生成/透传失败发随机 UUID 兜底调用不中断并行把下游打爆max_parallel分批限流整个方案出问题git revert一条命令回到串行版本代码直接替换无隐藏开关核心原则任何外部依赖服务发现 / A2A 调用 / 鉴权 / 监控失败都转成结构化失败结果填回去绝不让异常中断整个流程。六、上线节奏分 4 步走每步独立可验证打地基数据模型加字段retryable/idempotency_key 新增配置模型改核心子 Agent 调用层解析 retryable 透传幂等键 重试逻辑改编排并行派活 重试决策 Planner 生成幂等键 工厂接线补测试 验证10 类单测全绿pytest/ruff/mypy三轨通过 →小流量灰度先跑多意图场景对比串行 vs 并行的 P95 和成功率七、汇报时可以直接念的 3 分钟版本这个设计是完善我们中枢 Agent 的任务分发链路解决 4 个问题第一慢。之前多个子任务明明互不相干却排队串行执行总耗时是叠加的。而且并行能力代码其实早就写好了只是一直没被调用。现在改成能并行的一起派、有依赖的按顺序来耗时从累加变成取最慢那个。第二笨。子任务失败后系统不分原因一律重新规划整个计划。但网络抖动和会议室已满是两回事——前者再试一次就好后者试一百次也没用。现在系统会读子 Agent 返回的可重试标记可重试的走重试不可重试的直接给用户降级提示不再浪费资源。第三有重复执行风险。像订会议室这类写操作重试时下游会当成两次新请求。现在给每个任务分配唯一幂等键重试复用同一个键下游认键去重重复请求只执行一次。第四改参数要发版。超时、重试次数、并发数原来写死在代码里现在提到配置文件运维随时可调关键参数还支持热更新。目标是编排成功率 ≥90%、端到端 P95 ≤3s。整个方案是纯内部优化对外接口完全不变出问题一条git revert就能回到旧版。八、领导可能追问 怎么答Q这个改动动了多少代码风险大吗涉及 10 个源文件但都是内部编排逻辑对外 HTTP 接口一个都没变对下游子 Agent 的 A2A 契约也没变只是新增了幂等键的正确用法。回滚只需一条命令。Q并行会不会把下游打爆有max_parallel默认 5限流而且任务分组本身就限制了并发度——只有 Planner 判定互不依赖的任务才会进同一组并行。Q为什么现在才做之前为什么串行Planner 的分组能力是后来才补上的ac00dd3当时 Dispatcher 没跟上导致分组信息一直被闲置。这次是补齐这块拼图。Q做完能达到 3 秒吗得看具体场景。单意图场景本来就不慢收益主要在多意图场景——原来串行叠加现在取最慢那个。我们会先小流量灰度用真实数据对比串行 vs 并行的 P95 再全量。Q这次不做什么明确不做人工确认环节HITL、补偿事务回滚、入站 Redis 持久化。这几块是后续独立的议题不在本次范围。需要我帮你做成PPT 大纲每页一句话 一图或者一页纸汇报摘要吗这两个写法控制的是同一个东西dispatch_step节点执行完后图是继续循环回dispatch_step还是走向reflector结束分发。区别在于怎么判断还有没有待处理的 step。旧写法lambdas:dispatch_stepifs.get(step_index,0)len(s.get(plan,[]))elsereflector用索引计数器step_index判断已处理步数 计划总步数 → 继续否则结束。隐含假设plan 是线性平铺的step 按顺序逐个执行step_index从 0 一直 1 到len(plan)。致命前提step_index必须在每次dispatch_step后正确自增且 plan 里没有分组/并行/跳过概念。一旦 plan 带 group同组并行、或某些 step 被跳过step_index的线性递增就和实际执行对不上。新写法lambdas:dispatch_stepifs.get(current_group,1)_max_group(s)elsereflector# _max_group(s) max(step[group] for step in s[plan]) if s.get(plan) else 0用分组编号current_group判断当前组号 ≤ 计划里最大组号 → 继续否则结束。隐含假设plan 是按group分层的同group的 step 是一批可并行/同轮current_group每轮 1 走到最大组号为止。好处不依赖总共几步被执行了只依赖现在在第几组、最大第几组。即使某组里有 step 被跳过/并行组号推进逻辑不受影响不会因 step 计数错位而提前终止或死循环。核心差异一句话旧新判断维度步数索引step_indexvs 总步数组号current_groupvs 最大组号适用结构线性顺序 plan分组/并行 plan风险点plan 非纯线性时step_index错位 → 误终止或漏执行依赖 plan 每项有group字段且current_group每轮正确 1新写法更稳健的前提是你的 plan 确实是按group组织的每个 step 带group整数且dispatch_step每轮把current_group推到下一个还存在的组号。如果 plan 其实仍是线性、没有 group 概念那新写法需要保证每个 step 都有group且连续编号否则_max_group算出来的上限和实际轮数对不上可能少跑或多跑。如果你贴一下dispatch_step节点里current_group/step_index是怎么更新的我能判断这两个写法在你的实际流程里会不会行为分叉。
返回列表