Skip to content

并行与 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.addlist 追加 / 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 → merge

worker_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 汇聚到一个节点。
  • 并行汇聚字段必须挂 reduceroperator.add / add_messages / 自定义)。
  • MapReduce 固定数量用条件边列表,动态数量用 Send(下一篇)。
  • 并行节点之间不应有数据依赖,执行顺序非确定。
  • 真并发要异步图 + async def 节点。

下一篇 动态图与 Send 讲运行时才知道数量的并行。