Skip to content

Send:节点扇出的 PUSH 任务原语

源码版本1.2.9

职责

Send 是 LangGraph 的「向指定节点投递任务」原语,定义在 libs/langgraph/langgraph/types.py(Send class:664)。一个 Send(node="foo", arg={"k": v}) 表示「在下一步用 arg 作为输入调用节点 foo」,跟普通节点触发的区别在于:普通节点(PULL)是从通道订阅关系被动触发,Send(PUSH)是上一节点主动声明「我要让 foo 跑一次,带上这个特定 arg」。

它的位置在条件边和 TASKS 通道之间:add_conditional_edges(START, lambda s: [Send("gen", {"x": 1}), Send("gen", {"x": 2})]) 把条件函数返回的 Send 列表写进 TASKS 通道(goto as sends:62),下一轮 prepare_next_taskstasks_channel.get() 里把它们取出来,每个 Send 组装成一个独立的 PregelExecutableTask(PUSH tasks:421)。

设计动机

  • map-reduce 的天然表达:经典场景是「同一节点并行跑 N 次,每次输入不同」。Send("generate", {"subject": s}) 在条件边里返回一个 list,每个 Send 一次独立调用,Pregel 自然并发执行,reducer 把 N 次结果合回主状态。不需要在图里手动展开 N 条边。
  • PUSH / PULL 二分:节点的触发分两类——PULL(被动,由订阅的通道更新触发)、PUSH(主动,由 Send 显式投递)。prepare_next_tasks 先消费 TASKS 通道里的 PUSH 任务,再按 trigger_to_nodes 算 PULL 任务(PUSH tasks:440)`。
  • arg 不必是主状态:Send.arg 可以是任意对象,不强制符合图的主状态 schema——目标节点的 input_schema 决定 arg 的形状,允许「向同一节点投递异构输入」,例如 Send("review", ReviewContext(doc=...)) vs Send("review", DefaultContext())
  • SendCommand.goto 等价:Command(goto=Send(...))map_command 里走同一条 TASKS 通道(goto as sends:62),所以条件边和节点返回值两种入口最终都汇到 PUSH 任务调度,统一抽象。
  • 每个 Send 是独立 task:Send 列表里 N 个元素就是 N 个独立 PregelExecutableTask,有各自的 task_id 和 path,可以分别 retry、分别出错处理——不会因为其中一个失败影响其他并发的 Send。

关键文件

  • Send class:664Send(node, arg, *, timeout=None) 数据类,__slots__ = ("node", "arg", "timeout")
  • Send __init__:718 — 构造函数,timeoutTimeoutPolicy.coerce
  • Send __hash__:739hash((node, arg, timeout)),作为 futures dict key 的去重依据。
  • prepare_next_tasks:392 — 主调度函数,先消费 PUSH 任务再算 PULL 候选节点。
  • PUSH tasks:440 — 从 TASKS 通道取出 Send 序列,每个调 prepare_single_task 组装成 PregelExecutableTask
  • PULL tasks:470 — PULL 路径,按 trigger_to_nodes 算出哪些节点该跑,不消费 TASKS
  • goto as sends:62map_commandCommand.gotoSend 写进 TASKS 通道,统一 PUSH 路径。
  • add_conditional_edges:969 — 条件边入口,条件函数返回 SendSend 列表。
  • PregelExecutableTask:627 — 任务的可执行形态,持有 name / input / proc / writes / triggers / path 等。
  • PregelRunner.tick:176 — 并发执行本步所有 PUSH + PULL 任务,Send 投递出来的任务在这里被跑掉。

数据流

Send 真正被消费是在 prepare_next_tasks(prepare_next_tasks:392)的开头:

python
input_cache: dict[INPUT_CACHE_KEY_TYPE, Any] = {}
checkpoint_id_bytes = binascii.unhexlify(checkpoint["id"].replace("-", ""))
null_version = checkpoint_null_version(checkpoint)
tasks: list[PregelTask | PregelExecutableTask] = []
# Consume pending tasks
tasks_channel = cast(Topic[Send] | None, channels.get(TASKS))
if tasks_channel and tasks_channel.is_available():
    for idx, _ in enumerate(tasks_channel.get()):
        if task := prepare_single_task(
            (PUSH, idx),
            None,
            checkpoint=checkpoint,
            checkpoint_id_bytes=checkpoint_id_bytes,
            checkpoint_null_version=null_version,
            pending_writes=pending_writes,
            processes=processes,
            channels=channels,
            managed=managed,
            config=config,
            step=step,
            stop=stop,
            for_execution=for_execution,
            store=store,
            checkpointer=checkpointer,
            manager=manager,
            input_cache=input_cache,
            cache_policy=cache_policy,
            retry_policy=retry_policy,
        ):
            tasks.append(task)

TASKS 通道是 Topic[Send]——一个支持多次写入、多次读取的通道类型(Topic:23)。每个 Sendtasks_channel.get() 里被取出来,(PUSH, idx) 作为 task_id 前缀传给 prepare_single_task,后者从 processes[send.node] 拿出对应的 PregelNode 配置(含 proc / writers / retry_policy 等),把 send.arg 作为 input 组装成 PregelExecutableTaskSend.arg 不进主状态通道——它直接作为 task.input 传给节点 proc,节点函数拿到的是这个 arg,不是图主状态的全量快照。

Send 进入 TASKS 通道有两条路径:(1) add_conditional_edges 配的条件函数返回 Send 列表,attach_branch 把它们写进 TASKS;(2) 节点函数返回 Command(goto=Send(...)),map_command_io.pySend 写进 TASKS(goto as sends:62)。两条路径殊途同归,都是写进 TASKS 通道、下一轮 prepare_next_tasks 消费。

边界与失败

  • Send 必须 target 已注册节点:processes[send.node] 不存在时 prepare_single_task 返回 None,该 Send 被静默丢弃——不会 raise,但任务不跑。调试时如果发现 Send 没生效,先检查目标节点名是否在 add_node 里注册过。
  • Send.arg 不走状态 reducer:Sendarg 直接作为节点 input,不经过主状态通道的 reducer 合并。如果目标节点期望读主状态,需要在该节点函数里通过 config[CONFIG_KEY_READ]state: State 参数拿——arg 只是「额外输入」。
  • Send 列表去重靠 hash:Send.__hash__(node, arg, timeout) 三元组算(Send __hash__:739),如果同一节点的两个 Send 的 arg 是不可哈希对象(如 dict),会直接 TypeError——Send_call 的 dedup 路径里被当作 dict key 用。
  • 同一节点可被 PULL 和 PUSH 同时触发:如果 Send("foo", arg) 投递 foo,同时 foo 的某个 trigger 通道被更新,prepare_next_tasks 会把 PUSH 任务和 PULL 任务都列出来,PregelRunner.tick 一并发跑——这是允许的,但要注意节点函数可能被同一步调用两次。
  • Send.timeout 单独生效:Send(node, arg, timeout=10) 只对该 PUSH 任务生效,不影响目标节点的默认 timeout 配置(Send __init__:718)。
  • TASKS 通道是 Topic 不是 LastValue:Topic 支持多次写入累积(Topic:23),所以同一超步多个节点都能往里写 Send,下一轮一次性消费完——这是 map-reduce 的基础。

小结

Send 是 PUSH 任务的载体,跟普通 PULL 节点互为补充。条件边和 Command.goto 两条入口汇到同一个 TASKS 通道,由 prepare_next_tasks 统一调度。要追 task 怎么被并发跑看 /func/concurrency,要追条件边到节点的关系看 libs/langgraph/langgraph/graph/_branch.py。对照官方资料:LangGraph 文档 · README