子图与并行任务

前几章的图都是"一条主线"。真实业务往往是多条流水线的组合:文档要逐篇走"提取→校验→入库",数据要拆成多路并行处理再汇总。LangGraph 两招解决:已编译的图可以当节点塞进另一个图(子图复用),节点可以同时发散出多个任务(fan-out)再汇合(fan-in)。本章用一个综合可跑样例串起全部概念。

子图:把一段流程封装成"节点"

一张编译好的 StateGraph 放进另一张图,就成了后者的一个节点。适合封装重复子流程:主图只管编排,子流程内部细节被隔离。子图与外层通过状态字段通信——主图传入它需要的键,它把产出通过同样的键写回。

并行:fan-out / fan-in 与 map

  • fan-out(发散):一个节点同时连出多条路,多个任务并行执行。
  • fan-in(汇合):多个并行任务完成后,汇到同一个节点做汇总。
  • map(映射):对一批数据逐个执行同一子流程。LangGraph 用 Send 动态生成 N 个任务,等价于把列表"map"到子流程上,这是官方推荐的 map-reduce 写法。共享状态方面,同一轮并行节点读到同一份状态,各自写回后由 Reducer(add 等)合并,谁也不会覆盖谁。

综合样例:文档批处理(提取→校验→入库)

from operator import add
from typing import Annotated, TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send

DOCS = [
    "会议纪要\n确定下周一发布 v2.0,由后端同学负责上线。",
    "短",                        # 故意放一篇不合格的
    "周报\n本周完成登录联调,下周开始做支付模块。",
]

# ---------- 子图:单篇文档 提取 -> 校验 -> 生成入库记录 ----------
class DocState(TypedDict):
    doc: str
    title: str
    body: str
    ok: bool
    results: Annotated[list, add]     # Reducer:并行结果在此累加

def extract(state):
    lines = [s for s in state["doc"].splitlines() if s.strip()]
    title = lines[0] if lines else "(无标题)"
    body = "".join(lines[1:])
    print(f"    [提取] 标题={title!r},正文 {len(body)} 字")
    return {"title": title, "body": body}

def validate(state):
    ok = len(state["body"]) >= 8          # 正文不足 8 字判为不合格
    print(f"    [校验] {'通过' if ok else '不合格'}")
    return {"ok": ok}

def to_report(state):
    mark = "入库" if state["ok"] else "退回"
    return {"results": [f"{state['title']} | {mark}"]}

doc_builder = StateGraph(DocState)
doc_builder.add_node("extract", extract)
doc_builder.add_node("validate", validate)
doc_builder.add_node("to_report", to_report)
doc_builder.add_edge(START, "extract")
doc_builder.add_edge("extract", "validate")
doc_builder.add_edge("validate", "to_report")
doc_builder.add_edge("to_report", END)
doc_sub = doc_builder.compile()           # 编译成可复用子图

# ---------- 主图:列表 fan-out 到子图,再 fan-in 汇总 ----------
class BatchState(TypedDict):
    docs: list
    results: Annotated[list, add]

def fan_out(state):
    # 每篇文档一个 Send:动态并行执行子图(= map)
    return [Send("process_one", {"doc": d, "results": []}) for d in state["docs"]]

def summary(state):
    ok_n = sum("入库" in r for r in state["results"])
    print(f"\n共 {len(state['docs'])} 篇,入库 {ok_n} 篇:")
    for r in state["results"]:
        print("  -", r)

builder = StateGraph(BatchState)
builder.add_node("distribute", fan_out)
builder.add_node("process_one", doc_sub)   # 子图直接当节点用
builder.add_node("summary", summary)
builder.add_edge(START, "distribute")
builder.add_conditional_edges("distribute", fan_out)   # 返回 Send 列表 = 动态并行
builder.add_edge("process_one", "summary")             # 全部跑完才汇合(fan-in)
builder.add_edge("summary", END)
graph = builder.compile()

graph.invoke({"docs": DOCS, "results": []})

运行输出(并行任务的实际顺序可能略有差异):

    [提取] 标题='会议纪要',正文 16 字
    [校验] 通过
    [提取] 标题='短',正文 0 字
    [校验] 不合格
    [提取] 标题='周报',正文 17 字
    [校验] 通过
共 3 篇,入库 2 篇:
  - 会议纪要 | 入库
  - 短 | 退回
  - 周报 | 入库

这个样例里发生了什么

  • 三篇文档由 fan_out 生成三个 Send,并行进入同一个子图节点 process_one(各带自己的 doc,results 初始为空列表)。
  • 每篇在子图内部依次走 extract→validate→to_report,返回一条"标题 | 结果"记录。
  • 各路的 results 通过 add Reducer 累加成一份完整列表,全部完成后才触发 summary 做汇总——这就是 map-reduce 的骨架。
  • 想让它更"AI"?把 extract 换成模型节点(如第 18 章写法),或在校验不通过时接一个第 19 章的人工审批闸门即可,图的骨架不用动。

补充:静态并行(固定条数)不必用 Send,从一个节点连出多条边即可,例如"查天气"与"查日历"两条边都从 START 出发,汇合节点会等两边都完成——与上面 Send 动态发散是同一套 fan-out/fan-in 语义,适用于分支数在编译期就确定的场景。

常见问题

  • 一次 Send 太多任务会怎样?注意控制并发规模:可分批发送(处理完一批再发下一批)并留意执行上限,生产批处理建议限流分页。
  • 并行任务能写同一个普通字段吗?不能依赖——普通字段是覆盖语义,后完成者会覆盖先完成者;汇总数据必须用带 Reducer 的字段(如本例的 results)。
  • 子图与外层的 State 必须完全一样吗?不必。子图有自己的状态定义,父状态只透传两边都声明的键;私有键(title、body 等)只存在于子图内部,天然隔离。
  • 想让子图某环节变成 AI?把对应节点换成模型调用(第 18 章写法),子图对外接口不变,主图无需改动。

小结:编译好的图可以当作节点嵌入另一张图(子图复用),Send 把一批任务动态发散到子流程并靠 Reducer 汇合结果(map-reduce);主图负责编排与汇总、子图负责单条流水线细节,复杂批处理因此可以被拆成小而清晰的图。

笔记加载中…