Appearance
动态图与 Send
上一篇的并行要求"运行前知道数量"。但很多场景是运行时才知道——比如检索到 5 篇文档要并行打分,下次检索到 3 篇。这时用
SendAPI 动态创建并行执行。
一、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 节点读取的一致。如果 grade 读 state["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 时排序。