PregelLoop:超步循環驅動器
職責
PregelLoop 是 LangGraph 執行時的「節拍器」。Pregel 本身只負責組裝資源、暴露 invoke / astream 入口,真正一圈圈跑圖的是一個 PregelLoop 實例:它持有當前檢查點 (checkpoint)、通道 (channel) 集合、pending_writes、step 計數器、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 的活;它只決定「這一步該看誰、跑完怎麼收尾、下一步還要不要繼續」。
設計動機
為什麼把循環邏輯拆成 tick 和 after_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_writesfuture,_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_nodes、should_interrupt。out_of_steps:607-609—self.step > self.stop時進入out_of_steps狀態並回傳False,外層while退出。done branch:653-655—self.tasks為空時進done狀態,這是圖的「自然死亡」。draining branch:657-659—control.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列表裡,命中就 raiseGraphInterrupt。after_tick method:683-726— 收 writes →apply_writes→ emit values → 清 pending →_put_checkpoint→interrupt_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_writesfuture。SyncPregelLoop:1469/AsyncPregelLoop:1722— 兩個子類別,覆蓋__enter__/__exit__、accept_push、schedule_error_handler等同步/非同步差異部分。
資料流
tick 的入口先做步數上限檢查,接著調 prepare_next_tasks 把「下一步要跑哪些 task」算出來。這步是整個超步的種子——後面所有動作都是對這批 task 的加工。
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-629。self.step > self.stop 是「步數上限」退出條件——stop 在 PregelLoop.__init__ 時被設成 recursion_limit,遞迴深度超了就停在 out_of_steps。prepare_next_tasks 拿當前 checkpoint、pending_writes、節點表、通道表算出本步要執行哪些 task,塞進 self.tasks。for_execution=True 表示這些 task 是要真跑的(而非 dry-run 用於串流預覽)。
緊接著是三個並列的退出條件——空 task / drain / 續跑貼回 writes——再走 should_interrupt:
# 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-671。interrupt_before 的檢查只在有 pending_writes 續跑之後才做——因為續跑時通道版本號已經反映了上次的寫入,should_interrupt 才能正確判斷「自上次 interrupt 以來有沒有新更新」。
after_tick 在 runner 把 task 全部跑完後呼叫,負責收尾:
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_replaying 置 False、_put_checkpoint({"source": "loop"}) 落檢查點,然後跑 interrupt_after 判斷(after_tick 尾部:714-724)。
整個超步的生命週期可以畫成這樣:
邊界與失敗
out_of_steps不是錯誤:step > stop直接回傳False,status 被設成out_of_steps,外層while自然退出。Pregel上層會根據 status 判斷是 throwRecursionError還是靜默收尾——這是「軟上限」語意,不是異常路徑(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_writesfuture,下個 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