Topic 與 BinaryOperatorAggregate:累積式通道與歸約器
職責
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_tasks 的 Topic(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-809—Pregel.__init__給__pregel_tasks通道硬編碼Topic(Send, accumulate=False),不允許使用者覆蓋。add_messages:60-66— 典型 reducer:按訊息 ID 合併兩個列表,新 ID 追加、舊 ID 替換。
資料流
BinaryOperatorAggregate.update 是 reducer 通道的心臟(update:123-144):
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 需要兩個輸入);之後每個值若不是 Overwrite 就 self.value = operator(self.value, value)。如果出現 Overwrite,之後所有非 overwrite 值都被丟棄——這保證「覆蓋」語意不會被後續 reducer 又改回去。
StateGraph 怎麼把 Annotated[list[AnyMessage], add_messages] 編譯成 BinaryOperatorAggregate,看 _is_field_binop:1890-1908:
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):
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-143if 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-1907用inspect.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-channel。add_messages 這種 reducer 在 /graph/state-graph 裡被 Annotated[..., reducer] 宣告,執行時由 /pregel/algo 的 apply_writes 呼進 update。Topic(accumulate=False) 作為 __pregel_tasks 通道是 Send 扇出的載體,具體見 /subgraph/send-command。
對照官方資料:LangGraph 文件 · README