Day 16 / 共 60 天 · 阶段3 控制流

Send 动态扇出与 map-reduce

条件边和 Command 都是"选哪个下游"。Send 更进一步:它能在运行时决定拉起多少个下游、每个还带不同的输入。这正是 map-reduce——把一个列表炸成 N 个并行任务、跑完再聚合。今天读 types.pySend 和官方 bench 例子,看这套动态扇出的底层。

📍 你在 60 天里的位置(阶段3 · 控制流 D13-18)
D13 边 D14 Branch D15 Command D16 Send D17 递归上限 D18 START/END
💡 一个类比先兜住今天(餐厅出餐) 普通边像"点一道菜、上一道菜"。Send一桌客人点了 5 份不同口味的披萨,主厨一声令下同时开 5 个炉子——每个炉子(同一个"烤披萨"节点)拿到不同的配料单(各自的 arg),并行开烤,最后 5 份一起端上桌汇总。Send("烤披萨", {"口味":"榴莲"}) 就是"给烤披萨节点派一单,配料是榴莲"。要几份就发几个 Send——数量运行时才知道(客人现场点)。这就是"动态扇出"。
L01

痛点:下游数量运行时才知道

🤔 痛点:我事先不知道要并行几个分支"给每个用户提到的主题各生成一个笑话"——主题有几个?运行时读了 state 才知道,可能 2 个也可能 20 个。普通边和条件边只能连到建图时就写死的固定节点,没法表达"动态开 N 个并行任务、每个喂不同数据"。你总不能建图时就写 20 个一样的节点吧。
💡 本质:Send = "运行时生成的、带自定义输入的一次性任务"Send(node, arg) 表示"请额外跑一次 node,它的输入是 arg(不是主状态!)"。路由函数返回一个 [Send, Send, ...] 列表,列表多长就并行开多少个任务。每个任务拿到自己那份 arg,互不干扰。跑完各自往共享通道写结果,靠 reducer 聚合(Day 09 的 operator.add)。数量由列表长度动态决定——这就是它比条件边强的地方。
条件边扇出Send 扇出
下游数量固定(path_map 里那几个)运行时动态(列表长度)
每个下游的输入都是同一份主状态各自不同的 arg
典型用途if/else 分流map-reduce、并行子任务
L02

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 全相等才算相等。
💡 本质:Send 是"数据包",不是"边"普通边/条件边描述的是"图的固定拓扑";Send 是运行时临时生成的"投递给某节点的数据包"。所以它不存在 builder 的任何容器里,而是路由函数当场造出来、当场发出去。图里甚至可以有一个"没有任何静态入边、只靠 Send 触发"的节点。
L03

官方例子:主题列表炸成 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"]}
控制流:map-reduce 扇出/聚合 STARTsubjects=[..] continue_to_jokes造 Send 列表(map) generate_joke(cats) generate_joke(dogs) …(N 个并行) 聚合operator.add
图注:一个函数造 N 个 Send(map),引擎并行拉起 N 份同一节点,各写共享通道,reducer 聚合(reduce)。
L04

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 数量,在这里落实。
💐 设计取舍①:为什么 Send 要绕道 TASKS 通道,而不直接调节点? 因为 LangGraph 的执行是严格分超步的(BSP,阶段 4)——所有"下一步要跑什么"必须先落进通道、在超步边界统一收集,才能被检查点保存、才能支持中断恢复。如果 Send 绕过通道直接调节点,就破坏了"每一步都可存档、可重放"的核心保证。代价是多一层通道中转,换来的是动态扇出也能崩溃恢复——跑到一半挂了,重启能从"还剩哪些 Send 没处理"继续。这是可靠性优先的取舍。
L05

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_idpacket.node + idx + step 等算出稳定的任务 id——同样的 Send 在重放时算出同样的 id,保证幂等(不会因重启重复执行同一 Send,阶段 6 会深挖)。
🚫 坑:Send 的 arg 结构要匹配目标节点的输入 schemaSend 喂给节点的是 arg不是主图状态。所以目标节点读到的字段完全取决于你在 arg 里放了什么。fanout 例子里目标子图专门声明了 JokeInput = {"subject": str} 来接这个单数 arg。如果 arg 结构和目标节点期望的对不上,节点里就会 KeyError。发 Send 前想清楚"目标节点需要什么,我 arg 里就给什么"。
另外还有个 sanitize_untracked_values_in_send_algo.py:1442):Send 的 arg 若含 UntrackedValue(不持久化的字段),检查点前会先剔除,避免把"不该存的东西"写进存档。这是 Send 与持久化协作的一个细节。
L06

reduce 聚合与常见边界

map 出去 N 个任务,它们并行往同一个通道写结果。能安全并行写,全靠 Day 07-10 讲的 reducer。没有 reducer 会怎样?

🚫 最经典的坑:聚合字段忘了加 reducer如果 jokes 声明成普通 list[str](没 Annotated[..., operator.add]),那 N 个并行任务同一超步写同一通道,就会触发 InvalidUpdateError("多个任务并发更新同一通道但通道不支持")。因为默认 LastValue 通道只允许一个写者。map-reduce 的聚合字段必须配累加型 reducer——这是 Send 用法里第一大翻车点。

👶 小白:如果某个 Send 任务失败了,其它并行任务会被连累吗?

👨‍🏫 老师:默认情况下,同一超步里任一任务抛异常会导致这一步失败(可配重试策略缓解,阶段 4 Day 25 讲 RetryPolicy)。但因为一切都在通道里、有检查点,重启能从这一超步重来,已成功写入的不丢。另外每个 Send 可单独设 timeout(L02 见过),给"某个子任务特别慢"留了单独兜底的口子。

📝 补充:Send 也能配合 Command 用 Day 15 我们见过 Command(goto=Send(...))。也就是说节点可以既更新状态、又动态扇出return Command(update={"step":"mapping"}, goto=[Send("worker", x) for x in items])。Command 负责"更新+发起",Send 负责"每个任务喂什么",两者正交组合。
L07

今日小结 + 动手 + 明日预告

🧠 今天你应该能回答

  • 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 的必要性)
明天预告 · Day 17:条件边能形成循环(ReAct 的 tools→agent),Send 能一次炸出成千上万个任务——万一循环停不下来、或扇出爆炸怎么办?Day 17 讲递归上限 recursion_limitGraphRecursionError:引擎如何用"超步计数 + 上限"给失控的图踩刹车。
← Day 15 Command 对象 Day 17 · 循环与递归上限 →