并行与Send

一份报告要同时查「库存」「价格」「合规」三个数据源——串行太慢,手写 asyncio.gather 又难和状态机对齐。Send 让路由函数返回多个动态任务,LangGraph 在同一 superstep 并行调度,再用 reducer 合并。


1. 定位

维度 内容
角色 动态并行扇出:map 阶段
输入 → 输出 Send(node, partial_state) 列表 → 多路 invoke 同一 worker 节点
核心 API Sendadd_conditional_edges 返回 list[Send]
依赖 LangChain 无;worker 内可再调 Tool

2. 图拓扑

节点表

节点名 职责 读 State 写 State
fanout 生成 Send 列表 topics —(路由)
worker 处理单个 topic topic results(reducer 追加)
aggregate 汇总 results summary

边表

目标 类型 说明
START fanout 固定
fanout worker 条件 返回 [Send("worker", {"topic": t})]
worker aggregate 固定 全部 worker 完成后
aggregate END 固定

3. invoke 生命周期与 superstep

1
2
3
Superstep 1: fanout 执行,router 返回 3 个 Send → 并行启动 3 个 worker
Superstep 2: 3 个 worker 各自返回 partial update → reducer 合并 results
Superstep 3: aggregate 读合并后的 results → END

并行节点的写入必须通过 Annotated[…, reducer] 合并,否则后写覆盖先写。


4. 原理

4.1 Send 语义

Send("worker", {"topic": "A"})worker 投递一份子 state(可与全局 state 字段子集不同,取决于 schema)。

4.2 与静态并行区别

固定 add_edge 无法表达「列表长度运行时决定」;Send 是 动态 map

4.3 reducer 是并行前提

列表字段用 operator.add 或自定义 reducer 追加,计数可用求和 reducer。


5. 最小可运行示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
import operator
from typing import Annotated, TypedDict

from langgraph.constants import Send
from langgraph.graph import StateGraph, START, END


class State(TypedDict):
topics: list[str]
topic: str # worker 子上下文
results: Annotated[list[str], operator.add]
summary: str


def fanout(state: State) -> list[Send]:
return [Send("worker", {"topic": t}) for t in state["topics"]]


def worker(state: State) -> dict:
t = state["topic"]
return {"results": [f"ok:{t}"]}


def aggregate(state: State) -> dict:
joined = ",".join(sorted(state["results"]))
return {"summary": joined}


builder = StateGraph(State)
builder.add_node("worker", worker)
builder.add_node("aggregate", aggregate)
builder.add_conditional_edges(START, fanout, ["worker"])
builder.add_edge("worker", "aggregate")
builder.add_edge("aggregate", END)

graph = builder.compile()
out = graph.invoke({
"topics": ["a", "b"],
"topic": "",
"results": [],
"summary": "",
})
print(out["summary"]) # ok:a,ok:b(顺序可能因并行而变)

6. 执行追踪

输入:topics=["a","b"]

Superstep 活跃节点 results
1 fanout → 2× worker 并行 [] → 并行写入
2 worker 合并后 ["ok:a","ok:b"](顺序不定)
3 aggregate summary="ok:a,ok:b"

重要配置参数

参数 类型 / 默认 作用与影响 参考起点 配置指导
Send(node, arg) dataclass 动态任务 map 列表 arg 字段要在 State 有定义
add_conditional_edges(..., ["worker"]) list 声明 Send 目标 单 worker 名 与 Send 第一参数一致
Annotated[list, operator.add] reducer 并行追加 结果列表 无 reducer 会覆盖
worker 子 state dict 仅传必要字段 topic 避免传大对象
compile() 并行调度 默认 IO 密集考虑 async 篇
fanout 返回 [] 空列表 无 worker 边界处理 需连 aggregate 或 END

7. 易踩坑

  1. 并行写字段无 reducer:只保留最后一个 worker 的结果。
  2. Send 目标节点名拼错:compile 通过但 invoke 失败。
  3. 假设 results 顺序:并行完成顺序不确定,汇总前需 sort 或带 key。

小结

  • Send = 运行时决定并行份数的 map;同一 superstep 多实例跑同一节点。
  • 并行写入必须配 reducer;汇总节点等全部 worker 完成再跑。
  • 适合多检索、多文件解析;reduce 阶段用单独 aggregate 节点。

参考链接

-------------本文结束感谢您的阅读-------------