Skip to content

BaseChannel 抽象:チャネル即状態プリミティブ

源码版本1.2.9

役割

LangGraph の実行時において、「チャネル (channel)」は状態 (state) の最小単位です。StateGraph をコンパイルして得られる Pregel インスタンスは、内部的にユーザーの dict を直接保持するわけではなく、一連の BaseChannel インスタンスを保持します。各チャネルは1つのキーを持ち、値の格納・書き込みのマージ・読み出しの公開・現在状態のチェックポイント (checkpoint) 形式へのシリアライズを自分で担います。BaseChannel はすべてのチャネルクラスが共有する抽象基底クラスです(BaseChannel:19)。

このクラス自体はデータを持ちません(__slots__keytyp だけ、__slots__:22-26)。サブクラスが実装すべき4つの抽象メソッドだけを定めます。型属性 ValueType / UpdateType、および3つのアクション update / get / from_checkpoint(abstract methods:60-99)です。さらに checkpoint / is_available / consume / finish / copy の5つはデフォルト実装を持ち、どちらも get() にフォールバックするか False を返して no-op になります。

言い換えると、BaseChannel は「状態をどう格納し、どう変え、どう痕跡を残すか」の3つを一度にインターフェースとして釘付けにします。Pregel エンジンは update / get / checkpoint などのメソッドを呼ぶだけで、背後で LastValue がマージしているのか、Topic なのか、BinaryOperatorAggregate なのかを知る必要がありません。だからこそエンジンと状態形状は完全に疎結合にできます。

設計動機

なぜ状態を「チャネル」というオブジェクトにするのか。ノード間で dict をそのまま渡すのではなく?

  • 書き込みマージのセマンティクスを統一:同じスーパーステップ (superstep) 内で複数のノードが同じキーに並行書き込みする可能性があります。dict では「リストは append、スカラーは上書き、カウンタは加算」という異なるマージルールを表現できません。チャネルは update(values: Sequence[Update]) -> bool で戦略を一箇所に集約し、サブクラス自身が fold の方法を決めます。
  • ステップ内の不変なビュー:ノードが読むのはステップ開始時のスナップショットで、同じステップで他のノードが書いた値は次のステップまで見えません。get は現在値を読むだけ、update はステップ境界で apply_writes だけが呼びます(apply_writes:317-323)。これが競合を構造的に排除します。
  • シリアライズ可能:各チャネルは checkpoint() でシリアライズ可能な値を(checkpoint:49-58)、from_checkpoint で自分自身を再構築できます(from_checkpoint:60-65)。BaseCheckpointSaver 系はこの2つのインターフェース頼りで、スレッド全体の状態を SQLite/Postgres に書き出します。
  • ライフサイクルフック:consumefinish の2メソッド(consume / finish:101-121)は、特殊チャネル(Topic(accumulate=False)LastValueAfterFinish など)にステップ境界や実行終端で自身の状態を変える機会を与えます。デフォルトが no-op であることは、大多数のチャネルがこの層のセマンティクスを必要としないことを示します。
  • 型の自己記述:ValueType / UpdateType の2プロパティ(ValueType / UpdateType:28-36)により、コンパイル時にも実行時にも「このチャネルは何を格納し、何の書き込みを受け付けるか」を逆参照できます。StateGraph_is_field_channel:1862-1887 でこれを使い、Annotated[..., SomeChannel] 宣言を識別します。

主要ファイル

  • BaseChannel class:19-26 — 抽象基底クラス定義。ジェネリックパラメータ Value / Update / Checkpoint__slots__keytyp のみ。
  • ValueType / UpdateType:28-36 — 2つの抽象プロパティ。サブクラスは格納値型と書き込み値型を宣言するために使います。
  • checkpoint:49-58 — デフォルト実装は self.get() をそのまま返し、空チャネルは MISSING を返します。
  • from_checkpoint:60-65 — 抽象メソッド。1つの checkpoint 値から等価な新チャネルを構築し、サブクラスが実装必須。
  • get:69-73 — 抽象読み出しインターフェース。空チャネルは EmptyChannelError を投げます。
  • is_available:75-85 — デフォルト実装は get() を呼び EmptyChannelError を捕捉。サブクラスは通常より高速な判定に再定義します。
  • update:89-99 — 抽象書き込みインターフェース。Pregel が各スーパーステップの終わりに1回呼び、順序は任意。True はチャネルが変更されたことを意味します。
  • consume / finish:101-121 — ライフサイクルフック。デフォルトは no-op で False を返します。
  • apply_writes:315-345 — Pregel が当ステップの全タスクの書き出しをチャネルにマージする唯一の入口。update / consume / finish はすべてここで一括スケジュールされます。
  • _is_field_channel:1862-1887StateGraph はコンパイル時に isinstance(item, BaseChannel)Annotated[..., channel] を channel インスタンスに変換。

データフロー

チャネルへの書き込みは apply_writes の中で行われます。Pregel は当ステップの全タスクの writes をチャネル別に集約し、update を呼びます(apply_writes:317-323):

python
# Apply writes to channels
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():
    if chan in channels:
        if channels[chan].update(vals) and next_version is not None:
            checkpoint["channel_versions"][chan] = next_version
            # unavailable channels can't trigger tasks, so don't add them
            if channels[chan].is_available():
                updated_channels.add(chan)

# Channels that weren't updated in this step are notified of a new step
if bump_step:
    for chan in channels:
        if channels[chan].is_available() and chan not in updated_channels:
            if channels[chan].update(EMPTY_SEQ) and next_version is not None:
                checkpoint["channel_versions"][chan] = next_version
                if channels[chan].is_available():
                    updated_channels.add(chan)

2つ目の段落に注意してください。当ステップでどのタスクからも書かれなかったチャネルにも update(EMPTY_SEQ) が呼ばれます。これは BaseChannel インターフェースの暗黙の契約で、サブクラスの update は空シーケンスを受け付け、通常は「変更なし」を意味する False を返します。Topic(accumulate=False) のような「毎ステップクリア」セマンティクスはこの空呼び出しで実現されます。空 update で前ステップの値を消し、次ステップでは見えなくなります。

Pregel の1回の tick におけるチャネルのライフサイクル全体は次のとおりです:

境界と失敗

  • EmptyChannelError:get() はチャネルに一度も書き込まれていないときに投げます(get:70-73)。is_available のデフォルト実装はこの例外を catch して可用性を判定します。サブクラスがより安く判定できるなら再定義すべきです。
  • MISSING センチネル:checkpoint()get()EmptyChannelError を投げたとき None ではなく langgraph._internal._typing.MISSING を返します(checkpoint fallback:55-58)。これにより None は正当な格納値として「空」と区別して扱えます。
  • updateFalse を返す:False はチャネルが変更されていないことを意味し、Pregel は channel_versions を更新しません。後続の prepare_next_tasks もこのチャネルを起点にどのノードもトリガーしません(update call:319)。これが Pregel の収束判定の基盤です。
  • 空シーケンス呼び出し:apply_writes は未変更チャネルにも update(EMPTY_SEQ) を呼びます(EMPTY_SEQ update:329)。BaseChannel.update の契約は空入力で False を返すことですが、Topic(accumulate=False) はここで自身をクリアし True を返すことがあります。サブクラスは「空入力 = 何もしない」を前提にしてはいけません。
  • consumefinish はデフォルトで no-op(consume / finish:101-121):False を返すので、大多数のチャネルはライフサイクル管理に参加しません。apply_writesbump_step の終端で finish を呼び(finish call:338)、LastValueAfterFinish のような特殊チャネルだけが「実行終了まで見えない」セマンティクスにこれを使います。
  • copy のデフォルトは checkpoint 経由:BaseChannel.copy のデフォルトは self.from_checkpoint(self.checkpoint()) です(copy:40-47)。サブクラスが checkpoint を構築するコスト高い場合(大きな list を深コピーするなど)、copy を再定義してより安い実装を提供すべきです。Topic.copy がその例です。

まとめ

BaseChannel は「状態をどう格納し、どうマージし、どうシリアライズするか」を一度にインターフェースとして釘付けにします。Pregel エンジンは update → is_available → get → checkpoint → finish の固定順序で呼ぶだけで、任意のチャネル実装の上で動きます。具体的な実装は /channels/last-value/channels/topic-binop を参照してください。チャネルが Pregel から読み書きされる過程は /pregel/algo に、チェックポイントが checkpoint() の戻り値をどう永続化するかは /checkpoint/base-saver に続きます。

公式資料:LangGraph 文档 · README