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