Skip to content

PregelLoop:超步循環驅動器

源码版本1.2.9

職責

PregelLoop 是 LangGraph 執行時的「節拍器」。Pregel 本身只負責組裝資源、暴露 invoke / astream 入口,真正一圈圈跑圖的是一個 PregelLoop 實例:它持有當前檢查點 (checkpoint)、通道 (channel) 集合、pending_writesstep 計數器、status 狀態機,以及一份 tasks 字典——這些都是單個超步 (superstep) 內的全部上下文。外層 while loop.tick(): ... loop.after_tick() 的形態就建立在它身上。

它的兩個核心方法是 tick(tick:599)和 after_tick(after_tick:683)。前者負責「跑之前」:看是否超出步數上限、調 prepare_next_tasks 把這一步要執行的節點裝進 self.tasks、檢查中斷、處理 drain_requested、把上次斷點續跑留下的 pending_writes 重新貼回成功節點上。後者負責「跑之後」:收齊所有 task.writes,調 apply_writes 把它們合進通道,清掉 pending,落檢查點,再判斷要不要 interrupt_after。兩者一前一後,把「讀取輸入 → 排程 → 執行 → 寫回 → 持久化」這個 BSP 循環夾在中間。

換句話說,PregelLoop 不親自跑節點——那是 PregelRunner 的活;它只決定「這一步該看誰、跑完怎麼收尾、下一步還要不要繼續」。

設計動機

為什麼把循環邏輯拆成 tickafter_tick 兩段,而不是一個 step() 方法搞定?

  • 執行權交還給呼叫者:tick 回傳 True 後,呼叫方(Pregel.astream 的實作體,見 sync main loop:2964-2984)自己負責驅動 runner.tick 把這一批 task 跑完,再回頭調 loop.after_tick。這樣 PregelLoop 就不需要持有 PregelRunner 的參照,兩個物件之間是鬆耦合的協作關係,而不是嵌套。
  • 狀態機顯式化:status 欄位有 "input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps" 七種取值(status Literal:256-264)。每個狀態都對應一種退出原因,外部程式碼靠它判斷該 resume、該 raise、還是該 break,不用猜回傳值語意。
  • 斷點續跑語意清晰:is_replaying 只在第一步 tick 為真——表示當前這一步是「重放」上次中斷前已經跑完的 task,等 after_tick 跑完就關掉(is_replaying = False:716)。配合 _reapply_writes_to_succeeded_nodes 把檢查點裡殘留的成功 writes 貼回記憶體 task,就能讓成功節點不重跑、失敗節點進 error handler。
  • Drain 是外部干預入口:RunControl.drain_requested 由 stream 消費方設定(比如客戶端要中途停),tick 一看到就進 draining 狀態並回傳 False,讓外層 while 自然退出,而不是在節點中途硬殺。
  • Exit-mode delta 寫入與檢查點時序:_delta_write_futs(_delta_write_futs:207)收集所有 delta 通道的 put_writes future,_checkpointer_put_after_previous 必須先 drain 它們再寫下個 checkpoint——保證「寫產生 checkpoint,而不是 checkpoint 蓋過寫」的因果順序。

關鍵檔案

  • class PregelLoop:158 — 類別定義,欄位幾乎全是裸 attribute,沒有 property 包裝,便於子類別 SyncPregelLoop / AsyncPregelLoop 直接覆蓋。
  • status Literal:256-264 — 七種狀態列舉,"out_of_steps" 是步數上限觸發,"draining" 是外部請求停止,"done" 是正常結束。
  • tick method:599-681 — 一個完整 tick 的全部邏輯:步數檢查、prepare_next_tasks、空任務判 done、drain、_reapply_writes_to_succeeded_nodesshould_interrupt
  • out_of_steps:607-609self.step > self.stop 時進入 out_of_steps 狀態並回傳 False,外層 while 退出。
  • done branch:653-655self.tasks 為空時進 done 狀態,這是圖的「自然死亡」。
  • draining branch:657-659control.drain_requested 為真時進 draining 狀態。
  • reapply + resume handlers:662-664 — 續跑時先把 pending writes 貼回成功節點,再調 _resume_error_handlers_if_applicable 給失敗節點準備 handler task。
  • interrupt_before:667-671 — 調 should_interrupt 檢查節點是否落在 interrupt_before 列表裡,命中就 raise GraphInterrupt
  • after_tick method:683-726 — 收 writes → apply_writes → emit values → 清 pending → _put_checkpointinterrupt_after 檢查。
  • _reapply_writes_to_succeeded_nodes:736-749 — 把 pending writes 貼回記憶體 task,但跳過 ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME 這四種控制訊號。
  • _resume_error_handlers_if_applicable:751-816 — 掃 ERROR_SOURCE_NODE 標記,給失敗節點構造 error-handler task,讓 runner 跳過原任務直接跑 handler。
  • put_writes:415-508 — task 產出 writes 的入口,去重控制通道、累積 null task、起 checkpointer 的 put_writes future。
  • SyncPregelLoop:1469 / AsyncPregelLoop:1722 — 兩個子類別,覆蓋 __enter__/__exit__accept_pushschedule_error_handler 等同步/非同步差異部分。

資料流

tick 的入口先做步數上限檢查,接著調 prepare_next_tasks 把「下一步要跑哪些 task」算出來。這步是整個超步的種子——後面所有動作都是對這批 task 的加工。

python
def tick(self) -> bool:
    """Execute a single iteration of the Pregel loop.

    Returns:
        True if more iterations are needed.
    """

    # check if iteration limit is reached
    if self.step > self.stop:
        self.status = "out_of_steps"
        return False

    # prepare next tasks
    self.tasks = prepare_next_tasks(
        self.checkpoint,
        self.checkpoint_pending_writes,
        self.nodes,
        self.channels,
        self.managed,
        self.config,
        self.step,
        self.stop,
        for_execution=True,
        manager=self.manager,
        store=self.store,
        checkpointer=self.checkpointer,
        trigger_to_nodes=self.trigger_to_nodes,
        updated_channels=self.updated_channels,
        retry_policy=self.retry_policy,
        cache_policy=self.cache_policy,
    )

這段來自 tick 頭部:599-629self.step > self.stop 是「步數上限」退出條件——stopPregelLoop.__init__ 時被設成 recursion_limit,遞迴深度超了就停在 out_of_stepsprepare_next_tasks 拿當前 checkpoint、pending_writes、節點表、通道表算出本步要執行哪些 task,塞進 self.tasksfor_execution=True 表示這些 task 是要真跑的(而非 dry-run 用於串流預覽)。

緊接著是三個並列的退出條件——空 task / drain / 續跑貼回 writes——再走 should_interrupt:

python
# if no more tasks, we're done
if not self.tasks:
    self.status = "done"
    return False

if self.control is not None and self.control.drain_requested:
    self.status = "draining"
    return False

# if there are pending writes from a previous loop, apply them
if not self.is_replaying and self.checkpoint_pending_writes:
    self._reapply_writes_to_succeeded_nodes(self.tasks)
    self._resume_error_handlers_if_applicable()

# before execution, check if we should interrupt
if self.interrupt_before and should_interrupt(
    self.checkpoint, self.interrupt_before, self.tasks.values()
):
    self.status = "interrupt_before"
    raise GraphInterrupt()

這段來自 tick 中段:652-671interrupt_before 的檢查只在有 pending_writes 續跑之後才做——因為續跑時通道版本號已經反映了上次的寫入,should_interrupt 才能正確判斷「自上次 interrupt 以來有沒有新更新」。

after_tick 在 runner 把 task 全部跑完後呼叫,負責收尾:

python
def after_tick(self) -> None:
    # finish superstep
    writes = [w for t in self.tasks.values() for w in t.writes]
    self._delta_channels_with_overwrite.update(
        ch
        for ch, v in writes
        if isinstance(self.specs.get(ch), DeltaChannel) and _get_overwrite(v)[0]
    )
    # all tasks have finished
    self.updated_channels = apply_writes(
        self.checkpoint,
        self.channels,
        self.tasks.values(),
        self.checkpointer_get_next_version,
        self.trigger_to_nodes,
    )

這段來自 after_tick 頭部:683-698。先把所有 task 的 writes 攤平,順手統計哪些 delta 通道發生了 overwrite(影響 sparse replay 的起始值);再調 apply_writes 把這些 writes 真正合進通道,回傳值 updated_channels 會被下一步的 prepare_next_tasks 用來加速「找下一個該觸發的節點」。after_tick 的尾部還會清空 checkpoint_pending_writes、把 is_replayingFalse_put_checkpoint({"source": "loop"}) 落檢查點,然後跑 interrupt_after 判斷(after_tick 尾部:714-724)。

整個超步的生命週期可以畫成這樣:

邊界與失敗

  • out_of_steps 不是錯誤:step > stop 直接回傳 False,status 被設成 out_of_steps,外層 while 自然退出。Pregel 上層會根據 status 判斷是 throw RecursionError 還是靜默收尾——這是「軟上限」語意,不是異常路徑(out_of_steps:607-609)。
  • draining 不能 emit lifecycle events:_emit_graph_lifecycle_event 顯式拒絕在 draining 狀態下被呼叫(draining guard:384-385),因為 drain 是使用者主動停,不算 interrupt / resume 語意。
  • 續跑時 pending writes 含控制訊號要跳過:_reapply_writes_to_succeeded_nodes 必須跳過 ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME(skip control signals:746-747),否則失敗過的 task 會被當成成功的,跑過的 handler 又被重新觸發。
  • Resume 時的 error handler 重新排程:_resume_error_handlers_if_applicable 只對帶 error_handler_node 的節點生效(handler_node check:791-793),沒有 handler 的失敗節點保持原樣,由 runner 的 _should_stop_others 走 panic 路徑。
  • put_writes 的 null task 累積語意:NULL_TASK_ID 的 writes 不覆蓋而是累積(null task accumulate:422-431),用於 input writes 這種「不屬於任何節點但要進 checkpoint」的寫入。
  • Delta channel 的 future 必須先於下個 checkpoint 落盤:_delta_write_futs 收集了所有 delta 通道的 put_writes future,下個 checkpoint 落盤前會先 drain 它(_delta_write_futs:201-207)——顛倒順序會讓 sparse replay 看不到產生 checkpoint 的那些寫入。

小結

PregelLoop 是一個狀態機驅動的超步節拍器,把「準備 task → 跑 → 收尾」三段夾住 PregelRunner,自己不碰執行。理解了 tick / after_tick 兩段、status 七種狀態、以及 is_replaying + _reapply_writes_to_succeeded_nodes 的續跑機制,就抓住了 LangGraph 執行模型的主軸。Runner 怎麼把 task 真正跑起來見 /pregel/runner;演算法層 prepare_next_tasks / apply_writes / should_interrupt 的內部細節見 /pregel/algo;整體怎麼組裝見 /pregel/pregel。通道的寫回語意見 /channel/base-channel

對照官方資料:LangGraph 文件 · README