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 の仕事です。PregelLoop は「このステップで誰を見るか、走り終わったらどう収尾するか、次のステップを続けるか」だけを決めます。

設計動機

なぜループロジックを tickafter_tick の二段に分け、一つの step() メソッドで済ませないのか?

  • 実行権を呼び出し元に返す:tickTrue を返した後、呼び出し側(Pregel.astream の実装体、sync main loop:2964-2984 参照)が自ら runner.tick を駆動してこのバッチの task を走らせ、その後 loop.after_tick を呼びます。これにより PregelLoopPregelRunner への参照を持たず、両オブジェクトは疎結合な協力関係にあって入れ子になりません。
  • ステートマシンの明示化: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 を自然に終了させ、ノード途中で強制 kill しません。
  • Exit-mode delta 書き込みとチェックポイント時序:_delta_write_futs(_delta_write_futs:207)はすべての delta チャネルの put_writes future を集め、_checkpointer_put_after_previous は次の checkpoint を書く前にこれらを必ず drain します——「書き込みが 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-671should_interrupt でノードが interrupt_before リストに入るかチェックし、命中すれば GraphInterrupt を raise。
  • after_tick method:683-726 — writes 収集 → apply_writes → values emit → 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-816ERROR_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-629 からの引用です。self.step > self.stop は「ステップ数上限」の離脱条件——stopPregelLoop.__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 に入ります:

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-671 からの引用です。interrupt_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 に基づき RecursionError を throw するか黙って収尾するかを判断します——これは「ソフト上限」セマンティクスで、例外パスではありません(out_of_steps:607-609)。
  • draining では lifecycle events を emit できない:_emit_graph_lifecycle_eventdraining 状態での呼び出しを明示的に拒否します(draining guard:384-385)。drain はユーザ能動の停止で、interrupt / resume セマンティクスではないからです。
  • 再開時の pending writes に制御信号が含まれる場合はスキップ:_reapply_writes_to_succeeded_nodesERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME を必ずスキップします(skip control signals:746-747)。さもないと、失敗した task が成功扱いされ、走った handler が再び触発されます。
  • Resume 時の error handler 再スケジュール:_resume_error_handlers_if_applicableerror_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