Appearance
并行与 MapReduce
有些任务可以同时做:并行总结 5 段文本、并行给 3 个文档打分。LangGraph 支持节点扇出(fan-out)到多个并行执行,再扇入(fan-in)汇聚结果。本篇讲两种并行方式和一个经典 MapReduce 模式。
一、并行执行的场景
| 场景 | 说明 |
|---|---|
| 多源检索 | 同时查 3 个知识库,再合并 |
| 并行打分 | 给多个候选答案同时评分 |
| MapReduce | 大量输入先 map 并行处理,再 reduce 汇总 |
| 加速 | 把串行的独立步骤并行跑,省时间 |
并行的核心难点是:多个执行要往同一状态字段写结果,必须用 reducer 汇聚,否则冲突。
二、扇出与扇入
扇出(fan-out):一个节点执行完,下一步同时去多个节点。扇入(fan-in):多个并行节点都完成后,汇聚到一个节点继续。
mermaid
flowchart LR
START --> A[分发] --> B1[worker1]
A --> B2[worker2]
A --> B3[worker3]
B1 --> C[汇聚]
B2 --> C
B3 --> C
C --> END用条件边返回列表实现多路分发
add_conditional_edges 的路由函数可以返回一个列表,表示"同时去这些节点":
python
from typing import TypedDict, Annotated
from operator import add
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
query: str
results: Annotated[list[str], add] # 关键:reducer 汇聚
def fanout_dispatch(state) -> list[str]:
# 返回列表,表示同时分发到三个 worker
return ["worker_a", "worker_b", "worker_c"]
def worker_a(state):
return {"results": ["A的结果"]}
def worker_b(state):
return {"results": ["B的结果"]}
def worker_c(state):
return {"results": ["C的结果"]}
def merge(state):
return {"results": state["results"] + ["已合并"]}
g = StateGraph(State)
g.add_node("worker_a", worker_a)
g.add_node("worker_b", worker_b)
g.add_node("worker_c", worker_c)
g.add_node("merge", merge)
g.add_conditional_edges(START, fanout_dispatch, ["worker_a", "worker_b", "worker_c"])
# 三个 worker 都汇聚到 merge
g.add_edge("worker_a", "merge")
g.add_edge("worker_b", "merge")
g.add_edge("worker_c", "merge")
g.add_edge("merge", END)
app = g.compile()
result = app.invoke({"query": "x", "results": []})
print(result["results"]) # ['A的结果', 'B的结果', 'C的结果', '已合并']要点:
- 路由函数返回列表而非单个字符串。
results字段必须挂Annotated[list[str], add],否则三个 worker 同时写会报InvalidUpdateError。- 三个 worker 是并行执行的(在异步图里真正并发)。
三、MapReduce 模式
MapReduce:map 阶段对每个输入项并行处理,reduce 阶段汇总。固定数量用上面的条件边列表即可;动态数量(运行时才知道几个)用 Send API,见 动态图与 Send。
下面是固定 MapReduce:并行总结多段文本再合并。
python
from typing import TypedDict, Annotated
from operator import add
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
chunks: list[str] # 输入:多段文本
summaries: Annotated[list[str], add] # map 输出汇聚
final: str # reduce 输出
# 为每个 chunk 造一个 worker 节点(这里固定 3 个)
def make_summarizer(idx):
def _sum(state):
chunk = state["chunks"][idx]
return {"summaries": [f"摘要{idx}:{chunk[:6]}..."]}
return _sum
def reduce_node(state):
return {"final": " || ".join(state["summaries"])}
g = StateGraph(State)
g.add_node("sum0", make_summarizer(0))
g.add_node("sum1", make_summarizer(1))
g.add_node("sum2", make_summarizer(2))
g.add_node("reduce", reduce_node)
def dispatch(state) -> list[str]:
return ["sum0", "sum1", "sum2"]
g.add_conditional_edges(START, dispatch, ["sum0", "sum1", "sum2"])
g.add_edge("sum0", "reduce")
g.add_edge("sum1", "reduce")
g.add_edge("sum2", "reduce")
g.add_edge("reduce", END)
app = g.compile()
chunks = ["第一段很长的文本...", "第二段也很长的文本...", "第三段依旧很长..."]
result = app.invoke({"chunks": chunks, "summaries": [], "final": ""})
print(result["final"])
# 摘要0:第一段很... || 摘要1:第二段也很... || 摘要2:第三段依旧...mermaid
flowchart LR
START --> D{dispatch}
D --> S0[sum0]
D --> S1[sum1]
D --> S2[sum2]
S0 --> R[reduce]
S1 --> R
S2 --> R
R --> END这种固定写法要求运行前就知道有几个 chunk。运行时才知道数量(比如检索结果个数不定)要用 Send,见下一篇。
四、并行汇聚的 reducer 选择
并行节点都写同一个字段,reducer 决定怎么合并:
| reducer | 行为 | 适合 |
|---|---|---|
operator.add | list 追加 / int 累加 | 收集所有结果 |
add_messages | 消息追加(带 id 去重) | 消息历史 |
| 自定义函数 | 任意合并逻辑 | 取最大值、去重等 |
自定义 reducer 示例(取分数最高):
python
def keep_best(left, right):
# left/right 都是 [(score, text), ...]
combined = (left or []) + (right or [])
return max(combined, key=lambda x: x[0]) if combined else []
class State(TypedDict):
candidates: Annotated[list, keep_best]五、执行顺序的非确定性
并行节点执行顺序不保证。如果你依赖"worker_a 必须先于 worker_b",那就不该并行。并行节点之间不应有数据依赖,只通过汇聚 reducer 交换数据。
六、常见踩坑
踩坑 1:并行写同字段没挂 reducer
三个 worker 都返回 {"results": [...]},但 results 没挂 reducer,LangGraph 报 InvalidUpdateError: At key 'results': Can receive only one value per step。并行汇聚字段必须挂 reducer。
踩坑 2:扇入后漏接某个 worker
python
g.add_edge("worker_a", "merge")
g.add_edge("worker_b", "merge")
# 忘了 worker_c → mergeworker_c 执行完没去向,图结构非法,compile 会报错或 merge 等不到它。扇入要接全所有扇出的 worker。
踩坑 3:并行节点读共享状态时机错乱
并行 worker 同时读 state["chunks"] 是安全的(只读),但如果某个 worker 还修改了 chunks,其他 worker 读到的可能是修改前或后——非确定。并行节点只读共享输入,写各自的汇聚字段。
踩坑 4:以为同步图会真并行
同步 invoke 里并行节点其实是"顺序执行"的(Python GIL + 同步循环),只是逻辑上并行。要真正并发用 astream/ainvoke + async def 节点。
踩坑 5:map 数量动态却用固定写法
chunk 个数运行时才确定,却写死了 sum0/sum1/sum2。运行时 chunk 有 5 个就只能处理 3 个。动态数量必须用 Send,见 动态图与 Send。
七、小结
- 扇出用条件边返回列表,扇入用多条
add_edge汇聚到一个节点。 - 并行汇聚字段必须挂 reducer(
operator.add/add_messages/ 自定义)。 - MapReduce 固定数量用条件边列表,动态数量用 Send(下一篇)。
- 并行节点之间不应有数据依赖,执行顺序非确定。
- 真并发要异步图 +
async def节点。
下一篇 动态图与 Send 讲运行时才知道数量的并行。