ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

LangGraph动态分支与并行Map-Reduce:Command与Send实战

LangGraph动态分支与并行Map-Reduce:Command与Send实战 第一次把 LangGraph 用在真实业务上的人大概都经历过同一个阶段图能跑通但一遇到分支数量要等到运行时才知道的场景就卡住了。静态边只能表达A 之后走 B 或者 C一旦你需要把一份文档切成 37 块、每块独立打分再合并或者审核轮次取决于模型输出静态图就开始变得又长又丑。LangGraph 给出的答案是两个原语Command负责动态控制流Send负责运行时扇出两者配合就是标准的并行 Map-Reduce 骨架。这篇文章围绕这三件事展开——Command 怎么用、Send 怎么把固定分支变成动态分支、以及我在它们身上踩过的那几个坑。如果你已经会写add_edge和add_conditional_edges想往前走一步这篇内容就是给你准备的如果你还没上手也能看懂因为关键概念我都会先讲清楚再上代码。1. 分支数量由运行时决定时静态边为什么一定会失控1.1 一个几乎每个人都会遇到的失控场景先看个具体例子。假设你要做一个长文摘要流水线把文章按 500 字切块每块送进模型提取要点最后把所有要点汇总成一份结构化摘要。切块数量取决于文章长度可能是 3 块也可能是 120 块。用最朴素的静态图写法你得这样想如果只有一块怎么办两块怎么办这种思路根本没法写下去因为你不可能为 120 种情况各画一条边。有人会退一步把处理逻辑塞进一个大节点里用for循环串行跑完所有块。这样确实能跑但代价是失去了并发、失去了每个块独立的失败重试、失去了断点续跑而且任何一个块抛异常整条流水线全挂。这就是静态图的边界。add_edge和add_conditional_edges表达的是从一个节点到有限几个候选节点它们无法表达从一个输入生成 N 个并行分支。你需要的不是更多边而是一个能在运行时构造分支的原语。1.2 Command 和 Send 各自守哪一块地盘很多教程把这两个东西放在一起讲导致初学者分不清什么时候用哪个。我的经验是记一句话Command 管去哪Send 管开几路。原语核心作用从哪里返回典型场景Command一次返回同时完成状态更新与路由跳转节点函数的返回值意图路由、审核回环、子图跳回父图Send运行时决定并行分支数量每路带私有状态条件边的路由函数或节点返回值Map-Reduce、批量打分、分块处理interruptCommand(resume)暂停图并在外部输入后继续节点内调用 / 恢复时传入人工审核、敏感操作确认一个重要区别Command是替换了原来的边语义——返回它之后节点后面挂的普通边就不生效了而Send是补充扇出它通常配合一个归并节点形成 barrier。理解这一点后面很多报错你就能自己推断出原因。提示不要把Send理解成多线程调用一个函数。它是图调度层面的扇出每个分支是独立的 superstep 任务共享同一份状态通道但拿到的是各自的私有输入。2. 环境与版本基线不少玄学报错其实是版本太旧2.1 依赖安装与版本对照Command、Send、interrupt这几个原语在不同版本里的导入路径和可用性差别不小我在两个项目里就遇到过同一份代码在 A 环境能跑、B 环境报 ImportError的情况。先把基线对齐能省掉大量排查时间。python -m venv .venv source .venv/bin/activate pip install -U langgraph langchain-core python -c import langgraph; print(langgraph.__version__)几个判断依据可以对着自查Send和Command都在langgraph.types里如果你的版本里from langgraph.types import Send报错基本就是版本太旧。interrupt和Command(resume...)属于较新的能力依赖 checkpoint 机制必须给图配checkpointer。RetryPolicy可能从langgraph.types或langgraph.graph导出拿不准就先在交互环境里dir()一下比翻文档快。Python 版本建议 3.10 以上。原因不是 LangGraph 强制而是Annotated配合TypedDict写 reducer 时3.9 的类型语法会让你写得很难受而且很多第三方库的新版本已经不再支持 3.9。2.2 一个最小可运行骨架在动Command和Send之前先把下面这个骨架跑通后面所有实验都在它上面改from typing import Annotated, TypedDict import operator from langgraph.graph import StateGraph, START, END from langgraph.checkpoint.memory import InMemorySaver class State(TypedDict): text: str items: Annotated[list[str], operator.add] def split(state: State) - dict: return {items: [state[text][i:i 200] for i in range(0, len(state[text]), 200)]} def merge(state: State) - dict: return {text: | .join(state[items])} builder StateGraph(State) builder.add_node(split, split) builder.add_node(merge, merge) builder.add_edge(START, split) builder.add_edge(split, merge) builder.add_edge(merge, END) app builder.compile(checkpointerInMemorySaver()) print(app.invoke({text: x * 900}, {configurable: {thread_id: t1}}))这里有两处值得留意。第一items字段用了Annotated[list[str], operator.add]这是 reducer 的声明方式现在看起来多余但等Send扇出之后它就是唯一的收口机制。第二compile时带了checkpointer并传了thread_id——只要你后面想用interrupt或者断点续跑这两个东西缺一不可提前养成习惯比事后补要好。3. Command把改状态和改路由压成一次返回3.1 Command 的字段构成与执行时机Command的完整形态大致是这样from typing import Literal from langgraph.types import Command def router(state: State) - Command[Literal[fast_path, slow_path, need_human]]: score len(state[text]) if score 100: return Command(update{items: [short]}, gotofast_path) if score 5000: return Command(gotoslow_path) return Command(gotoneed_human)四个参数各管一摊update要写回状态通道的增量写法和你平时节点返回 dict 完全一样reducer 照样生效。goto目标节点名可以是字符串也接受节点名组成的列表列表意味着并行跳到多个节点。graph指定跳转的目标图Command.PARENT用来从子图跳到父图不写默认在当前图内找。resume配合interrupt使用不是路由参数恢复执行时传。Command只能在节点函数的返回值位置使用。这一点是硬约束我见过不少人试图把它塞进条件边的路由函数里结果得到的是一个被当成字典处理的对象行为完全不对。3.2 它比先改状态再走条件边紧凑在哪不用Command的实现方式通常是这样节点里先返回状态更新再挂一个条件边根据更新后的状态决定去哪。这要求路由函数必须能读到刚写入的状态也就是先写后判两段式。Command把这两段压成一段好处有两个。一是逻辑内聚判定条件和它引发的状态变更写在一起Review 的时候一眼能看清不用在两个函数之间来回跳。二是避免了一类隐蔽 bug——如果你用了 reducer 且两个分支都往同一个字段追加内容两段式写法下条件边可能读到的是部分合并的中间态而Command的update和目标节点是同一批被调度的语义上干净很多。当然也有取舍。Command的路由逻辑对图结构是隐式的你在get_graph().draw_ascii()里看不到那条边。所以我遵守一条规则路由分支超过 5 个、或者需要团队共享时仍然用条件边 一张显式的路由表只有强耦合、分支少、状态变更和路由判断必须一起做的场景才上Command。3.3 人工审核回环interrupt 配 Command(resume)这是Command最有价值的一个用法。传统做法是让人工审核这一步把状态存到数据库再起一个新图靠业务代码把两边粘起来。现在可以直接在图内暂停from langgraph.types import interrupt, Command def review(state: State) - Command[Literal[publish, rewrite]]: decision interrupt({draft: state[text], question: 是否通过审核}) if decision.get(approved): return Command(gotopublish, update{items: [approved]}) return Command(gotorewrite, update{items: [decision.get(comment, )]}) app builder.compile(checkpointerInMemorySaver()) config {configurable: {thread_id: review-001}} first app.invoke({text: 草稿内容}, config) # 此时图停在 review 节点first 里能拿到 interrupt 抛出的载荷 second app.invoke(Command(resume{approved: False, comment: 第二段太笼统}), config)几个实操要点。interrupt是把当前节点的执行状态整个挂起恢复时会从这个节点开头重新执行所以节点里interrupt之前不能有副作用重的操作比如发邮件、写数据库否则恢复时会再来一遍。真需要副作用把它挪到interrupt之后的节点或者用幂等键兜住。另外thread_id是恢复的锚点同一个thread_id才能续上不同thread_id就是另一条独立的执行链。3.4 子图跳父图graphCommand.PARENT子图里想直接跳到父图的某个节点写Command(goto父图节点名, graphCommand.PARENT)。我踩过的坑是跳转目标名必须在父图里真实存在写错的话报错信息指向的是父图的编译期而不是你写错的这一行排查时容易被误导。习惯做法是把父图节点名提成模块级常量两边引用同一个常量。4. Send 与 Map-Reduce把扇出交给运行时4.1 Send 的私有状态语义Send的构造是Send(节点名, 传给该节点的输入)。关键在于第二个参数它是这个分支的私有状态不是主状态的补丁。也就是说当Send(score, {idx: 0, chunk: ...})被调度时score节点收到的 state 就是{idx: 0, chunk: ...}主状态里的其他字段在它眼里根本不存在。这个设计和很多人的直觉相反——直觉上会觉得分支能读到全局状态只是自己那份输出会被合并回去。实际上读的不是全局写回去的才是通过 reducer 合并。为什么要这样设计因为并行分支之间必须隔离。如果每个分支都能看到完整的可变状态那 reducer 合并时到底是哪个版本说了算就说不清了还会出现读到别的分支半成品的脏读。私有输入 统一收口是这个模型能保证结果确定性的前提。4.2 reducer 是扇入的唯一合法收口扇出之后N 个分支的返回值要写回主状态。如果目标字段没有 reducerLangGraph 会直接拒绝因为一个字段在一个 superstep 里被写了 N 次这件事必须有明确的合并规则。声明方式就是在TypedDict上用Annotatedfrom typing import Annotated import operator class State(TypedDict): chunks: list[str] scores: Annotated[list[tuple[int, float]], operator.add] counter: Annotated[int, operator.add]operator.add只适合列表拼接和数值累加这类满足结合律的场景。合并字典要自己写def merge_dict(left: dict, right: dict) - dict: out dict(left) out.update(right) return out class State(TypedDict): partials: Annotated[dict, merge_dict]自定义 reducer 有两个硬性要求函数签名是(旧值, 新值) - 合并值并且必须满足结合律和交换律。第二条特别容易忽略。我写过一版 reducer 用list.append风格实现结果并发分支的执行顺序一变输出顺序就变了同一份输入跑两次结果不一致定位花了很久。改成两边都返回新列表、直接相加之后才稳定。4.3 完整的分块并行打分实例把上面这些拼起来就是一个能跑的 Map-Reduceimport operator from typing import Annotated, TypedDict from langgraph.graph import StateGraph, START, END from langgraph.types import Send class State(TypedDict): doc: str chunks: list[str] scores: Annotated[list[tuple[int, float]], operator.add] summary: str class ChunkInput(TypedDict): idx: int chunk: str def split(state: State) - dict: doc state[doc] return {chunks: [doc[i:i 500] for i in range(0, len(doc), 500)]} def fanout(state: State): return [Send(score, {idx: i, chunk: c}) for i, c in enumerate(state[chunks])] def score(state: ChunkInput) - dict: # 真实场景这里换成模型调用示例用长度做伪打分 value min(len(state[chunk]) / 500.0, 1.0) return {scores: [(state[idx], value)]} def merge(state: State) - dict: top sorted(state[scores], keylambda x: x[1], reverseTrue)[:3] return {summary: f高价值块索引: {[i for i, _ in top]}} builder StateGraph(State) builder.add_node(split, split) builder.add_node(score, score) builder.add_node(merge, merge) builder.add_edge(START, split) builder.add_conditional_edges(split, fanout) builder.add_edge(score, merge) builder.add_edge(merge, END) app builder.compile() result app.invoke({doc: 内容 * 800}) print(result[summary])这段代码里有三个值得讲的设计点。第一Send的 payload 里带了idx。因为 reducer 是operator.add合并顺序由调度决定不保证和你扇出的顺序一致。如果你把顺序当成坐标含义去用比如第 3 个结果对应第 3 块一定会出错。带上索引在归并时sorted一下问题就消失了。这是我在真实项目里改过好几次的地方。第二score - merge这条普通边是隐式的 barrier。扇出之后score节点有 N 个待执行实例merge不会在第一个实例完成后就触发而是等所有score实例都结束才执行一次。理解这一点你就知道归并逻辑为什么可以放心地假设数据齐了。第三归并节点的输入是完整主状态。merge拿到的是所有 reducer 合并后的scores而不是某个分支的私有状态。这个不对称是 Map-Reduce 在 LangGraph 里的基本形状。4.4 并行度、顺序与失败隔离Send列表的长度就是扇出宽度LangGraph 会一次性把它们都调度出去。这意味着一件事你没法用图配置去限制并发数。文档切出 2000 块就是 2000 个分支同时待跑模型 API 的限流会被瞬间打满。我的做法是在分发函数里自己分批先切块再把块按批打包每批作为一个Send单元节点内部串行处理该批。批次大小按你的下游限流来定我一般取 8 到 16。这样扇出宽度可控而且批次本身还能复用重试策略。失败处理方面默认情况下任何一个分支抛异常都会让整个 superstep 失败。想做到部分失败不影响整体就在节点内部把异常吞掉并返回一个带标记的结果比如{scores: [(idx, -1.0)]}然后在归并阶段过滤掉。这样做的前提是你的 reducer 能接受这种失败占位。from langgraph.types import RetryPolicy builder.add_node( score, score, retry_policyRetryPolicy(max_attempts3, retry_on(TimeoutError, ConnectionError)), )RetryPolicy的可用性和导入路径随版本变化用之前先确认一下你的版本支持。它解决的是偶发失败自动重试解决不了整体超时——全局超时要靠上游asyncio.wait_for之类的机制兜住。5. 踩坑实录报错信息背后的真实原因5.1 InvalidUpdateError并行写同一个键却没有 reducer完整的报错大致是InvalidUpdateError: At key scores: Can receive only one value per step. Use an Annotated key to handle multiple values.这个信息其实挺直白但真正容易踩的场合是我明明加了 reducer怎么还报。我当时的情况是State里声明了Annotated[list, operator.add]但节点函数返回时用的是自定义类而不是listreducer 拿到的类型对不上照样炸。排查链路我一般这样走确认报错里的 key 名和你TypedDict里的字段名完全一致关系字段拼错很常见。确认该字段确实用了Annotated不是只写了list[str]。确认节点返回的值类型和 reducer 能处理的类型一致——operator.add对list list有效对list dict直接抛异常。如果 reducer 是自定义的先在纯 Python 里单独测一遍输入两个真实样本看输出是不是你期望的。5.2 Send 之后主状态没变因为 payload 是私有的这个坑的表现很迷惑图跑完了没报错merge节点也执行了但读出来的结果是空的或者只有默认值。原因通常有两种。一是节点返回的键名写错了。Send分支节点的返回值要写回主状态键名必须是主状态里存在的字段名写错的话那个值会被静默丢弃部分版本会给 warning但不一定。二是误以为在分支里用Send的 payload 改了状态。前面说过payload 是私有输入它不参与主状态合并。如果分支节点收到{idx: 0, chunk: ...}然后直接返回这个 dict那idx和chunk会被当成状态更新去写主状态——如果主状态里恰好没有这两个字段就没效果如果有同名字段但没有 reducer就撞上 5.1 的报错。所以分支节点返回什么一定要照主状态的字段来。我的习惯是在分支节点上加类型注解def score(state: ChunkInput) - dict输入类型和输出键名分开写清楚虽然 Python 不强制但对排查帮助极大。5.3 同一 superstep 里出现两个 Command并行分支里只有一个节点能返回Command这是规则。我遇到的场景是扇出后每个分支都想根据自己的结果决定下一步跳哪于是每个分支都返回了Command(goto...)。结果是不确定的图的行为变得难以解释。正确做法是把决定去哪这件事抽出来。分支只返回结果数据用一个显式的归并节点收集所有分支的结果在归并节点里统一做判断由它返回一个Command。这样一来扇出和路由在时间上是分离的语义清晰也不会撞规则。5.4 recursion_limit 与 thread_id 引发的假死recursion_limit默认 25。一个包含模型生成 - 工具调用 - 模型生成的回环流程很容易几轮就撞上。撞上时的报错是GraphRecursionError信息还算明确调大就行app.invoke(input_data, {configurable: {thread_id: t1}, recursion_limit: 100})更隐蔽的是thread_id相关的问题。不带thread_id却配了checkpointer某些版本下interrupt之后的恢复会找不到现场表现为图停了但怎么都续不上。另外重跑同一个thread_id时图默认会从上次的 checkpoint 继续而不是从头开始——如果你在本地反复调试同一个thread_id看到的可能是上一次残留的状态会误以为逻辑写错了。调试期我固定用时间戳生成thread_id发布后再换成业务侧稳定的 ID。6. 把两个原语拼成一个能上线的调度骨架6.1 骨架长什么样在真实项目里Command和Send通常是组合出现而不是各写各的。我目前用得最顺的结构是这样START - classify分类/意图识别返回 Command 动态分流 - plan规划决定要处理哪些单元写入 units 字段 - fanout条件边按 units 生成 Send 列表 - worker并行处理返回部分结果 - reduce归并返回 Command 决定是否需要再审 - END 或 回到 review这个骨架的关键在于把职责切干净分类节点只负责去哪规划节点只负责有哪些单元扇出函数只负责把这些单元变成 Send工作节点只负责处理一个单元并返回结果归并节点只负责看齐数据后做决策。每一层都只做一件事出了问题时定位范围就很小。注意规划节点必须把单元列表写进主状态的独立字段比如units不能在扇出函数里临时算。原因是扇出函数可能被调用多次比如恢复执行时每次都要拿到同一份单元列表否则分支数量会漂移。6.2 调试与可观测性并行执行的调试体验和串行完全不同最关键的是能看到每个分支到底跑了什么、返回了什么。用stream_mode按更新粒度看for chunk in app.stream(input_data, stream_modeupdates, config{configurable: {thread_id: t1}}): for node_name, payload in chunk.items(): print(node_name, payload)updates模式会逐个节点、逐个分支地把更新吐出来扇出后你能看到score这个名字反复出现每次带一份不同的部分结果。想看完整状态就用values模式想看模型消息流就用messages。我一般在开发期用updates盯数据形状压测时切回values看整体。另外app.get_graph().draw_ascii()对Send的展示是一个节点指向自己看不出动态宽度这点要有心理预期别指望从静态图里读出扇出规模。6.3 上线前的自检清单这套东西我踩坑踩下来的经验浓缩成一份上机前会过一遍的清单检查项为什么必须确认所有会被并行写入的字段都有 reducer少一个就是InvalidUpdateError自定义 reducer 满足结合律与交换律否则并发顺序会影响最终结果扇出 payload 里带了稳定索引归并顺序不等于扇出顺序扇出宽度做过分批控制避免打满下游限流分支节点内没有重副作用操作重试或恢复时会重复执行interrupt之前的逻辑幂等恢复时该节点会从头执行归并节点之后再出现Command并行分支里只能有一个Command生效调试用thread_id与生产分离否则会续上历史状态干扰判断recursion_limit按最坏回环轮数估算默认 25 往往不够我个人在实际项目里体会最深的一点是LangGraph 的并行能力本身不难用难的是想清楚哪些数据需要在分支之间共享、哪些必须隔离。一旦你接受了分支读私有输入、写走 reducer 收口这个单向模型Send的绝大多数行为就变得可预测了。反过来如果你总想绕过这个模型去共享可变状态坑会一个接一个而且报错信息通常不会直接告诉你根源在哪。最后再分享一个排查小技巧遇到并行相关的问题先别急着看图结构把扇出函数的返回值print出来看看——是几个Send、每个 payload 的键是什么。我遇到过的并行 bug 里超过一半的问题在这一步就暴露了。
返回列表