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