Skip to content

动态图与 Send

上一篇的并行要求"运行前知道数量"。但很多场景是运行时才知道——比如检索到 5 篇文档要并行打分,下次检索到 3 篇。这时用 Send API 动态创建并行执行。

一、Send 解决什么问题

Send 对象用于在运行时决定要并行启动几个执行、各自传什么输入。对比:

方式数量决定时机各路输入
条件边返回列表编译时固定都拿同一份状态
Send运行时动态每路可传不同输入

典型场景:动态 map——运行时才知道处理几个输入项,每个输入项并行处理。

二、Send 对象

Send(node_name, state) 表示"向 node_name 节点发送一份 state 作为输入"。路由函数返回一个 Send 列表,框架就会为每个 Send 启动一个并行执行。

python
from langgraph.types import Send

def dynamic_dispatch(state):
    # 运行时根据 state 里有多少文档,就发几个并行任务
    sends = []
    for i, doc in enumerate(state["docs"]):
        # 每个文档作为独立输入发给 grade 节点
        sends.append(Send("grade", {"doc": doc, "idx": i}))
    return sends

关键区别:条件边返回列表时,每个目标节点拿的是同一份完整父状态;而 Send 可以为每个目标单独构造输入。这让"每个并行实例处理不同数据"成为可能。

三、与条件边返回列表的区别

python
# 方式 A:条件边返回列表(共享同一份状态)
def route(state) -> list[str]:
    return ["worker_a", "worker_b"]   # 两个 worker 都拿 state

# 方式 B:Send 列表(各自独立输入)
def route(state) -> list[Send]:
    return [
        Send("worker", {"item": "任务1"}),
        Send("worker", {"item": "任务2"}),
        Send("worker", {"item": "任务3"}),   # 同一个节点,三个并行实例,不同输入
    ]

方式 B 的妙处:同一个节点 worker 被并行实例化 3 次,每次输入不同。这正是 map 的语义。

四、完整示例:动态并行给文档打分

检索到 N 篇文档,每篇并行打分,再汇总。N 运行时才知道。

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

# 父图状态
class State(TypedDict):
    query: str
    docs: list[str]                          # 运行时才知道有几个
    scores: Annotated[list[int], add]         # 并行打分结果汇聚
    summary: str

# 每个并行实例的状态(Send 传给 grade 节点的输入)
class GradeInput(TypedDict):
    doc: str
    idx: int

def retrieve(state):
    # 假装检索,返回数量不定的文档
    docs = [f"文档{i}" for i in range(3)]   # 这次 3 篇,下次可能 5 篇
    return {"docs": docs}

def grade(state: GradeInput) -> dict:
    # 每个 grade 实例只拿到自己那篇文档
    doc = state["doc"]
    # 假装打分:文档名越长分越高
    score = len(doc) * 10
    return {"scores": [score]}   # 写回父图的 scores(reducer 汇聚)

def summarize(state):
    return {"summary": f"共 {len(state['docs'])} 篇,总分 {sum(state['scores'])}"}

# 动态分发:根据 docs 数量发 N 个 Send
def dispatch(state) -> list[Send]:
    return [Send("grade", {"doc": doc, "idx": i})
            for i, doc in enumerate(state["docs"])]

g = StateGraph(State)
g.add_node("retrieve", retrieve)
g.add_node("grade", grade)
g.add_node("summarize", summarize)

g.add_edge(START, "retrieve")
# retrieve 后动态分发到 grade(数量不定)
g.add_conditional_edges("retrieve", dispatch, ["grade"])
# 所有 grade 实例完成后汇聚到 summarize
g.add_edge("grade", "summarize")
g.add_edge("summarize", END)
app = g.compile()

result = app.invoke({"query": "LangGraph", "docs": [], "scores": [], "summary": ""})
print(result["summary"])
# 共 3 篇,总分 90  (文档0=30, 文档1=30, 文档2=30)
mermaid
flowchart LR
    START --> R[retrieve]
    R -->|Send x N| G1[grade 实例1]
    R -->|Send x N| G2[grade 实例2]
    R -->|Send x N| Gn[grade 实例N]
    G1 --> S[summarize]
    G2 --> S
    Gn --> S
    S --> END

retrieve 里文档数量为 5,框架自动发 5 个并行 grade,无需改图结构。这就是动态并行的威力。

五、add_conditional_edges 配合 Send

add_conditional_edges 的第三参数(去向列表)在有 Send 时仍建议传,用于可视化校验。但 Send 的目标节点名可以不在静态列表里也能跑(框架运行时才解析)。规范起见列出可能目标。

python
g.add_conditional_edges("retrieve", dispatch, ["grade"])

六、何时用 Send,何时用固定并行

情况选择
运行前就知道并行数量条件边返回列表(更简单)
运行时才知道数量Send
每路输入不同Send
每路输入相同(共享状态)条件边返回列表
需要给同一节点启动多个实例Send(条件边列表对同名节点只能去一次)

简单原则:数量或输入是动态的 → Send;静态且共享输入 → 条件边列表。

七、Send 与 reducer 的配合

每个 Send 实例写回的字段仍然走父图的 reducer。上面例子里所有 grade 实例都返回 {"scores": [score]},靠 Annotated[list[int], add] 汇聚成完整列表。没有 reducer,多个实例同时写同一字段会冲突报错。这是 Send 最常踩的坑。

八、常见踩坑

踩坑 1:Send 的输入与目标节点 state 不兼容

Send("grade", {"doc": ..., "idx": ...}) 传的字典,字段要和 grade 节点读取的一致。如果 gradestate["document"] 但 Send 传的是 doc,节点拿到 None 或 KeyError。

python
# ❌ 字段名对不上
Send("grade", {"document": doc})   # grade 里读 state["doc"]

# ✅
Send("grade", {"doc": doc})

踩坑 2:汇聚字段没挂 reducer

多个 Send 实例并行写 scores,没挂 Annotated[..., add],报 InvalidUpdateError。和上一篇并行踩坑一样,任何多路写的字段都要 reducer

踩坑 3:dispatch 返回空列表

如果 state["docs"] 为空,dispatch 返回 [],图会直接结束(没有并行实例,也不会去 summarize)。需要空输入也走后续逻辑时,加个判断:

python
def dispatch(state):
    if not state["docs"]:
        return "summarize"   # 没文档直接去汇总
    return [Send("grade", ...) for ...]

踩坑 4:Send 传了整个父状态

python
Send("grade", state)   # ❌ 把整个父状态当输入,可能含不可序列化或冗余字段

Send 的输入应该只包含目标节点需要的那部分数据,精简且字段匹配。

踩坑 5:以为 Send 能控制执行顺序

Send 创建的并行实例执行顺序非确定。如果 reduce 阶段需要按原顺序排列结果,在 grade 返回里带上 idx,reduce 时按 idx 排序,而不是依赖到达顺序。

python
def grade(state) -> dict:
    return {"scores": [(state["idx"], score)]}  # 带 idx

def summarize(state):
    ordered = sorted(state["scores"], key=lambda x: x[0])  # 按原顺序排
    ...

九、小结

  • Send(node, state) 在运行时动态创建并行执行,每个实例可传不同输入。
  • 与条件边返回列表的区别:Send 支持动态数量 + 异构输入 + 同节点多实例。
  • 典型场景:动态 map(运行时才知道处理几个)。
  • 汇聚字段必须挂 reducer,输入字段要与目标节点读取的一致。
  • 并行顺序非确定,需要顺序时在数据里带索引,reduce 时排序。

至此 进阶特性 模块完成,可继续看 智能体实战项目