1. 扇出是什么 #

日常写图,最开始都是一条直线:START → A → B → END。真实任务却经常是「同一件事要同时干好几份」——一篇长文切成三段各自摘要、一个问题拆成三个子问题分别检索、一份订单同时去查库存和查物流。

这种「从一个节点同时触发多个下游」的结构,就叫扇出(fan-out)。名字来自电路:一根线分成好几根。

拆开之后,你有两种收场:

所以扇入不是另一种并行技术,只是扇出的合拢。没有扇出,谈扇入没有意义;有扇出,却不一定非要扇入。

1.1. 你学完能做什么 #

学完你应该能亲手写出三种由易到难的扇出,并且说清它们各自解决什么问题:

  1. 静态扇出:分支写死在图上,add_edge 连两次就行
  2. 列表扇出:运行时才决定这次并行跑几个节点,但每个节点看到的仍是同一份状态
  3. Send 扇出:运行时既决定跑几个,又给每个工人一份不同的输入(map-reduce)

另外两件配套的事必须一起会:

2. 前置知识 #

扇出本身只是「边怎么连」。真正卡住的,往往是三个配角:图怎么搭、什么叫同时、并行写入怎么合并。先把这三块补上,后面的扇出代码才读得懂。

2.1. 图的三件套:状态、节点、边 #

LangGraph 的核心概念只有三个:

概念 一句话 对应代码
状态 State 一个字典,整张图运行期间所有节点共享它 class Mini(TypedDict): ...
节点 Node 一个普通函数:读状态,返回要改的字段 def greet(state) -> dict: ...
边 Edge 规定「这个节点跑完,下一个跑谁」 builder.add_edge("a", "b")

用一句话串起来:

图 = 一堆函数(节点) + 它们的执行顺序(边) + 它们之间传递的数据(状态)。

可以把它想成工厂流水线:节点是工位上的工人,边是传送带,状态是那个一路传下去的托盘。工人不直接把零件递给下一个人,而是放回托盘——所以后面要并行时,多个工人可以同时从托盘上取东西。

下面是一张最小的图:入口进 greet,它读到名字,写回一句问候,然后结束。

# TypedDict 用来声明「状态这个字典里有哪些键、各是什么类型」
from typing import TypedDict

# END 是终点哨兵,START 是入口哨兵,StateGraph 是图的构建器
from langgraph.graph import END, START, StateGraph


# 状态里只有一个字符串字段
class Mini(TypedDict):
    # 用户名,后面会被改写成问候语
    text: str


# 节点就是普通函数:参数是当前状态,返回「要改哪些字段」
def greet(state: Mini) -> dict:
    # 只返回要更新的部分,图会把它合并进完整状态
    return {"text": f"你好,{state['text']}"}


# 新建一张未编译的图,绑定状态类型
builder = StateGraph(Mini)
# 把 greet 函数注册成名叫 "greet" 的节点
builder.add_node("greet", greet)
# 入口:图一开始就跑 greet
builder.add_edge(START, "greet")
# greet 跑完就结束
builder.add_edge("greet", END)
# 编译成可以 invoke 的对象
graph = builder.compile()
# 传入初始状态,打印最终状态
print(graph.invoke({"text": "小明"}))

运行输出:

{'text': '你好,小明'}

记住一件事就够:节点返回的是「增量」,不是完整状态。 greet 只返回 {"text": ...},图负责把它合并回去。扇出时多个节点各自返回增量,图再在同一步里把它们合在一起——这就是下一节 reducer 要管的事。

2.2. 超步:什么叫「同时跑」 #

LangGraph 不是「一个函数 return 了立刻跑下一个」。它按超步(superstep)往前推:

  1. 找出这一步该跑的所有节点
  2. 把它们都跑完(每个节点看到的都是这一步开始前的状态)
  3. 把所有返回值合并进状态
  4. 根据边决定下一步跑谁

所以「并行」在这里有一个非常具体的含义:

同一超步里的节点,彼此看不见对方刚刚写下的值。

这和「两个函数真的在两个 CPU 上同时跑」不是一回事。默认用普通 def 时,LangGraph 仍可能一个接一个调用它们,只是合并发生在全部跑完之后。想靠扇出缩短墙钟时间,节点得写成 async def,那是后话。这篇用同步函数,方便你用 print 看清结构。

怎么亲眼看见「同一超步」?用 stream_mode="debug",看每个任务的 step 编号。具体数字放到 §4.2 再跑。现在只要记住:

2.3. Reducer:并行写入怎么合成一份 #

节点返回增量之后,图要决定「新值和旧值怎么合」。默认规则是覆盖:后来的整份换掉原来的。

串行时覆盖往往看不出来问题——A 写完,B 再写,本来就有先后。并行时两个节点在同一超步里写同一个键,图不知道该听谁的,会直接抛 InvalidUpdateError。

解决办法是给这个字段挂一个 reducer:一个「旧值 + 新值 → 合并结果」的函数。列表最常用的是 operator.add,效果就是拼接:

旧值 ["start"]  +  新值 ["天气好"]  →  ["start", "天气好"]

写法是 Python 自带的 Annotated:在类型右边贴一张便签,便签上写合并规则。

# operator.add 对列表的含义是拼接,不是算术加法
import operator
# Annotated 用来给类型附加额外信息,这里附加的是 reducer
from typing import Annotated, TypedDict

# END / START / StateGraph 是搭图的三件套
from langgraph.graph import END, START, StateGraph


# 状态:一个会累加的日志字段
class S(TypedDict):
    # 挂了 operator.add:每个节点返回的列表都会被拼到后面
    log: Annotated[list[str], operator.add]


# 第一个节点往 log 里追加一条
def a(state: S) -> dict:
    # 返回单元素列表,reducer 负责拼到已有 log 后面
    return {"log": ["A 来过"]}


# 第二个节点也往 log 里追加一条
def b(state: S) -> dict:
    # 同样返回单元素列表
    return {"log": ["B 来过"]}


# 新建图
builder = StateGraph(S)
# 注册节点 a
builder.add_node("a", a)
# 注册节点 b
builder.add_node("b", b)
# 入口进 a
builder.add_edge(START, "a")
# a 跑完接着 b(这是串行,用来先看清 reducer 的效果)
builder.add_edge("a", "b")
# b 跑完结束
builder.add_edge("b", END)
# 编译
graph = builder.compile()
# 初始 log 是空列表,最终应看到两条
print(graph.invoke({"log": []}))

运行输出:

{'log': ['A 来过', 'B 来过']}

如果这里不写 Annotated[..., operator.add],b 返回的 ["B 来过"] 会把 ["A 来过"] 整段盖掉,最后只剩 B。扇出时两个节点同时写,连盖掉的机会都没有,直接报错。并行写入的字段必须有 reducer,这是扇出能跑通的前提。

3. 扇出是什么 #

前置知识齐了,可以给「扇出」下定义了。

3.1. 一个生活里的类比 #

点外卖时,商家不会先炒完菜再去做饮料:后厨和饮品台同时开工,都好了才装袋给你。

对应到图上:

              ┌─ weather ─┐
START → start ┤           ├→ merge → END
              └─ news ────┘

start 是接单,weather 和 news 是两个同时干活的工位,merge 是出餐口。

3.2. 扇出、分支、扇入差在哪 #

这三个词很容易混,差在「下一步跑几个」以及「要不要再合回来」:

下一步跑几个 典型问法 怎么连边
分支 只跑 1 个 「退货走 A,否则走 B」 条件边返回一个节点名
扇出 同时跑 N 个 「库存和物流一起查」 从同一点连出多条边,或返回一个列表
扇入 N 个都跑完再跑 1 个 「两份结果收拢成一份报告」 多条边指向同一个下游节点

对照着记:

扇入没有单独的函数。你只要写两次 add_edge("weather", "merge") 和 add_edge("news", "merge"),LangGraph 就会等两边都到齐再跑 merge。§4.1 是只扇出不合拢,§4.3 才把扇入接上。

3.3. 三种实现,由易到难 #

后面三节按这个顺序讲,建议不要跳:

  1. 静态扇出(§4):编译前就知道有哪几条并行边,add_edge 写死
  2. 列表扇出(§6):运行时才知道这次要跑哪几个节点,但大家看到同一份状态
  3. Send 扇出(§7):运行时才知道跑几个,而且每人一份不同输入

选型很简单:

4. 静态扇出:边写死,分支固定 #

最简单的扇出不需要条件边,也不需要 Send。从同一个节点连出两条(或更多)固定边,下游就会进入同一个超步。

4.1. 从同一个节点连出两条边 #

场景:用户给一个城市名,同时去「查天气」和「查新闻」。两个节点各往 log 里记一条,最后一起出现在结果里。

# operator.add 给 log 做列表拼接
import operator
# Annotated 挂 reducer,TypedDict 声明状态
from typing import Annotated, TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph


# 状态:城市名 + 一条会累加的日志
class S(TypedDict):
    # 普通字段,谁后写谁覆盖;这里只读不抢着写
    topic: str
    # 并行节点都往这里追加,所以必须挂 reducer
    log: Annotated[list[str], operator.add]


# 起点:只记一笔「开始了」
def start_node(state: S) -> dict:
    # 打印出来,方便对照执行顺序
    print("   [start] 执行")
    # 往 log 追加一条
    return {"log": ["start"]}


# 并行工人甲:根据城市写一条天气
def weather(state: S) -> dict:
    # 打印,证明这个节点确实跑了
    print("   [weather] 执行")
    # 读的是 start 跑完之后的状态,topic 还在
    return {"log": [f"天气:{state['topic']} 今天晴"]}


# 并行工人乙:根据城市写一条新闻
def news(state: S) -> dict:
    # 打印
    print("   [news] 执行")
    # 同样读 topic,往 log 追加
    return {"log": [f"新闻:{state['topic']} 有新进展"]}


# 新建图
builder = StateGraph(S)
# 注册起点
builder.add_node("start", start_node)
# 注册天气节点
builder.add_node("weather", weather)
# 注册新闻节点
builder.add_node("news", news)
# 入口进 start
builder.add_edge(START, "start")
# 从 start 连到 weather:第一条扇出边
builder.add_edge("start", "weather")
# 从 start 再连到 news:第二条扇出边。两条都连上,就是扇出
builder.add_edge("start", "news")
# weather 自己结束后去 END
builder.add_edge("weather", END)
# news 自己结束后去 END
builder.add_edge("news", END)
# 编译
graph = builder.compile()
# 跑一次,看 log 里是否两条都在
print("结果:", graph.invoke({"topic": "北京", "log": []}))

运行输出:

   [start] 执行
   [news] 执行
   [weather] 执行
结果: {'topic': '北京', 'log': ['start', '新闻:北京 有新进展', '天气:北京 今天晴']}

两件事值得看:

  1. 你没有写任何「并行」关键字。 并行完全来自「同一个节点连出了两条边」。
  2. news 可能出现在 weather 前面。 同一超步内的顺序不稳定,不要依赖「先注册的先跑」。要稳定顺序,就不要扇出,改回串行。

4.2. 用 stream 看见同一超步 #

invoke 只给你最终状态。想看中间发生了什么,用 stream:默认每个节点产出一个 chunk。再打开 debug 模式,就能看到超步编号。

# operator.add 给 log 做列表拼接
import operator
# Annotated 挂 reducer,TypedDict 声明状态
from typing import Annotated, TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph


# 和上一节相同的状态
class S(TypedDict):
    # 城市名
    topic: str
    # 累加日志
    log: Annotated[list[str], operator.add]


# 起点节点
def start_node(state: S) -> dict:
    # 追加一条 start
    return {"log": ["start"]}


# 天气节点
def weather(state: S) -> dict:
    # 追加天气
    return {"log": [f"天气:{state['topic']} 今天晴"]}


# 新闻节点
def news(state: S) -> dict:
    # 追加新闻
    return {"log": [f"新闻:{state['topic']} 有新进展"]}


# 新建图,结构与上一节完全一样
builder = StateGraph(S)
# 注册三个节点
builder.add_node("start", start_node)
# 注册天气
builder.add_node("weather", weather)
# 注册新闻
builder.add_node("news", news)
# 入口
builder.add_edge(START, "start")
# 扇出边 1
builder.add_edge("start", "weather")
# 扇出边 2
builder.add_edge("start", "news")
# 天气出口
builder.add_edge("weather", END)
# 新闻出口
builder.add_edge("news", END)
# 编译
graph = builder.compile()

# 先用默认 stream:每个节点一个 chunk
print("=== 默认 stream ===")
# 每循环一次就是一个节点刚刚写下的增量
for chunk in graph.stream({"topic": "北京", "log": []}):
    # 打印节点名和它返回的字典
    print("  ", chunk)

# 再用 debug 模式看超步编号
print("=== debug 超步 ===")
# stream_mode="debug" 会打出更细的事件
for chunk in graph.stream({"topic": "北京", "log": []}, stream_mode="debug"):
    # 只看 task 事件,它告诉你「这一步要跑哪个节点」
    if chunk["type"] == "task":
        # step 是超步号,name 是节点名
        print(f"   step={chunk['step']} node={chunk['payload']['name']}")

运行输出:

=== 默认 stream ===
   {'start': {'log': ['start']}}
   {'news': {'log': ['新闻:北京 有新进展']}}
   {'weather': {'log': ['天气:北京 今天晴']}}
=== debug 超步 ===
   step=1 node=start
   step=2 node=news
   step=2 node=weather

start 的 step 是 1,news 和 weather 的 step 都是 2。编号相同就是同一超步,这就是扇出的实证。

默认 stream 里,并行的两个节点各出一个 chunk,不会合成一个。这也说明:观察进度时按「节点」看,不要指望一步只出现一条。

4.3. 扇入:两条线再汇合 #

上一节两个工人各自走到 END,图在两边都结束时收工。更常见的是:两边都跑完,再进一个汇总节点。做法是让 weather 和 news 都连到同一个 merge。

LangGraph 会等两条边都到达,再跑 merge。merge 读到的 log,已经包含两个工人追加的内容。

# operator.add 给 log 做列表拼接
import operator
# Annotated 挂 reducer,TypedDict 声明状态
from typing import Annotated, TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph


# 状态不变
class S(TypedDict):
    # 城市名
    topic: str
    # 累加日志
    log: Annotated[list[str], operator.add]


# 起点
def start_node(state: S) -> dict:
    # 记一笔
    print("   [start] 执行")
    # 追加 start
    return {"log": ["start"]}


# 天气工人
def weather(state: S) -> dict:
    # 记一笔
    print("   [weather] 执行")
    # 追加天气
    return {"log": [f"天气:{state['topic']} 今天晴"]}


# 新闻工人
def news(state: S) -> dict:
    # 记一笔
    print("   [news] 执行")
    # 追加新闻
    return {"log": [f"新闻:{state['topic']} 有新进展"]}


# 汇总节点:等两个工人都到齐再跑
def merge(state: S) -> dict:
    # 打印此时已经攒下的 log,证明扇入发生在合并之后
    print("   [merge] 执行, 此时 log =", state["log"])
    # 再追加一条 merge,表示汇总完成
    return {"log": ["merge"]}


# 新建图
builder = StateGraph(S)
# 注册四个节点
builder.add_node("start", start_node)
# 注册天气
builder.add_node("weather", weather)
# 注册新闻
builder.add_node("news", news)
# 注册汇总
builder.add_node("merge", merge)
# 入口
builder.add_edge(START, "start")
# 扇出到天气
builder.add_edge("start", "weather")
# 扇出到新闻
builder.add_edge("start", "news")
# 天气跑完去 merge(扇入的一条边)
builder.add_edge("weather", "merge")
# 新闻跑完也去 merge(扇入的另一条边)
builder.add_edge("news", "merge")
# 汇总完结束
builder.add_edge("merge", END)
# 编译
graph = builder.compile()
# 跑一次
print("结果:", graph.invoke({"topic": "北京", "log": []}))

运行输出:

   [start] 执行
   [news] 执行
   [weather] 执行
   [merge] 执行, 此时 log = ['start', '新闻:北京 有新进展', '天气:北京 今天晴']
结果: {'topic': '北京', 'log': ['start', '新闻:北京 有新进展', '天气:北京 今天晴', 'merge']}

merge 打印的时候,两条工人日志已经在了。执行顺序一定是:先 start,再两个工人(顺序不定),最后才是 merge。这就是扇出 + 扇入的完整形状。

5. 没配 reducer 会怎样 #

静态扇出看起来太顺利了。下面把 log 换成一个没有 reducer 的普通字段,让两个工人同时写它。这是扇出里最常见的报错,建议自己跑一次,下次看见 InvalidUpdateError 就知道该往哪改。

# TypedDict 声明状态,这次故意不引入 Annotated
from typing import TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph


# 两个并行节点都要写 note,但 note 没有 reducer
class Plain(TypedDict):
    # 城市名,只读
    topic: str
    # 普通字符串:同一超步里只能被写一次
    note: str


# 新建图
builder = StateGraph(Plain)
# 起点随便写一下 note,不参与冲突
builder.add_node("start", lambda s: {"note": "start"})
# 天气工人写 note
builder.add_node("weather", lambda s: {"note": "天气好"})
# 新闻工人也写 note,和天气抢同一个键
builder.add_node("news", lambda s: {"note": "有新闻"})
# 入口
builder.add_edge(START, "start")
# 扇出
builder.add_edge("start", "weather")
# 扇出第二条
builder.add_edge("start", "news")
# 天气出口
builder.add_edge("weather", END)
# 新闻出口
builder.add_edge("news", END)
# 编译能通过:compile 不检查「会不会并行抢键」
graph = builder.compile()
# 运行时才会爆
try:
    # 同一超步里 weather 和 news 都返回 note
    print(graph.invoke({"topic": "北京", "note": ""}))
except Exception as e:
    # 打印异常类型和第一行信息
    print(type(e).__name__, str(e).splitlines()[0])

运行输出:

InvalidUpdateError At key 'note': Can receive only one value per step. Use an Annotated key to handle multiple values.

报错已经把做法写在脸上了:Use an Annotated key。把 note 改成列表并挂 operator.add,或者改成两个工人写不同的键(一个写 weather_note,一个写 news_note),都可以。

对照记一句:

串行写同一字段 并行写同一字段
没 reducer 后写的盖掉先写的(不报错) 直接报错
有 reducer 按规则合并 按规则合并

扇出出问题,第一反应永远是:他们是不是在抢同一个没配 reducer 的键?

6. 列表扇出:这次跑几个可以变 #

静态扇出的边是写死的:每次都是天气 + 新闻。如果用户有时只要天气、有时三个都要,编译期就没法把边画死。这时让条件边返回一个节点名列表,列表里的节点会进入同一个超步。

6.1. 运行时决定跑哪几个 #

路由函数读状态里的 tasks,把它翻译成节点名。出口表(path_map)列出所有可能的下游;返回值是这一次的子集。

# operator.add 给 log 做列表拼接
import operator
# Annotated 挂 reducer,TypedDict 声明状态
from typing import Annotated, TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph


# 状态:用户点了哪些任务 + 日志
class BranchS(TypedDict):
    # 用户勾选的任务名,比如 ["天气", "路况"]
    tasks: list[str]
    # 累加日志,用来观察实际跑了谁
    log: Annotated[list[str], operator.add]


# 分发节点:本身不决定去哪,只记一笔
def dispatch(state: BranchS) -> dict:
    # 追加 dispatch,证明分发节点跑过了
    return {"log": ["dispatch"]}


# 路由函数:根据 tasks 翻译出节点名列表
def pick(state: BranchS) -> list[str]:
    # 中文任务名 → 节点名
    mapping = {"天气": "weather", "新闻": "news", "路况": "traffic"}
    # 只保留认识的任务,翻译成节点名
    return [mapping[t] for t in state["tasks"] if t in mapping]


# 天气节点
def weather(state: BranchS) -> dict:
    # 记下自己的名字
    return {"log": ["weather"]}


# 新闻节点
def news(state: BranchS) -> dict:
    # 记下自己的名字
    return {"log": ["news"]}


# 路况节点
def traffic(state: BranchS) -> dict:
    # 记下自己的名字
    return {"log": ["traffic"]}


# 新建图
builder = StateGraph(BranchS)
# 注册分发节点
builder.add_node("dispatch", dispatch)
# 注册三个候选工人
builder.add_node("weather", weather)
# 注册新闻
builder.add_node("news", news)
# 注册路况
builder.add_node("traffic", traffic)
# 入口进分发
builder.add_edge(START, "dispatch")
# 条件边:pick 返回列表;第三个参数是出口全集
builder.add_conditional_edges("dispatch", pick, ["weather", "news", "traffic"])
# 三个工人跑完都结束
builder.add_edge("weather", END)
# 新闻出口
builder.add_edge("news", END)
# 路况出口
builder.add_edge("traffic", END)
# 编译
graph = builder.compile()
# 三个都要
print("三项:", graph.invoke({"tasks": ["天气", "新闻", "路况"], "log": []}))
# 只要两个
print("两项:", graph.invoke({"tasks": ["天气", "路况"], "log": []}))

运行输出:

三项: {'tasks': ['天气', '新闻', '路况'], 'log': ['dispatch', 'news', 'traffic', 'weather']}
两项: {'tasks': ['天气', '路况'], 'log': ['dispatch', 'traffic', 'weather']}

第二次没有 news。出口表是可能的全集,返回值是这一次的子集。 三个候选节点都要先 add_node 注册好,不能在运行时发明一个没注册的名字。

和静态扇出比,列表扇出多出来的能力只有一句:分支集合可以随输入变。每个工人看到的仍是同一份状态(都看得到完整的 tasks)。一人一份不同输入,是下一节 Send 的事。

6.2. 返回空列表会悄悄结束 #

tasks 若是空的,pick 就返回 []。空列表的含义是「一个下游都不触发」,图就此结束,不报错。

# operator.add 给 log 做列表拼接
import operator
# Annotated 挂 reducer,TypedDict 声明状态
from typing import Annotated, TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph


# 复用上一节的状态
class BranchS(TypedDict):
    # 任务清单
    tasks: list[str]
    # 累加日志
    log: Annotated[list[str], operator.add]


# 分发节点
def dispatch(state: BranchS) -> dict:
    # 记一笔
    return {"log": ["dispatch"]}


# 路由:空 tasks 就会得到空列表
def pick(state: BranchS) -> list[str]:
    # 翻译表
    mapping = {"天气": "weather", "新闻": "news", "路况": "traffic"}
    # 空输入 → 空列表
    return [mapping[t] for t in state["tasks"] if t in mapping]


# 三个工人(这次一个都不会跑,但图上必须注册)
def weather(state: BranchS) -> dict:
    # 若被调用会记一笔
    return {"log": ["weather"]}


# 新闻工人
def news(state: BranchS) -> dict:
    # 记一笔
    return {"log": ["news"]}


# 路况工人
def traffic(state: BranchS) -> dict:
    # 记一笔
    return {"log": ["traffic"]}


# 新建图
builder = StateGraph(BranchS)
# 注册分发
builder.add_node("dispatch", dispatch)
# 注册天气
builder.add_node("weather", weather)
# 注册新闻
builder.add_node("news", news)
# 注册路况
builder.add_node("traffic", traffic)
# 入口
builder.add_edge(START, "dispatch")
# 条件边,出口表仍是三个
builder.add_conditional_edges("dispatch", pick, ["weather", "news", "traffic"])
# 三个出口
builder.add_edge("weather", END)
# 新闻出口
builder.add_edge("news", END)
# 路况出口
builder.add_edge("traffic", END)
# 编译
graph = builder.compile()
# 传入空任务列表
print("空:", graph.invoke({"tasks": [], "log": []}))

运行输出:

空: {'tasks': [], 'log': ['dispatch']}

只有 dispatch,三个工人都没跑,也没有任何异常。写「按条件过滤出要跑的分支」时要小心:过滤到一个都不剩,图就悄悄收工。 稳妥做法是空列表时返回一个兜底节点名,比如 "idle",在那里给用户一句「你没有勾选任何任务」。

7. Send:每个工人拿到不同的输入 #

列表扇出有一个限制:所有工人看到的是同一份状态。这在「天气 / 新闻 / 路况」这种「工种不同、输入相同」的场景刚好够用。换成「把 10 个问题分给 10 个工人,每人只处理一个」,同一份状态就不够了——每个工人需要自己那一条。

7.1. 为什么普通列表不够 #

假如状态是 {"questions": ["是什么", "怎么用", "注意什么"]},你返回 ["worker", "worker", "worker"],三个实例看到的仍是整份 questions。工人自己还得想办法知道「我是三个里的第几个」,既别扭又容易抢同一个字段。

Send 就是为这个场景准备的:每发一份,都带上这个工人该看到的那一小段数据。

7.2. Send 的两个参数 #

Send 来自 langgraph.types,长这样:

Send("节点名", {这个工人能看到的全部状态})

两个要点,建议读两遍:

  1. 第二个参数就是工人眼里的全部状态。 主图里的其它字段,工人默认看不见。需要用的,必须亲手塞进这个字典。
  2. 工人返回的增量,仍按字段名合并回主图。 所以工人虽然输入是「一份专属小字典」,输出的键(比如 answers)必须在主状态里存在,并且并行写入时要有 reducer。

可以把它想成发快递:面单上写清送到哪个工位(节点名),箱子里只放这个工位需要的零件(专属输入)。工位干完把成品放回仓库里对应的货架(主状态的 reducer 字段)。

7.3. 完整的拆分 → 扇出 → 汇总 #

下面是一个不调模型的 map-reduce:

  1. Map / 拆分:把主题拆成三个子问题
  2. Fan-out / 扇出:每个子问题一份 Send,三个 worker 实例并行
  3. Reduce / 扇入:worker 都连到 reduce,把答案拼成一段报告

工人函数读的是 WorkerIn(只有 question 和 topic),不是主状态。这正是 Send 的意义。

# operator.add 用来把每个工人返回的 answers 拼成一张大列表
import operator
# Annotated 挂 reducer,TypedDict 声明两种状态
from typing import Annotated, TypedDict

# 搭图用的三件套
from langgraph.graph import END, START, StateGraph
# Send 表示「派一个节点实例去处理这一小段数据」
from langgraph.types import Send


# 主图状态:主题、子问题清单、并行答案、最终报告
class MapState(TypedDict):
    # 用户给的大主题
    topic: str
    # 拆分出来的子问题
    questions: list[str]
    # 每个工人追加一条答案,必须有 reducer
    answers: Annotated[list[str], operator.add]
    # 汇总后的完整报告
    report: str


# 工人看到的输入:只有 Send 里塞进来的字段
class WorkerIn(TypedDict):
    # 自己负责的那一个子问题
    question: str
    # 主题也带上,写答案时用得着
    topic: str


# 拆分节点:根据主题造三个子问题
def split_topic(state: MapState) -> dict:
    # 演示用:固定拆成三个问法,真实项目里这里可以调模型
    qs = [
        # 第一个子问题
        f"{state['topic']}:是什么",
        # 第二个子问题
        f"{state['topic']}:怎么用",
        # 第三个子问题
        f"{state['topic']}:注意什么",
    ]
    # 写入 questions;report 先清空
    return {"questions": qs, "report": ""}


# 路由函数:为每个子问题构造一份 Send
def fan_out(state: MapState) -> list:
    # 列表推导:一个问题对应一个 Send
    return [
        # 目标节点都叫 worker;第二参数是这个实例能看到的全部内容
        Send("worker", {"question": q, "topic": state["topic"]})
        # 遍历拆分结果
        for q in state["questions"]
    ]


# 工人:只处理自己那一个问题
def worker(state: WorkerIn) -> dict:
    # 打印工人实际看到的字典,方便确认「看不见主状态里的其它键」
    print(f"   [worker] 看到 {state}")
    # 模拟回答:真实项目里这里调模型或检索
    return {"answers": [f"{state['question']} → 模拟回答"]}


# 汇总:等所有工人到达后再跑
def reduce(state: MapState) -> dict:
    # 把多条答案拼成一段话
    body = ";".join(state["answers"])
    # 写入报告字段
    return {"report": f"关于「{state['topic']}」的汇总:{body}"}


# 新建图,绑定的是主状态
builder = StateGraph(MapState)
# 注册拆分
builder.add_node("split", split_topic)
# 注册工人(会被 Send 多次进入)
builder.add_node("worker", worker)
# 注册汇总
builder.add_node("reduce", reduce)
# 入口先进拆分
builder.add_edge(START, "split")
# 拆分完由 fan_out 动态扇出;出口表里写出 worker
builder.add_conditional_edges("split", fan_out, ["worker"])
# 每个工人实例跑完都去 reduce,reduce 会等全部到齐
builder.add_edge("worker", "reduce")
# 汇总完结束
builder.add_edge("reduce", END)
# 编译
graph = builder.compile()
# 造一份完整初始状态:reducer 字段给空列表,普通字段给空字符串
print(
    "结果:",
    graph.invoke(
        {
            # 大主题
            "topic": "扇出",
            # 子问题待拆分节点填写
            "questions": [],
            # 答案从空列表开始累加
            "answers": [],
            # 报告待汇总节点填写
            "report": "",
        }
    ),
)

运行输出:

   [worker] 看到 {'question': '扇出:是什么', 'topic': '扇出'}
   [worker] 看到 {'question': '扇出:怎么用', 'topic': '扇出'}
   [worker] 看到 {'question': '扇出:注意什么', 'topic': '扇出'}
结果: {'topic': '扇出', 'questions': ['扇出:是什么', '扇出:怎么用', '扇出:注意什么'], 'answers': ['扇出:是什么 → 模拟回答', '扇出:怎么用 → 模拟回答', '扇出:注意什么 → 模拟回答'], 'report': '关于「扇出」的汇总:扇出:是什么 → 模拟回答;扇出:怎么用 → 模拟回答;扇出:注意什么 → 模拟回答'}

对照着看三处:

这就是 map-reduce 在 LangGraph 里的最小形状。把 worker 里的「模拟回答」换成真正的检索或模型调用,就是生产里最常见的用法。

7.4. 漏带字段会立刻报错 #

Send 的第二参数漏了工人要用的键,工人一读就会 KeyError。这一点比「静默跑错」友好:漏了就炸,不会假装成功。

# operator.add 给 answers 做拼接
import operator
# Annotated / TypedDict
from typing import Annotated, TypedDict

# 搭图三件套
from langgraph.graph import END, START, StateGraph
# Send 用来投递专属输入
from langgraph.types import Send


# 主状态
class MapState(TypedDict):
    # 待处理条目
    items: list[str]
    # 并行结果
    answers: Annotated[list[str], operator.add]


# 工人故意同时读 item 和 topic,但 Send 只给了 item
def worker(state: dict) -> dict:
    # 这里访问 topic:Send 里没有就会 KeyError
    return {"answers": [state["item"] + "|" + state["topic"]]}


# 新建图
builder = StateGraph(MapState)
# 注册工人
builder.add_node("worker", worker)
# 从 START 直接扇出;每个 Send 只塞了 item,没塞 topic
builder.add_conditional_edges(
    START,
    # 路由:给每个 item 发一份残缺的 Send
    lambda s: [Send("worker", {"item": i}) for i in s["items"]],
    # 出口表
    ["worker"],
)
# 工人出口
builder.add_edge("worker", END)
# 编译
graph = builder.compile()
# 运行并抓异常
try:
    # items 有一条,足够触发工人
    print(graph.invoke({"items": ["a"], "answers": []}))
except Exception as e:
    # 打印类型和缺的那个键
    print(type(e).__name__, "|", e)

运行输出:

KeyError | 'topic'

修好的方法只有一个:工人要用的字段,Send 里都带上。 不要指望它能「顺便」读到主状态。

8. 常见坑 #

前面的报错都见过了,这里收成一张对照表,方便出问题时对照。每一条都是扇出场景里会真实碰到的,不是边角知识。

现象 原因 怎么办
InvalidUpdateError: Can receive only one value per step 同一超步里多个节点写同一个没 reducer 的字段 给该字段加 Annotated[..., operator.add],或改成写不同的键
最终列表少了某条 串行覆盖,或以为扇出会自动追加 并行字段必须有 reducer;串行覆盖是默认行为
KeyError: '某字段'(Send 场景) 工人读了 Send 里没给的键 把该字段放进 Send 的第二个字典
工人一个都没跑,也不报错 路由函数返回了空列表 空列表时改返回一个兜底节点
并行节点的打印顺序每次不一样 同一超步不保证顺序 不要依赖顺序;要顺序就改串行
墙钟时间并没有变短 普通 def 仍是一个接一个调用 结构并行 ≠ 时间并行;要重叠耗时需 async def

最后一条再强调一次。本教程所有例子都是同步函数:扇出改变的是合并语义(同一超步、彼此看不见对方的写入),不是自动让两个 time.sleep(1) 变成只等 1 秒。入门先把结构跑对,再考虑异步。

9. 动手练习 #

只看不改,印象留不住。下面五道都建立在已经能跑的代码上,改完对照 §8 那张表。

  1. 复现抢键。 把 §4.1 的 log 改成普通 str(去掉 Annotated),跑一次,确认出现 InvalidUpdateError。再改回来。
  2. 给静态扇出加第三个工人。 在 start 上再连一个 traffic 节点,三路都扇入 merge。看 merge 打印时 log 里是不是三条都在。
  3. 列表扇出返回空列表。 用 §6.2 的图,给 pick 加一句:列表为空时返回 "idle",并注册一个 idle 节点往 log 里写 "请先勾选任务"。
  4. Send 漏字段。 把 §7.3 里 Send 的 "topic": state["topic"] 删掉,看工人报什么错。改回去。
  5. 自己换一份清单。 把主题改成一篇文章的三个段落(用三个字符串模拟),工人返回「这一段的字数」,reduce 打印总字数。不调模型,纯 Python 就能做完。

10. 小结 #

扇出不是新 API,是一种连边方式:

  1. 从一个节点同时触发多个下游,它们进入同一个超步,彼此看不见对方刚写的值。
  2. 静态扇出用多次 add_edge;列表扇出用条件边返回节点名列表;Send 扇出在列表的基础上给每个工人一份专属输入。
  3. 扇入就是多条边指向同一个下游。下游会等所有上游到齐,再读已经合并好的状态。
  4. 并行写入的字段必须有 reducer。 没有就 InvalidUpdateError。列表用 operator.add 就够入门。
  5. Send 的第二参数是工人能看到的全部状态。 漏带字段会 KeyError,主图其它键不会自动出现。
  6. 空列表等于一个下游都不走,图悄悄结束。过滤分支时要给空结果留兜底。
  7. 同步函数的扇出不保证更快,只保证合并语义正确。先把图画对,再谈异步加速。

三种写法怎么选:分支固定用静态;人数会变但输入相同用列表;「清单里每一项一个工人」用 Send。入门把 §4、§5、§7.3 三份代码跑通、改通,扇出就算学会了。