恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
多 Agent 串行流水线:把一个任务拆成可重试、可续跑的 Pipeline 节点链
首页
资讯中心
/
多 Agent 串行流水线:把一个任务拆成可重试、可续跑的 Pipeline 节点链
多 Agent 串行流水线:把一个任务拆成可重试、可续跑的 Pipeline 节点链
发布时间:2026/10/8 18:52:20
把一个中等复杂度的 Agent 任务比如「调研 → 大纲 → 撰写 → 校验 → 发布」全塞进一个 Prompt越写越长的提示词很快失控中间状态无法校验、某一步出错只能整段重跑出了问题连是哪个环节的锅都定位不到。常见解法是继续加约束——在 Prompt 里写「请分步骤思考」或者依赖模型自带的 ReAct 循环。但这类做法的问题在于步骤隐含在自然语言里没有明确的输入输出边界也就无法单独重试、单独计费、单独观测上下文一长前面的指令还会被后面的内容冲淡。本文分享一套可直接落地的串行流水线方案数据契约Pydantic Schema 断点续跑每节点产物落盘 重试降级与批内并发节点化封装后可直接搬进任何 Python 项目。一、范式选择先判断任务值不值得拆节点不是所有任务都该上多 Agent先看三个特征再决定步骤可枚举任务能写出固定的先后顺序不需要运行时动态规划路径。每步可验收某一步的产出有明确的对错标准例如格式合法、字段齐全、引用存在。失败可局部重跑某一步失败时前面的产出仍可复用不必从零开始。问题三项里只要有一项不满足比如路径要动态决策硬套串行只会更脆。治理三项都满足时串行流水线是心智成本最低、日志最清晰的多 Agent 架构。**核心结论**顺序稳定 局部可验收才值得拆成节点链否则先考虑路由或图式结构。二、模式对比四种协作拓扑各自的适用边界把主流协作范式并排看一遍就能判断什么时候该选串行 Pipeline模式拓扑结构适合场景主要风险单 Agent 长 Prompt一段式指令流步骤少、上下文短指令互相打架、无中间校验串行 Pipeline节点 A → B → C 固定顺序步骤稳定、每步可验收单点阻塞导致全链等待路由分发分类器选分支输入类型混杂误分类后难回溯图式多 Agent节点带条件边与回环需要人工介入与反复修订调试与状态管理成本高串行的代价任一节点卡住整条链就停在原地所以必须配超时与重试。串行的收益数据流单向、状态可落盘任何一次执行都能被完整复现。**核心结论**串行 Pipeline 的本质是用「放弃动态调度」换取「可预测与可审计」。三、开工准备目录骨架与依赖一次装好动手前先把依赖、目录结构和运行约定一次配齐后面每个节点都能独立跑mkdir-ppipeline/{nodes,runs,logs}cdpipeline python-mvenv .venvsource.venv/bin/activate pipinstall-Upydantic openai tenacitynodes/放节点一个节点一个文件只依赖上游传入的字典不偷偷读全局状态。runs/放产物每次执行一个子目录断点续跑与事后审计都靠它。logs/放日志每行一条 JSON方便后续按trace_id聚合。**核心结论**目录即流水线的物理形态节点、产物、日志三分离是后续一切能力的地基。四、接口契约用 Pydantic 把节点边界焊死流水线最怕接口含糊先用数据模型把每个节点的进出字段写死frompydanticimportBaseModel,FieldclassDraftIn(BaseModel):outline:list[str]Field(min_length1)audience:strgeneralclassDraftOut(BaseModel):content:strField(min_length50)citations:list[str][]defdraft_node(state:dict)-DraftOut:rawcall_model(draft,DraftIn.model_validate(state).model_dump())returnDraftOut.model_validate(raw)# 不合法就地抛错进参先校验上游少给字段时在入口报错而不是在几层调用之后炸出KeyError。出参必结构化下游拿到的永远是 Pydantic 对象模型的格式漂移被挡在节点内。契约即文档DraftIn/DraftOut直接生成接口说明新人看模型就懂数据流。**核心结论**契约先行坏数据在节点边界就被拦下这是串行链能稳定跑的前提。五、最小实现五节点串联的可运行流水线下面这条五段式流水线可直接运行节点顺序与数据流向一目了然importjson,os NODES[research,outline,draft,fact_check,publish]defcall_model(node:str,payload:dict)-dict:替换为你的真实模型调用返回结构化 dictreturn{node:node,ok:True,data:payload}defrun_pipeline(task:str,run_id:strr1)-dict:state,done{task:task,done:[]},[]os.makedirs(fruns/{run_id},exist_okTrue)fornodeinNODES:state[node]call_model(node,{input:task,up:state.get(done)})done.append(node)state[done]done# 每个节点跑完立即落盘这就是断点续跑的凭据json.dump(state,open(fruns/{run_id}/{len(done)}_{node}.json,w),ensure_asciiFalse,indent2)returnstate单向数据流节点只读state里编号更小的产物禁止回头改写避免隐式耦合。即跑即落盘不要等全链跑完才写文件中途断电也能看出停在哪一步。强约束生成 → 事实校验 → 低分校验自动回退重写。六、断点续跑用中间产物把长任务变成可恢复任务任务跑一半挂掉时靠中间产物落盘就能从断点接着跑不必从头再来产物即检查点每节点一个 JSON 文件文件名里的序号就是天然的执行顺序。启动先扫盘重新执行时先列出runs/{run_id}里已有的文件跳过已完成节点。节点需幂等同一输入重复执行必须得到等价结果否则续跑会写出不一致状态。契约校验产物续跑前用model_validate复读文件损坏的检查点直接作废重跑。问题长任务如万字报告单次失败重来的 token 成本可能占到整条链的一半。治理断点续跑把重跑范围缩到单节点失败代价从「整链」降为「一格」。**核心结论**能落盘的状态才是真状态内存里的中间结果在生产环境等于不存在。七、失败处理重试、降级与死信队列怎么分工先给错误分类再决定重试、降级还是丢进死信队列现象可能原因处理动作节点超时上游限流或网络抖动指数退避重试 3 次输出非 JSON模型格式漂移附格式示例后重发 1 次校验分过低上游上下文缺失回退上一节点重跑连续失败 3 次任务本身超能力写入死信队列并告警重试逻辑用一个装饰器实现指数退避避免把限流打成雪崩importtimefromfunctoolsimportwrapsdefretry(times3,base1.0):defdeco(fn):wraps(fn)defwrap(*args,**kwargs):foriinrange(times):try:returnfn(*args,**kwargs)exceptException:ifitimes-1:raise# 交给死信队列time.sleep(base*2**i)returnwrapreturndecoretry(times3,base1.0)deffact_check_node(state:dict)-dict:...**核心结论**重试只对瞬时错误有效逻辑错误必须靠回退上一节点解决两者不可混用。八、并发切分在串行骨架里安全地偷并行串行不等于慢阶段之间的顺序不能变但阶段内部可以扇出fromconcurrent.futuresimportThreadPoolExecutordeffan_out(items:list[dict],worker,workers6)-list[dict]:批内并发阶段间仍严格串行withThreadPoolExecutor(max_workersworkers)aspool:returnlist(pool.map(worker,items))# 结果顺序与输入一致只并叶子节点IO 密集的检索、摘要可以并发涉及全局状态的节点保持串行。控制扇出宽度workers建议先压在 4~8超出上游限流阈值只会换来 429。保留原始顺序map按输入序返回合并结果时无需再排序对齐。**核心结论**串行骨架 叶子节点扇出能在不破坏可预测性的前提下把耗时压下去。九、可观测一次执行要能被完整回答和复盘节点一多先要有能回答「哪段慢、哪段错、花了多少」的观测面字段含义排查用途trace_id整条流水线唯一 ID串起一次完整执行的日志node/attempt节点名与第几次尝试定位高频失败与重复重试latency_ms单节点耗时找出整链瓶颈所在tokens_in/out输入与输出 token按节点核算成本statusok / retry / dead区分正常、重试与死信每节点一条日志json.dumps输出到logs/一行一个事件grep trace_id即可复原现场。产物与日志对齐日志里的seq与runs/下的文件序号一致审计时可互相印证。**核心结论**观测不是加分项没有trace_id的流水线在生产环境等于黑盒。十、上线清单发版前逐项核对的检查表上线前按这份清单过一遍缺一条都别发版契约齐全每个节点的In/Out模型都已定义并通过model_validate。落盘与续跑任一节点中断后重跑能跳过已完成步骤并保持结果一致。重试有上限所有外部调用都带超时与次数上限超限进死信而非无限循环。限流有节流并发扇出数低于上游 QPS 上限压测下无 429 雪崩。日志可检索任取一个trace_id能在 5 分钟内复原整条链的执行过程。成本可观测单任务 token 成本有统计超阈值能自动告警。**核心结论**清单的价值在于把「架构对了」变成「生产上不出事」。结语串行 Pipeline 看起来是最朴素的多 Agent 架构但它把最难的三件事——接口边界、失败恢复、过程审计——用最直白的方式解决了节点契约挡住脏数据产物落盘让长任务可续跑日志与trace_id让每次执行都可复现。上面给出的目录骨架、Pydantic 模型、五节点实现、重试装饰器与扇出函数都是可直接复制的最小片段接上你自己的模型调用就能跑。当你需要扩展路由、回环或人工审批时也不必推倒重来——先让串行链稳定出结果再在个别节点外挂分支即可。先用串行把确定性跑稳再谈动态调度。