Send:节点扇出的 PUSH 任务原语
职责
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_tasks 在 tasks_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=...))vsSend("review", DefaultContext())。Send与Command.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:664—Send(node, arg, *, timeout=None)数据类,__slots__ = ("node", "arg", "timeout")。Send __init__:718— 构造函数,timeout走TimeoutPolicy.coerce。Send __hash__:739—hash((node, arg, timeout)),作为futuresdict 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:62—map_command把Command.goto的Send写进TASKS通道,统一 PUSH 路径。add_conditional_edges:969— 条件边入口,条件函数返回Send或Send列表。PregelExecutableTask:627— 任务的可执行形态,持有name/input/proc/writes/triggers/path等。PregelRunner.tick:176— 并发执行本步所有 PUSH + PULL 任务,Send投递出来的任务在这里被跑掉。
数据流
Send 真正被消费是在 prepare_next_tasks(prepare_next_tasks:392)的开头:
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)。每个 Send 在 tasks_channel.get() 里被取出来,(PUSH, idx) 作为 task_id 前缀传给 prepare_single_task,后者从 processes[send.node] 拿出对应的 PregelNode 配置(含 proc / writers / retry_policy 等),把 send.arg 作为 input 组装成 PregelExecutableTask。Send.arg 不进主状态通道——它直接作为 task.input 传给节点 proc,节点函数拿到的是这个 arg,不是图主状态的全量快照。
Send 进入 TASKS 通道有两条路径:(1) add_conditional_edges 配的条件函数返回 Send 列表,attach_branch 把它们写进 TASKS;(2) 节点函数返回 Command(goto=Send(...)),map_command 在 _io.py 把 Send 写进 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:Send的arg直接作为节点 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。