TL;DR: Send 是 LangGraph 的 map-reduce 原语:条件边返回一列 Send(节点, 输入),运行时为每项起一个并行任务。用它的三个纪律:reducer 合并顺序不保证,去重排序放到后续节点;并行区里的人审要走 id 映射;fan-out 数量要自己设上限。纪律背后是同一句话:并行把延迟压下来,把确定性收走,收走的部分要你自己补回去。
场景:500 份简历打分
一批文档、每个独立打分、最后汇总排名。这类需求一眼看过去是纯体力活,先按直觉写串行:循环 500 次,每次调一次模型打分。跑起来才发现慢得不能接受,假设单次打分 3 秒,500 份就是 25 分钟,用户等不起,定时窗口也未必够。而这 500 次调用彼此完全独立,谁也不依赖谁的输出,串行纯粹是代码形状造成的,不是业务约束。先把这个「彼此独立」确认清楚,是整篇的前提:如果第 100 份的打分要参考第 99 份的结果,那根本不是并行问题,是流程设计问题,先回去改流程。
并行的写法用 Send。条件边不再返回「下一个节点叫什么」,而是返回一列 Send,每一项指定「跑哪个节点、喂什么输入」,运行时为每项起一个并行任务。
动笔之前先校准一下预期:并行省的是时间,不是钱。500 次模型调用还是 500 次,token 账单一分不少;并行只是让它们同时发生。所以「并行了为什么成本没降」不是 bug,是这门生意的基本盘。另外提速也不是 500 倍,接口限流、共享配额、各任务耗时的参差都会把实际加速比压下来,经验上接近「并发槽位数 × 最慢任务的耗时」这个量级,槽位数又受下一节要讲的限流约束。预期校准了,再写代码。
import operator
from typing import Annotated, TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
class S(TypedDict):
docs: list[str]
scores: Annotated[list[int], operator.add] # 各分支结果往这里合并
def split(state: S) -> list[Send]:
# 条件边不返回节点名,返回 Send 列表:每项一个并行任务
return [Send("score", {"doc": d}) for d in state["docs"]]
def score(state: dict) -> dict:
# 注意:这里的 state 不是全局 S,是 Send 给的局部输入
return {"scores": [rate(state["doc"])]}
def merge(state: S) -> dict:
ranked = sorted(state["scores"], reverse=True)
return {"scores": []} # 汇总后的落库逻辑略
g = StateGraph(S)
g.add_node("score", score)
g.add_conditional_edges(START, split, ["score"])
g.add_edge("score", "merge")
g.add_edge("merge", END)
app = g.compile()
这段代码有四个读起来容易漏的细节。第一,split 挂在 add_conditional_edges 上,它本质是个条件边函数,只是返回值从节点名变成了 Send 列表。第二,每个 Send 自带输入,score 拿到的 state 不是全局的 S,是 Send 里给的那份局部字典,所以它的类型标注可以只是 dict,里面只有 doc,没有 docs。第三,score 的返回值是增量,{"scores": [x]} 是「往 scores 里加一个元素」的意思,前提是 scores 这个 key 声明了 operator.add 作为 reducer。第四,merge 不需要任何特殊标记,所有并行任务结束之后它自然被调用,等它跑的时候 scores 已经是合并完的完整列表。
运行时收到 500 个 Send,就同时跑 500 个 score 任务。每个任务的输入是 Send 里给的那份局部 state,不是全局状态。全部任务跑完(一个 superstep 结束),500 份增量 {"scores": [x]} 一起交给 operator.add 合并。map-reduce 的两个阶段对应得非常整齐:Send 是 map,reducer 是 reduce。
这个对应不只是修辞,它是排障时的分区依据。结果错了但每个任务的日志都对,问题在 reduce 侧:reducer 声明错了、顺序假设错了、去重没做。某个任务的输出从头就不对,问题在 map 侧:局部输入给错了、任务函数本身的逻辑问题。先分区再深挖,比在 500 份交错日志里大海捞针快一个量级。判断的依据也很明确:把 merge 节点收到的输入直接打印出来,如果这批增量本身就是错的,是 map 的锅;增量都对、合并结果不对,是 reduce 的锅。
superstep:并行真正发生的单位
要预判并行的行为,得先知道它在什么单位上发生。LangGraph 的执行按 superstep 推进:一个 superstep 里,运行时把当前要跑的所有任务(普通节点激活、或一整批 Send 任务)全部启动,然后等这一批全部结束,才把各自的增量合并进 state,进入下一个 superstep。这个「批次加栅栏」的模型决定了两件事。
第一,merge 的等待是天然的,不需要你计数。500 个任务哪怕先后的完成时间参差,栅栏都会等最后一个,然后一次性合并。你不需要 semaphore,不需要 future 列表,框架的调度模型替你做了同步。
第二,一批之内没有先后,一批之间有先后。如果 merge 之后又有条件边触发新一轮 Send,那是下一个 superstep 的事,用到的是已经合并完的 state。想在「部分任务完成」时就动状态,做不到,栅栏就是栅栏;这个限制反过来是安全保障:你看到的 state 永远是某个完整批次的合并结果,不会是五百个任务写了一半的中间态。还要补一句任务集合的确定时机:条件边函数在一个 superstep 开始时求值,返回的 Send 列表一次性定死这一批的任务集合,中途不会追加。也就是说「跑到一半看情况多发几个」这种动态需求,表达不了在一个 superstep 内,只能拆成多个 superstep,让每轮的条件边看到上一轮合并完的 state 再决定发多少。
recursion_limit 的账也记在这个单位上,默认 25 个 superstep。500 个 Send 只算一个 superstep,不因为任务多而多扣;但 map 之后如果还有层层嵌套的图(每个并行任务自己又是一个多步图),每一步都在消耗同一个预算。并行嵌套时遇到 GRAPH_RECURSION_LIMIT 报错,先数 superstep,别数任务。
为什么必须有栅栏?因为合并这个动作需要一个确定的输入集合。reducer 的语义是「把这一批的增量合并起来」,如果任务可以边跑边合并,reducer 每次看到的输入集合都不一样,合并结果依赖合并时机,也就不可复现。栅栏把「这一批」变成了一个确定不移的集合,同样的输入跑两遍,每个 superstep 的合并结果都一样。这是并行之下还能保留确定性的全部基础,代价是你要接受批内无进展可见。
用并行的心智读别人的 Send 代码时,可以带着这份栅栏意识去问三个问题:这一批的任务集合由什么决定(条件边在哪个 state 上求值)、每个任务的输入从哪来(Send 的局部字典)、合并后顺序是否有人假设了什么。三问过完,这段代码在生产会怎么表现,基本能推出来。
栅栏还有一个和第 03 篇直接相关的推论:快照存放在 superstep 边界上。一批任务跑到一半进程挂了,恢复时是从上一个边界开始重跑,这一批里已经成功的任务也会再执行一遍。所以并行任务的副作用必须幂等,这条和第 06 篇坑五是同一条纪律在并行场景的形态:打分可以重打(覆盖同一条记录),发钱不能重发。把「这批任务重跑一遍会怎样」当成并行设计的第一道审查题。
合并秩序:三件必须知道的事
顺序不保证。 500 个并行任务谁先完成不受你控制,合并进列表的顺序和提交顺序无关,甚至两次运行结果都不一样。这个不确定性会从最意想不到的地方漏出来:报告按 scores 列表顺序生成分页,页和页的归属乱掉;「取前三名发奖」取错了人,因为并列分数时谁在前谁在后是随机的;测试时断言列表内容用 == 比较整个列表,时绿时红。需要稳定顺序就在 merge 节点里排序,或者让 Send 的输入带上序号、结果里带回去:
def split(state: S) -> list[Send]:
return [
Send("score", {"doc": d, "idx": i}) # 输入带序号
for i, d in enumerate(state["docs"])
]
def score(state: dict) -> dict:
return {"scores": [{"idx": state["idx"], "v": rate(state["doc"])}]}
merge 里按 idx 排回去,顺序就稳定了。用 Python 的 sorted 加 key 函数就够,它本身是稳定排序,并列分数时保持既有相对次序,再要严格确定就给 key 里加第二排序字段(比如文档 id),彻底消除随机性。别在任何依赖「写入顺序」的逻辑上做假设,这类假设在测试环境(任务少、机器闲)往往成立,上了生产并行度一高就翻车。
写入冲突看 reducer。 并行任务写同一个 key,key 必须挂了 reducer,否则就是第 02 篇讲过的覆盖语义:同一个 superstep 里多个任务写同一个无 reducer 的 key,最后完成的那个把别人的全擦掉。这是并行场景最经典的数据丢失事故,而且不报错。并行场景的 state 设计先问一句:这个 key 的 reducer 是什么?答不上来就别并行写它。还有一类更隐蔽的:key 挂了 reducer,但任务之间有隐性依赖,比如两个任务都会「往列表里加同一个元素」(上游数据重复时),operator.add 来者不拒,合并结果里出现重复项。去重不能指望 reducer,放到 merge 节点里做,这是 TL;DR 里「去重排序放到后续节点」的完整含义。
reducer 的选择本身有个小清单。列表累加用 operator.add(对 list 是拼接),计数累加也是它(对 int 是加法),字典合并要用自定义函数(lambda a, b: {**a, **b} 之类),只取第一个成功结果的场景用「非空则跳过」的自定义函数。最常见的类型陷阱是同一个 key 一开始只有 int,后来某个任务开始返回 list,int + list 直接抛异常,而且抛在合并那一瞬间,堆栈指向 reducer 而不是写错的任务,排查时格外绕。类型注解 Annotated[list[int], operator.add] 里的期望类型是给人看的约定,运行时不会强制,所以约定破了要到合并时才炸。给 reducer 旁边写清期望类型、merge 里对元素结构做断言,是两道便宜的保险。
fan-out 没有免费的无限并发。 Send 列表长度没有硬性上限,但每个任务都要占运行时资源,LLM 接口有自己的限流。500 个 Send 直接怼上去,多半换来一堆 429,重试又放大流量,最坏情况是雪崩。实用做法两选一:把 fan-out 粒度调粗,一次 Send 处理一批文档,500 份变 50 批,并发压力降一个数量级;或者把限流做进任务内部,任务启动先拿配额,拿不到就等,把并发度控制在接口承受范围内。批大小怎么估有个朴素办法:接口给你的并发配额除以单任务平均耗时,得到单位时间能稳定处理的任务数,再用任务总数除它,得到预期完成时间的数量级。拿开篇的数字算一遍:假设接口并发配额是 60,单次打分 3 秒,500 个任务分 60 个槽跑,8 轮多一点,25 秒上下,和串行的 25 分钟比是 60 倍。但如果配额只有 5,500 个任务要跑 5 分钟,而你要为 5 路并发各写好重试和限流,值不值就得算算。这个数不必精确,但要有,它同时是 SLA 承诺和限流设计的依据。重试和补偿的完整设计在第 12 篇,那里和这里的组合拳是生产并行的标配。
并行区里的人审
并行分支里也有 interrupt 时(比如 50 份审批同时挂起),第 06 篇坑六的问题被放大:恢复值必须和 interrupt 一一对上,而并行的完成顺序本身就不稳定,「按顺序第几个」彻底不可用。
做法是给每个 interrupt 的 payload 带全局唯一 id(文档 id 加动作类型,比如 "resume-0231:score_review"),恢复时传一个 {id: 决策} 的映射,每个分支节点按自己的 id 取值。这个模式官方文档有对应说明,属于并行人审的标准解法。第 06 篇的五问检查表在这里全部适用,还要多问一句:并行区里挂起的 interrupt 数量有没有上限,五十个同时挂起是设计,五千个同时挂起是事故。映射模式还天然支持分批恢复:审批人先批了 30 份,resume 字典里就只有 30 个键,没拿到决策的分支继续挂着等,这个语义是顺的,但 merge 要等全部齐还是凑批处理,得是你明确的设计决定。
前端呈现也要跟着变。五十个待审不该渲染成五十个弹窗,而是一个聚合审批页:列表、勾选、一次提交,提交后端把整张 {id: 决策} 映射一次 resume 进去。交互形态和第 06 篇坑三的结论互相印证:一次收齐决策永远好过逐项打断。超时兜底同样按批设计,扫描器发现某批待审超过 SLA,resume 的映射里给没批的项带上 timeout: True 的默认拒绝,整批流程得以收尾,而不是永远悬着几个键。
数量固定的并行,不需要 Send
只想要固定数量的并行分支(比如两个评审互相独立打分),直接从同一个节点拉两条边就够了,两条边指向同一个节点名或不同节点,运行时自然并行,同样受 superstep 栅栏管,同样要求 reducer。不需要 Send。
Send 的存在价值是数量在运行时才知道:文档数、候选数、待审数是数据决定的。数量固定用静态边,数量动态用 Send,这个界线很清楚。两种方式放在一起看:
| 静态并行边 | Send | |
|---|---|---|
| 分支数量 | 写代码时确定 | 运行时由数据决定 |
| 每个分支的输入 | 共享全局 state | 每任务一份局部输入 |
| 典型场景 | 两三个独立评审 | 批量打分、批量审批 |
| 调试复杂度 | 低 | 随数量上升 |
中间还有个灰色用法值得一提:Send 不仅携带数量,还携带输入,所以它也能当「运行时路由」用。同一批文档里,大文件 Send 给分段处理的节点,小文件 Send 给直处理的节点,输入里带上各自需要的字段。这已经超出 map-reduce 的教科书用法,但它就是同一个原语,规则不变:局部输入、增量返回、reducer 合并。局部输入还带来一个隔离保证:任务之间不共享那份输入字典,一个任务改自己的 dict 不会影响别的任务,不需要担心并行改同一份可变对象的老问题。
失败与部分成功的设计
并行把延迟从「任务数 × 单任务耗时」压到接近「单任务耗时」,收益直观。成本在结果确定性上,而且集中爆发的位置是失败处理:500 个里 3 个失败,整体算成功还是失败?
这个问题没有默认答案,必须你在设计时选一个并写进 merge 逻辑。常见的三种策略:一是全有或全无,任何失败就整体重跑,简单但浪费,500 份里重跑 497 份成功的;二是记录部分成功,state 里加一个 failed key(同样要 reducer),失败的文档带原因进列表,merge 里按「成功率达标就出报告,附失败清单」处理,多数生产场景选这个;三是原地重试,失败的任务在下一个 superstep 里重新 Send 一次,配合幂等设计(打分调用可重入)能把成功率补上来。三种都要打分调用本身幂等做底,不然重试又会造出重复项。
重试还要和限流联动着设计,不然会互相打架:429 高发时,重试恰好撞在限流最紧的时刻,越重试越 429。所以重试要有预算(最多几次、指数退避、加抖动),并且失败率本身要当信号用:一批任务里失败率突然从 1% 跳到 40%,说明接口端在限流或故障,这时候正确的动作是整批降速,而不是让 200 个任务各自傻傻地重试。第 12 篇的错误分类里,「该重试的错误」和「该熔断的错误」是两类,并发放大了把后者当前者处理的风险。
把这些决定写下来,就是并行流程的设计文档:并发上限多少、失败算什么、重试几次、merge 按什么规则放行。并行不是加个参数,是把「一致性处理」从框架手里接过来自己写,写下来的那份决定,就是你的系统文档里最值钱的一页。
部分成功策略再往下推一层是阈值和死信。阈值指「成功率达标才出结果」:500 份里失败 3 份照常出报告,失败 300 份就该整体失败报警,阈值定在业务能接受的最低质量线上。死信指彻底失败的那部分去哪:进 state 里的失败清单只是第一步,还要有出口,进重试队列、进人工处理列表,或者至少发一条告警。没有出口的失败清单会安静地躺在 state 里,每次运行都复制一份,直到有人发现报告里的「失败 37 份」从来没人处理过。阈值和死信都不是框架概念,是你的业务决定,但它们必须有明确的主人,可以是 merge 节点里的一行判断,可以是跑批结束后的一个检查任务,唯独不能是「没人想过」。
并行图的观测与调试
并行图出问题时,最大的障碍是日志交错:50 个任务的输出混在一起,谁是谁分不清。三个便宜的手段能解决大部分痛苦。第一,任务内日志带 id:Send 输入里的那个唯一 id,在每个任务的每条日志里都带上,交错日志立刻变得可检索。第二,merge 节点做数量断言:进 merge 时先检查 len(scores) 是否等于预期任务数,不等就报错,这个断言能当场抓住「部分任务静默没写结果」的问题,比下游发现报告缺数据早得多。第三,单任务可重放:把单个 Send 的输入存下来,就能在图外单独复跑这一个任务,调试时不用把 500 个全跑一遍。第 14 篇的轨迹观测方法在并行图里全部适用,只是所有轨迹都要按任务 id 分组看。
预算观测要单说一句。并行最容易烧穿 token 和费用预算,因为量大而且同时发生:串行时 500 次调用按时间摊开,账单异常有时间窗能发现;并行时一次 superstep 就可能烧掉平时一天的量。按任务记账(每个任务的 token 消耗记进它的返回值或日志)、给单次运行设总预算、merge 前核账,三件事在并行场景从「锦上添花」变成「必须」。第 17 篇算成本账时,并行的用量曲线是单独的一类,规划容量时要按它的峰值而不是均值。
逐批发送:把 superstep 当节拍用
上一节说一个 superstep 的任务集合开始时定死,中途不能追加。这听起来像限制,其实是工具:想要「受控的动态并发」,把批次拆开,让每一轮 superstep 发送下一批。做法是 state 里维护一个游标,条件边每轮只发固定数量的 Send,merge 之后条件边看游标决定还有没有下一批:
def split(state: S) -> list[Send]:
start = state["cursor"]
batch = state["docs"][start:start + 50] # 每轮固定 50 个任务
return [Send("score", {"doc": d}) for d in batch]
def advance(state: S) -> dict:
return {"cursor": state["cursor"] + 50}
def more(state: S) -> str:
return "split" if state["cursor"] < len(state["docs"]) else END
批次走图结构的循环,限流天然被批大小控制,每一批结束都是一个 checkpoint 边界,挂了能从这一批的开头恢复,不用整批五百个重来。代码里的 S 要加一个 cursor: int,它每轮只有一个写入者(advance 节点),不用挂 reducer。代价是延迟从「一批全并发」变成「多批流水」,这是拿时间换可控性的典型交换。任务量不大时用不上,任务上千、接口配额紧张时,这就是标准写法。
这套写法还有一个衍生好处:逐批天然适合断点续跑的运营节奏。白天跑一半要停,晚上接着跑,游标停在哪,进度就在哪;要不要临时缩小批次(接口告警时从 50 降到 10),改一个数字就行。控制粒度到了批次这一级,运行一个大规模批处理就不再是「祈祷五百个任务一起顺利」,而是一段一段地推着走。
什么时候不该并行
并行不是默认优化,有三种情况并行是负资产。第一,任务之间其实有顺序依赖:打分要先分类、分类要先抽取结构,这种链式依赖硬拆成并行,就得在任务之间加等待、加重试、加中间状态,最后做出一个手写的串行流程,套着并行的壳。依赖在,就承认它,串行链或分层图(一层并行、层间串行)都更干净。第二,任务本身太碎:500 个 Send 每个只干 100 毫秒的活,调度和合并的开销占比过高,把 Send 粒度调粗(一个任务处理一批)几乎总是更好,上一节讲过批大小,这里补它的下限理由。第三,下游就等着汇总成一份东西、且没有时间压力:并行的收益是延迟,没有延迟压力就不要付确定性的成本。定时任务跑一小时无所谓,那就串行,代码简单一半,排查也简单一半。
判断的根子和第 01 篇一样是任务形状:任务彼此独立、数量动态、时间敏感,三项都占,Send 是对的工具;占得越少,越该考虑串行或粗粒度。还有一条隐性的自查:如果你发现自己在 merge 节点里写了大量等待、补齐、顺序修复的逻辑,这些逻辑的总量就是「这个需求其实想串行」的证据,代码在告诉你它的真实形状。
与子图叠加:并行的平方复杂度
Send 的任务里再跑子图(每个并行任务是一个完整的多步图),可行,规则也都成立,但要清楚复杂度是相乘的。快照:每个任务各自有 namespace、各自的快照链(第 07 篇的分层模型直接套用),一个 thread 里一次运行就可能产生几百份快照。流式:开了 subgraphs=True 之后事件量按任务数翻倍,前端过滤要按「层级加任务 id」两维做。递归预算:每个任务的每一步都消耗同一个 recursion_limit,预算要按「批数 × 单任务步数」重新算。人审:多个任务各自 interrupt,恢复走第 06 篇坑六的 id 映射,id 里最好连任务序号都带上。
这套叠加不是不能用,是每一项成本都要有意识地接受。经验法则是:先让 Send 跑普通函数节点跑通,再逐个把任务升级成子图,每升一级把观测和预算检查跟上。一步到位写「并行加嵌套加人审」的图,出问题时三层一起排查,定位成本远高于分步搭建。真到那一步,也是先看 map 侧还是 reduce 侧(分区方法在前面),再看任务内部,最后才看两层交界处。排查路径有层次,复杂度才压得住。
评审检查单:并行上线前五问
把全篇纪律收敛成五问,评审并行方案时逐条过。这五问的顺序也是设计顺序:先定 state 的合并规则,再定批次与并发,最后定失败语义,顺着答基本不用回头返工:
| 检查项 | 抓的问题 |
|---|---|
| 每个并行写的 key 有没有 reducer | 覆盖语义静默丢数据 |
| 顺序假设消灭了吗(排序、去重在 merge) | 结果顺序随机、重复项 |
| fan-out 上限和批大小是多少 | 429、雪崩、超预算 |
| 部分失败算什么(阈值、死信出口) | 报告缺数据没人发现 |
| 任务重跑一遍会怎样(幂等) | 快照恢复后副作用翻倍 |
五问的前两问是正确性,中间两问是容量设计,最后一问是和持久化机制的对接。都能答上,Send 的收益就归你;答不上哪问,问题就埋在哪问里。第 09 篇的多 agent 是并行的进一步组织化:Send 并行的是无状态任务,多 agent 并行的是有角色的执行体,下一篇讨论什么时候值得迈那一步。
延伸阅读
- Graph API 文档:Send 与条件边的官方定义
- Use Subgraphs 指南:map-reduce 里的并行与嵌套关系
- GRAPH_RECURSION_LIMIT 错误说明:并行嵌套时步数预算