Send 动态扇出与 map-reduce
条件边和 Command 都是"选哪个下游"。Send 更进一步:它能在运行时决定拉起多少个下游、每个还带不同的输入。这正是 map-reduce——把一个列表炸成 N 个并行任务、跑完再聚合。今天读 types.py 的 Send 和官方 bench 例子,看这套动态扇出的底层。
Send 像一桌客人点了 5 份不同口味的披萨,主厨一声令下同时开 5 个炉子——每个炉子(同一个"烤披萨"节点)拿到不同的配料单(各自的 arg),并行开烤,最后 5 份一起端上桌汇总。Send("烤披萨", {"口味":"榴莲"}) 就是"给烤披萨节点派一单,配料是榴莲"。要几份就发几个 Send——数量运行时才知道(客人现场点)。这就是"动态扇出"。痛点:下游数量运行时才知道
Send(node, arg) 表示"请额外跑一次 node,它的输入是 arg(不是主状态!)"。路由函数返回一个 [Send, Send, ...] 列表,列表多长就并行开多少个任务。每个任务拿到自己那份 arg,互不干扰。跑完各自往共享通道写结果,靠 reducer 聚合(Day 09 的 operator.add)。数量由列表长度动态决定——这就是它比条件边强的地方。| 条件边扇出 | Send 扇出 | |
|---|---|---|
| 下游数量 | 固定(path_map 里那几个) | 运行时动态(列表长度) |
| 每个下游的输入 | 都是同一份主状态 | 各自不同的 arg |
| 典型用途 | if/else 分流 | map-reduce、并行子任务 |
Send 类走读
源码 types.py:664,本体极简——就是"目标节点 + 要喂给它的参数"的封装:
# types.py:664
class Send:
__slots__ = ("node", "arg", "timeout") # types.py:711
node: str # 目标节点名
arg: Any # 喂给它的输入(可与主状态完全不同)
timeout: TimeoutPolicy | None
def __init__(self, /, node: str, arg: Any, *,
timeout: float | timedelta | TimeoutPolicy | None = None) -> None:
self.node = node
self.arg = arg
self.timeout = TimeoutPolicy.coerce(timeout) # 把数字/timedelta 统一成策略对象
def __hash__(self) -> int: # types.py:738
return hash((self.node, self.arg, self.timeout))
def __eq__(self, value: object) -> bool: # types.py:746
return (isinstance(value, Send) and self.node == value.node
and self.arg == value.arg and self.timeout == value.timeout)
__slots__用 __slots__ 而非普通 __dict__:省内存 + 禁止乱加属性。map-reduce 一次可能创建成百上千个 Send,每个省几十字节就很可观。node目标节点名——注意是字符串,运行时才拿去处理进程表里查真身。arg喂给目标节点的输入。它可以和主图状态结构完全不同!比如主状态有 subjects: list,而每个 Send 的 arg 是 {"subject": "cats"}——单数、只给一份。这是 map 阶段"切分"的体现。timeout可给这个单独任务设超时,不设就用目标节点默认的。TimeoutPolicy.coerce 把裸数字/timedelta 归一成策略对象(又见"归一化")。__hash__ / __eq__显式实现让 Send 可哈希、可比较。为什么?因为要放进通道、参与去重和检查点比对。两个 Send 的 node、arg、timeout 全相等才算相等。官方例子:主题列表炸成 N 个并行子图
官方性能基准 bench/fanout_to_subgraph.py:16 就是标准 map-reduce。核心那句路由函数:
# bench/fanout_to_subgraph.py:11
class OverallState(TypedDict):
subjects: list[str]
jokes: Annotated[list[str], operator.add] # ← reduce 靠这个累加器
# bench/fanout_to_subgraph.py:16 —— map:给每个主题发一个 Send
async def continue_to_jokes(state: OverallState):
return [Send("generate_joke", {"subject": s}) for s in state["subjects"]]
builder.add_conditional_edges(START, continue_to_jokes) # 条件边返回 Send 列表
builder.add_node("generate_joke", subgraphc) # 目标是一个子图
builder.add_edge("generate_joke", END)
[Send(...) for s in subjects]map 阶段:列表推导,主题有几个就造几个 Send。每个 Send 的 arg 是 {"subject": 单个主题}——把大列表切成一份份。add_conditional_edges(START, continue_to_jokes)把这个"造 Send 列表"的函数挂成 START 的条件边。回顾 Day 14 _finish:返回列表 → 一次扇出多个;返回项是 Send → 原样保留投递。两条规则在这里合体。target 是 subgraphc目标节点是个编译好的子图!Send 不仅能拉起普通节点,还能并行拉起 N 个子图实例,每个跑自己那份 subject。jokes: Annotated[list, operator.add]reduce 阶段:所有并行任务都往 jokes 写自己那条笑话,operator.add(Day 09)把它们拼成一个大列表。map 出去、reduce 回来,闭环。{"subjects": ["cats", "dogs"]}:→
continue_to_jokes 返回 [Send("generate_joke", {"subject":"cats"}), Send("generate_joke", {"subject":"dogs"})]→ 引擎并行拉起 2 个 generate_joke 子图,分别喂 cats / dogs
→ 各自产出
{"jokes": ["Joke about cats"]} / {"jokes": ["Joke about dogs"]}→
operator.add 聚合 → 最终 {"subjects":["cats","dogs"], "jokes":["Joke about cats","Joke about dogs"]}Send 落进 TASKS 通道
Send 被发出后,会写进一个特殊通道 TASKS。无论来自条件边(Day 14 _finish)还是 Command 的 goto(Day 15 map_command),最终都归到这里。回看 pregel/_io.py:67:
# pregel/_io.py:66(map_command 里处理 goto 是 Send 的分支)
for send in sends:
if isinstance(send, Send):
yield (NULL_TASK_ID, TASKS, send) # ← Send 写进 TASKS 通道
准备下一超步任务时,prepare_next_tasks 会先把 TASKS 通道里所有待处理的 Send 消费掉,每个变成一个 PUSH 任务。pregel/_algo.py:442:
# pregel/_algo.py:442
tasks_channel = cast(Topic[Send] | None, channels.get(TASKS))
if tasks_channel and tasks_channel.is_available():
for idx, _ in enumerate(tasks_channel.get()): # 遍历每个待处理 Send
if task := prepare_single_task(
(PUSH, idx), # ← PUSH 任务:由 Send 推来
...
):
TASKS 通道类型 Topic[Send]TASKS 是个 Topic 通道(Day 30)——能累积一批 Send,而不是覆盖。这样一次发的 N 个 Send 都留得住。(PUSH, idx)每个 Send 生成一个 PUSH 类型的任务路径,idx 是它在这批 Send 里的下标。PUSH vs PULL:PULL 任务是被边"拉"起来的(普通节点);PUSH 任务是被 Send"推"来的(动态任务)。这是引擎里两大任务来源。enumerate(tasks_channel.get())有几个 Send 就 enumerate 出几个 → 生成几个 PUSH 任务。并行度 = Send 数量,在这里落实。PUSH 任务如何被拉起
每个 Send 变成 PUSH 任务,由 prepare_push_task_send 把它组装成真正可执行的任务。pregel/_algo.py:938:
# pregel/_algo.py:938(裁剪核心)
def prepare_push_task_send(task_path, ..., channels, processes, step, ...):
if len(task_path) == 2:
idx = cast(int, task_path[1]) # 第几个 Send
if not channels[TASKS].is_available():
return
sends: Sequence[Send] = channels[TASKS].get()
if idx < 0 or idx >= len(sends):
return # ← 边界:下标越界,安静跳过
packet = sends[idx]
if not isinstance(packet, Send):
logger.warning(f"Ignoring invalid packet type {type(packet)} ...")
return # ← 边界:不是 Send,警告跳过
if packet.node not in processes:
logger.warning(f"Ignoring unknown node name {packet.node} ...")
return # ← 边界:目标节点不存在,警告跳过
proc = processes[packet.node] # 找到目标节点的执行进程
# ...(用 packet.node、idx 等算出确定性 task_id)
sends[idx]按下标取出这一个 Send。packet.arg 稍后会作为该节点这次运行的输入——注意是 arg 而非主状态,这是 Send 与普通节点最大不同。idx 越界 return边界①:检查点恢复等场景下 TASKS 可能对不上,越界就安静返回,不崩。不是 Send / 未知节点边界②③:脏数据(非 Send)或 Send 指向了图里不存在的节点,都只 warning + return,宁可漏跑一个也不让整图崩。防御式编程。确定性 task_id用 packet.node + idx + step 等算出稳定的任务 id——同样的 Send 在重放时算出同样的 id,保证幂等(不会因重启重复执行同一 Send,阶段 6 会深挖)。arg,不是主图状态。所以目标节点读到的字段完全取决于你在 arg 里放了什么。fanout 例子里目标子图专门声明了 JokeInput = {"subject": str} 来接这个单数 arg。如果 arg 结构和目标节点期望的对不上,节点里就会 KeyError。发 Send 前想清楚"目标节点需要什么,我 arg 里就给什么"。sanitize_untracked_values_in_send(_algo.py:1442):Send 的 arg 若含 UntrackedValue(不持久化的字段),检查点前会先剔除,避免把"不该存的东西"写进存档。这是 Send 与持久化协作的一个细节。reduce 聚合与常见边界
map 出去 N 个任务,它们并行往同一个通道写结果。能安全并行写,全靠 Day 07-10 讲的 reducer。没有 reducer 会怎样?
jokes 声明成普通 list[str](没 Annotated[..., operator.add]),那 N 个并行任务同一超步写同一通道,就会触发 InvalidUpdateError("多个任务并发更新同一通道但通道不支持")。因为默认 LastValue 通道只允许一个写者。map-reduce 的聚合字段必须配累加型 reducer——这是 Send 用法里第一大翻车点。👶 小白:如果某个 Send 任务失败了,其它并行任务会被连累吗?
👨🏫 老师:默认情况下,同一超步里任一任务抛异常会导致这一步失败(可配重试策略缓解,阶段 4 Day 25 讲 RetryPolicy)。但因为一切都在通道里、有检查点,重启能从这一超步重来,已成功写入的不丢。另外每个 Send 可单独设 timeout(L02 见过),给"某个子任务特别慢"留了单独兜底的口子。
Command(goto=Send(...))。也就是说节点可以既更新状态、又动态扇出:return Command(update={"step":"mapping"}, goto=[Send("worker", x) for x in items])。Command 负责"更新+发起",Send 负责"每个任务喂什么",两者正交组合。今日小结 + 动手 + 明日预告
🧠 今天你应该能回答
- Send 解决了条件边解决不了的什么问题?(下游数量运行时才知道、每个喂不同输入)
- Send 三个属性?为什么用
__slots__?(node/arg/timeout;省内存、量大) arg和主状态的关系?(可完全不同,是"切分出的一份"输入)- Send 发出后写进哪个通道、变成什么任务?(
TASKS通道 → PUSH 任务) - 为什么绕道 TASKS 通道而不直接调节点?(保证超步边界可存档/可恢复)
- map-reduce 的聚合字段必须配什么?(累加型 reducer,如
operator.add,否则 InvalidUpdateError)
✋ 10 分钟动手
# 1. 读 Send 类
sed -n '664,753p' libs/langgraph/langgraph/types.py
# 2. 读官方 fanout 基准(map-reduce 范本)
sed -n '1,60p' libs/langgraph/bench/fanout_to_subgraph.py
# 3. 跑一个最小 map-reduce(注意 jokes 的 reducer)
python - <<'PY'
from langgraph.graph import StateGraph, START, END
from langgraph.types import Send
from typing import TypedDict, Annotated
import operator
class S(TypedDict):
subjects: list; jokes: Annotated[list, operator.add]
def fan(state): return [Send("joke", {"s": x}) for x in state["subjects"]]
def joke(state): return {"jokes": [f"关于{state['s']}的笑话"]}
g = StateGraph(S)
g.add_node("joke", joke)
g.add_conditional_edges(START, fan) # 返回 Send 列表
g.add_edge("joke", END)
print(g.compile().invoke({"subjects": ["猫","狗","鱼"], "jokes": []}))
PY
# 4. 把 jokes 的 Annotated 去掉,观察 InvalidUpdateError(体会 reducer 的必要性)
GraphRecursionError:引擎如何用"超步计数 + 上限"给失控的图踩刹车。