- 发布日期
第 05 讲 · Channel 家族详解(下):屏障与特殊通道
NamedBarrierValue / AfterFinish / AnyValue / UntrackedValue 屏障与特殊通道机制
学习目标
- 理解
NamedBarrierValue如何实现"等所有上游到齐再放行"的同步屏障。 - 理解
...AfterFinish系列为什么要"延迟到 finish 才暴露"。 - 了解
AnyValue/UntrackedValue两个特殊通道的用途。
1. NamedBarrierValue:扇入同步屏障
channels/named_barrier_value.py。这是 LangGraph 实现"join / 汇合"的核心通道。
它持有一组期望的名字 names,和已到达的名字 seen。只有当 seen == names(全到齐) 时通道才"可用"。
update:记录谁到了
def update(self, values: Sequence[Value]) -> bool:
updated = False
for value in values:
if value in self.names:
if value not in self.seen:
self.seen.add(value)
updated = True
else:
raise InvalidUpdateError(
f"At key '{self.key}': Value {value} not in {self.names}"
)
return updated
- 每个上游写入自己的"名字",通道把它加进
seen。 - 写了不在
names里的值 → 报错(防止配置错误)。
get / is_available:未到齐就是"空"
def get(self) -> Value:
if self.seen != self.names:
raise EmptyChannelError()
return None
def is_available(self) -> bool:
return self.seen == self.names
- 没到齐时
get()抛EmptyChannelError→ 下游节点不会被触发。 - 到齐后
is_available()为True→ 触发下游。 - 注意它
get()返回None:屏障只负责"同步信号",不传数据,数据走别的通道。
consume:放行后清空,准备下一轮
def consume(self) -> bool:
if self.seen == self.names:
self.seen = set()
return True
return False
- 下游被触发执行后,
consume()把seen清空,使屏障可重复使用(循环图里很关键)。
它从哪来:多个普通边汇入一个节点
当你在 StateGraph 里给一个节点添加多条普通边(多个上游都指向它), 编译时会为它生成一个 NamedBarrierValue,names = 所有上游节点名。 于是该节点等所有上游都跑完才执行——这就是 join 语义。
对比:如果你希望"任一上游到了就触发",那是另一种语义(见下文 AnyValue / Topic)。
2. AfterFinish 系列:延迟到收尾才暴露
LastValueAfterFinish 与 NamedBarrierValueAfterFinish。它们比普通版多一个 finished 标志, get() 只有在 finished=True 时才返回值。
def finish(self) -> bool:
if not self.finished and self.value is not MISSING:
self.finished = True
return True
else:
return False
def get(self) -> Value:
if self.value is MISSING or not self.finished:
raise EmptyChannelError()
return self.value
为什么需要它?考虑这样的需求:某个值在 run 真正结束前不应该触发任何下游。 普通 LastValue 一写就可读、就会触发;而 AfterFinish 把"可读"推迟到 finish() (run 收尾,见第 3 讲方法表)被调用时,避免提前触发。
consume()在被读取后清空值(一次性)。- 典型用于"最终输出聚合"、"只在收尾阶段汇总"的场景。
3. AnyValue:取任意一个值(不报错)
channels/any_value.py。和 LastValue 相反:一步收到多个值也不报错,随便存一个 (语义上"它们应该都一样,或谁赢无所谓")。
适用:多个上游会写同一个幂等/相同的值,你只要其中一个,不想为此引入 reducer。
4. UntrackedValue:存值但不参与版本/检查点
channels/untracked_value.py。它存值供读取,但不进检查点、不驱动版本触发。
适用:放一些运行期临时、不需要持久化、也不该触发下游重算的辅助数据。 (与 EphemeralValue 区别:Ephemeral 是"只活一步",Untracked 是"不被追踪/不持久化"。)
特殊通道速查表
| 通道 | 核心语义 | 触发下游条件 | 典型场景 |
|---|---|---|---|
NamedBarrierValue | 等 names 全到齐 | seen==names | 多上游 join |
NamedBarrierValueAfterFinish | 到齐 + finish | finish 后 | 收尾阶段 join |
LastValueAfterFinish | 末值 + finish | finish 后 | 延迟暴露最终值 |
AnyValue | 任取一个,不报错 | 有值即可 | 幂等多源 |
UntrackedValue | 存值不追踪 | 不触发 | 临时辅助数据 |
使用场景:Map-Reduce 的 reduce 屏障
经典 fan-out/fan-in:把任务拆成 N 份并行(map),全部完成后汇总(reduce)。
- map 阶段:用
Send动态扇出 N 个子任务(第 19 讲)。 - reduce 节点:靠 reducer 通道(如
Annotated[list, operator.add])收集 N 份结果。 - 若 reduce 节点依赖"多个固定上游都完成",则用
NamedBarrierValue做同步屏障, 确保不会"少数上游先到就提前 reduce"。
理解了屏障,你就能解释"为什么我的汇总节点早跑了/没跑"这类问题。
动手实验
- 画一个菱形图:
A → B、A → C、B → D、C → D。给D加打印。 观察D是否等 B 和 C 都完成才执行一次(NamedBarrierValue 的 join 效果)。 - 把其中一条边删掉,观察
D触发时机变化。
阅读作业
- 精读
channels/named_barrier_value.py全文,重点对比NamedBarrierValue与NamedBarrierValueAfterFinish的get/consume/finish差异。 - 浏览
channels/any_value.py与channels/untracked_value.py,各用一句话总结其语义。
小结
NamedBarrierValue= 同步屏障,等所有上游到齐才放行,是 join 的实现。AfterFinish系列把"可读"推迟到 run 收尾,避免提前触发。AnyValue/UntrackedValue是两个特殊用途通道:容忍多值 / 不追踪。
下一讲:DeltaChannel——用增量 replay 省存储的进阶通道。