把一个中等复杂度的 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 .venv&&source.venv/bin/activate pipinstall-Upydantic openai tenacitynodes/放节点:一个节点一个文件,只依赖上游传入的字典,不偷偷读全局状态。runs/放产物:每次执行一个子目录,断点续跑与事后审计都靠它。logs/放日志:每行一条 JSON,方便后续按trace_id聚合。
**核心结论:**目录即流水线的物理形态,节点、产物、日志三分离是后续一切能力的地基。
四、接口契约:用 Pydantic 把节点边界焊死
流水线最怕接口含糊,先用数据模型把每个节点的进出字段写死:
frompydanticimportBaseModel,FieldclassDraftIn(BaseModel):outline:list[str]=Field(min_length=1)audience:str="general"classDraftOut(BaseModel):content:str=Field(min_length=50)citations:list[str]=[]defdraft_node(state:dict)->DraftOut:raw=call_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:"""替换为你的真实模型调用,返回结构化 dict"""return{"node":node,"ok":True,"data":payload}defrun_pipeline(task:str,run_id:str="r1")->dict:state,done={"task":task,"done":[]},[]os.makedirs(f"runs/{run_id}",exist_ok=True)fornodeinNODES:state[node]=call_model(node,{"input":task,"up":state.get("done")})done.append(node)state["done"]=done# 每个节点跑完立即落盘,这就是断点续跑的凭据json.dump(state,open(f"runs/{run_id}/{len(done)}_{node}.json","w"),ensure_ascii=False,indent=2)returnstate- 单向数据流:节点只读
state里编号更小的产物,禁止回头改写,避免隐式耦合。 - 即跑即落盘:不要等全链跑完才写文件,中途断电也能看出停在哪一步。
强约束生成 → 事实校验 → 低分校验自动回退重写。
六、断点续跑:用中间产物把长任务变成可恢复任务
任务跑一半挂掉时,靠中间产物落盘就能从断点接着跑,不必从头再来:
- 产物即检查点:每节点一个 JSON 文件,文件名里的序号就是天然的执行顺序。
- 启动先扫盘:重新执行时先列出
runs/{run_id}里已有的文件,跳过已完成节点。 - 节点需幂等:同一输入重复执行必须得到等价结果,否则续跑会写出不一致状态。
- 契约校验产物:续跑前用
model_validate复读文件,损坏的检查点直接作废重跑。
问题:长任务(如万字报告)单次失败重来的 token 成本可能占到整条链的一半。
治理:断点续跑把重跑范围缩到单节点,失败代价从「整链」降为「一格」。
**核心结论:**能落盘的状态才是真状态,内存里的中间结果在生产环境等于不存在。
七、失败处理:重试、降级与死信队列怎么分工
先给错误分类,再决定重试、降级还是丢进死信队列:
| 现象 | 可能原因 | 处理动作 |
|---|---|---|
| 节点超时 | 上游限流或网络抖动 | 指数退避重试 3 次 |
| 输出非 JSON | 模型格式漂移 | 附格式示例后重发 1 次 |
| 校验分过低 | 上游上下文缺失 | 回退上一节点重跑 |
| 连续失败 3 次 | 任务本身超能力 | 写入死信队列并告警 |
重试逻辑用一个装饰器实现,指数退避避免把限流打成雪崩:
importtimefromfunctoolsimportwrapsdefretry(times=3,base=1.0):defdeco(fn):@wraps(fn)defwrap(*args,**kwargs):foriinrange(times):try:returnfn(*args,**kwargs)exceptException:ifi==times-1:raise# 交给死信队列time.sleep(base*2**i)returnwrapreturndeco@retry(times=3,base=1.0)deffact_check_node(state:dict)->dict:...**核心结论:**重试只对瞬时错误有效,逻辑错误必须靠回退上一节点解决,两者不可混用。
八、并发切分:在串行骨架里安全地偷并行
串行不等于慢,阶段之间的顺序不能变,但阶段内部可以扇出:
fromconcurrent.futuresimportThreadPoolExecutordeffan_out(items:list[dict],worker,workers=6)->list[dict]:"""批内并发,阶段间仍严格串行"""withThreadPoolExecutor(max_workers=workers)aspool:returnlist(pool.map(worker,items))# 结果顺序与输入一致- 只并叶子节点:IO 密集的检索、摘要可以并发,涉及全局状态的节点保持串行。
- 控制扇出宽度:
workers建议先压在 4~8,超出上游限流阈值只会换来 429。 - 保留原始顺序:
map按输入序返回,合并结果时无需再排序对齐。
**核心结论:**串行骨架 + 叶子节点扇出,能在不破坏可预测性的前提下把耗时压下去。
九、可观测:一次执行要能被完整回答和复盘
节点一多,先要有能回答「哪段慢、哪段错、花了多少」的观测面:
| 字段 | 含义 | 排查用途 |
|---|---|---|
trace_id | 整条流水线唯一 ID | 串起一次完整执行的日志 |
node/attempt | 节点名与第几次尝试 | 定位高频失败与重复重试 |
latency_ms | 单节点耗时 | 找出整链瓶颈所在 |
tokens_in/out | 输入与输出 token | 按节点核算成本 |
status | ok / 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 模型、五节点实现、重试装饰器与扇出函数都是可直接复制的最小片段,接上你自己的模型调用就能跑。
当你需要扩展路由、回环或人工审批时,也不必推倒重来——先让串行链稳定出结果,再在个别节点外挂分支即可。
先用串行把确定性跑稳,再谈动态调度。