Skip to content

Topic 與 BinaryOperatorAggregate:累積式通道與歸約器

源码版本1.2.9

職責

LastValue 只會覆蓋,但實際 agent 狀態裡很多欄位是要累積的:訊息列表要追加、計數器要相加、待辦要合併。LangGraph 給這類需求提供兩條路:Topic(Topic class:23)是一個 pub/sub 風格的列表通道,BinaryOperatorAggregate(BinaryOperatorAggregate class:65)是一個接受任意二元運算子 (reducer) 的通用通道。兩者都直接或間接繼承自 BaseChannel(BaseChannel:19)。

BinaryOperatorAggregate 是使用者寫的 Annotated[list, add_messages] 這種宣告背後的實作——StateGraph 編譯時呼叫 _is_field_binop(_is_field_binop:1890-1908),把 Annotated 元資料裡最後那個雙參可呼叫物件當成 reducer,實例化成 BinaryOperatorAggregate(typ, reducer)。每次 update 都把新值用 operator 合併進當前值(update:123-144)。

Topic 更多用在內部通道:每個 Pregel 實例都自帶一個名為 __pregel_tasksTopic(Send, accumulate=False)(TASKS channel:809)用來裝 Send 扇出。它的特點是 accumulate=False 時每步清空——上一步的 Send 列表被 prepare_next_tasks 消費完就丟掉,下一步重新收集。

兩者共同點是「update 接受多個值並按某種規則 fold」,這是它們和 LastValue 的根本區別——後者只接受一個值。

設計動機

  • reducer 是狀態合併的通用介面:不同欄位的合併語意千差萬別——messages 要按 ID 去重追加、counter 要相加、tags 要聯集。與其給每種都做一個 channel 類別,不如提供一個通用通道接二元運算子 operator: (Value, Value) -> Value,讓使用者自己定義合併規則。add_messages(add_messages:60-66)就是這種 reducer:按訊息 ID 合併,新 ID 追加、舊 ID 替換。
  • 型別擦除後還原:_strip_extras(_strip_extras:22-28)把 Annotated / Required / NotRequired 這些 typing 包裝剝掉,再用 typ() 實例化預設值。collections.abc.Sequence 這種抽象型別不能直接實例化,所以在 typ fallback:83-88 裡映射回 list / set / dict(typ instantiation:82-92)。
  • Topic.accumulate 雙形態:accumulate=True 時是一個跨步累積的列表——update 把新值 extend 進去(update:77-85);accumulate=False 時每步先把舊值清空再 extend,變成「只看本步產生的值」。後者正是 TASKS 通道需要的:Send 一步內有效,被消費完就作廢。
  • _flatten 支援混合寫:Topic.update 接受 Value | list[Value] 混合序列(_flatten:15-20),所以節點可以一次寫一條 Send 也可以一次寫 [Send, Send],通道統一拍平。
  • Overwrite 旁路:reducer 通道還支援一種「直接覆蓋」語意:Overwrite 資料類別(Overwrite:980)或字典 {"__overwrite__": value}(OVERWRITE key:95)被 _get_overwrite 識別(_get_overwrite:31-51),在 update 裡繞過 reducer 直接賦值。每步只允許出現一次 Overwrite,第二次會拋 InvalidUpdateError(overwrite dedup:132-138)。

關鍵檔案

  • Topic class:23-41 — 類別定義,泛型參數 Sequence[Value] / Value | list[Value] / list[Value],建構函式接 accumulate 開關。
  • _flatten:15-20 — 工具函式:把混合 Value | list[Value] 序列拍平成 Iterator[Value]
  • Topic.copy:56-61 — 覆寫 copy 直接復用 values 列表的副本,不走 checkpoint 序列化。
  • Topic.checkpoint / from_checkpoint:63-75 — checkpoint 直接回傳 self.values 列表;from_checkpoint 還相容舊版 tuple 形式。
  • Topic.update:77-85 — 核心邏輯:accumulate=False 先清空再 extend,accumulate=True 直接 extend;空 flat 序列回傳 False
  • Topic.get / is_available:87-94 — 空列表拋 EmptyChannelError;is_available 覆寫成 bool(self.values)
  • Binop class:65-92 — 類別定義 + __init__:接 operator 二元可呼叫物件,用 typ() 初始化存值。
  • _get_overwrite:31-51 — 識別三種 Overwrite 表示形式(資料類別、sentinel-keyed dict、JSON 還原後的 dict)。
  • _operators_equal:54-62 — lambda 比較特殊處理:同名 lambda 視為相等,避免每次重編譯圖都把通道判為變更。
  • Binop.update:123-144 — reducer 的真正入口:首個值用來初始化,後續值逐個呼叫 operator(self.value, value),Overwrite 旁路例外。
  • _is_field_binop:1890-1908 — 編譯期識別 Annotated[..., reducer] 元資料裡的雙參可呼叫物件,實例化成 BinaryOperatorAggregate
  • TASKS channel:805-809Pregel.__init____pregel_tasks 通道硬編碼 Topic(Send, accumulate=False),不允許使用者覆蓋。
  • add_messages:60-66 — 典型 reducer:按訊息 ID 合併兩個列表,新 ID 追加、舊 ID 替換。

資料流

BinaryOperatorAggregate.update 是 reducer 通道的心臟(update:123-144):

python
def update(self, values: Sequence[Value]) -> bool:
    if not values:
        return False
    if self.value is MISSING:
        self.value = values[0]
        values = values[1:]
    seen_overwrite: bool = False
    for value in values:
        is_overwrite, overwrite_value = _get_overwrite(value)
        if is_overwrite:
            if seen_overwrite:
                msg = create_error_message(
                    message="Can receive only one Overwrite value per super-step.",
                    error_code=ErrorCode.INVALID_CONCURRENT_GRAPH_UPDATE,
                )
                raise InvalidUpdateError(msg)
            self.value = overwrite_value
            seen_overwrite = True
            continue
        if not seen_overwrite:
            self.value = self.operator(self.value, value)
    return True

讀這段要注意兩個分支:第一個值若通道是空的,直接當初始值(跳過一次 reducer 呼叫,因為 reducer 需要兩個輸入);之後每個值若不是 Overwriteself.value = operator(self.value, value)。如果出現 Overwrite,之後所有非 overwrite 值都被丟棄——這保證「覆蓋」語意不會被後續 reducer 又改回去。

StateGraph 怎麼把 Annotated[list[AnyMessage], add_messages] 編譯成 BinaryOperatorAggregate,看 _is_field_binop:1890-1908:

python
def _is_field_binop(typ: type[Any]) -> BinaryOperatorAggregate | None:
    if hasattr(typ, "__metadata__"):
        meta = typ.__metadata__
        if len(meta) >= 1 and callable(meta[-1]):
            sig = signature(meta[-1])
            params = list(sig.parameters.values())
            if (
                sum(
                    p.kind in (p.POSITIONAL_ONLY, p.POSITIONAL_OR_KEYWORD)
                    for p in params
                )
                == 2
            ):
                return BinaryOperatorAggregate(typ, meta[-1])
            else:
                raise ValueError(
                    f"Invalid reducer signature. Expected (a, b) -> c. Got {sig}"
                )
    return None

它用 inspect.signature 校驗最後一個 meta 必須是雙參可呼叫物件,簽名不匹配直接拋錯。這就是為什麼 add_messages(left, right) 這種兩參函式能當 reducer,而單參函式會被拒。

Topic.update 走的是另一條路徑,邏輯更短(update:77-85):

python
def update(self, values: Sequence[Value | list[Value]]) -> bool:
    updated = False
    if not self.accumulate:
        updated = bool(self.values)
        self.values = list[Value]()
    if flat_values := tuple(_flatten(values)):
        updated = True
        self.values.extend(flat_values)
    return updated

注意 accumulate=False 時的回傳值:如果舊值非空,即使新值序列空也回傳 True(因為清空本身就是一種改動),Pregel 會更新 channel_versions。這和 LastValue 空序列回傳 False 的行為相反,反映了不同通道對「空寫」的不同語意。

reducer 通道在一次 Pregel 超步裡的位置:

邊界與失敗

  • Overwrite 三種識別形式並存:overwrite forms:44-50 同時接受 Overwrite 資料類別、{"__overwrite__": value} 字典、以及 JSON 還原後的 {"type": "__overwrite__", "value": ...} 字典。第三個分支是為了讓經過 LangGraph API server 的 orjson 序列化往返後語意不丟。
  • 每步只能有一個 Overwrite:overwrite dedup:132-138 第二次出現直接拋 InvalidUpdateError,因為「覆蓋」語意在一步內多次出現是矛盾的。
  • Overwrite 之後的非 overwrite 值被靜默丟棄:discard after overwrite:142-143 if not seen_overwrite 決定是否呼叫 reducer,一旦 seen_overwrite=True,後續值既不入 reducer 也不報錯。這是有意的——「覆蓋」意味著不想保留之前的合併結果,也不需要後續值。
  • lambda reducer 等價性比較:_operators_equal:54-62 把所有 lambda 視為相等,因為 Python 的 lambda 每次 <lambda> 都是新物件,按身份比會導致 __eq__ 永遠不相等,影響圖快取。
  • Topic(accumulate=False) 空寫回傳 True:clear returns True:79-80——只要舊值非空,清空就視為改動。這會讓 apply_writes 更新 channel_versions 從而觸發依賴該通道的節點。如果使用者自己用 Topic,要意識到這一點:它不像 LastValue 那樣靜默處理空寫。
  • Topic 舊 checkpoint 相容:tuple compatibility:70-74 檢查 isinstance(checkpoint, tuple),舊版本的 checkpoint 存的是 (typ, values) 元組,新版改成直接存 list。這段是為了不破壞舊存檔。
  • reducer 簽名強制雙參:reducer signature check:1894-1907inspect.signature 數 positional 參數,不是兩個就報錯。帶 *args 這種可變參數的函式不會被接受——這是把「reducer 必須嚴格 (a, b) -> c」釘死。
  • add_messages 是個工廠:_add_messages_wrapper:41-57@_add_messages_wrapper 裝飾,實際函式可以在 add_messages(left, right) 直接呼叫,也可以 add_messages(format="langchain-openai") 回傳 partial。這種用法是 reducer 模式的常見擴展點——給 reducer 加執行時參數。

小結

BinaryOperatorAggregate 是「帶 reducer 的欄位」的統一實作,Topic 是「列表累積」的兩種形態:跨步累積或每步清空。它們和 /channels/last-value 構成了 LangGraph 通道系統的三件套,介面都來自 /channels/base-channeladd_messages 這種 reducer 在 /graph/state-graph 裡被 Annotated[..., reducer] 宣告,執行時由 /pregel/algoapply_writes 呼進 updateTopic(accumulate=False) 作為 __pregel_tasks 通道是 Send 扇出的載體,具體見 /subgraph/send-command

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