- 发布日期
第 11 讲 · 超步循环(下):Update 阶段
学习目标
- 看懂
after_tick()→apply_writes()如何把缓冲的写合并进通道。 - 理解版本 bump、consume、finish 的精确时机。
- 彻底理解第 1 讲的黄金法则"步内不可变、步间才可见"。
after_tick():Update 阶段主函数
_loop.py,after_tick()(676–714 行)。它在一个超步的所有节点执行完后被调用,做四件事:
- 收集所有任务的写:
task.writes(执行阶段缓冲在这里,第 13 讲)。 apply_writes(checkpoint, channels, tasks, ...)→ 真正改通道,返回updated_channels。- 若输出通道有更新 → emit
values流。 - 清空 pending writes,调
_put_checkpoint({"source": "loop"})存档(并step += 1)。 - 检查
interrupt_after(第 18 讲)。
updated_channels(本步真正变化的通道集合)会传给下一步的 prepare_next_tasks, 驱动 trigger_to_nodes 反向索引优化(第 10 讲)。
apply_writes:Update 的核心算法
_algo.py,apply_writes(232–345 行)。这是"批量合并写"的地方,对照第 3 讲通道方法表:
flowchart TD
A[apply_writes] --> B[按 task.path 排序<br/>保证确定性]
B --> C[更新 versions_seen<br/>记录每个节点见过的版本]
C --> D[consume: 被触发的通道<br/>用过即焚]
D --> E[按通道分组 writes]
E --> F["对每个通道 channel.update(vals)"]
F --> G[bump channel_versions<br/>版本号 +1]
G --> H[对未写通道发空 update<br/>通知新超步]
H --> I{可能是最后一步?}
I -->|是| J["channel.finish()"]
I --> K[返回 updated_channels]
逐步拆解:
1. 排序保证确定性
按 task.path 排序后再 apply,保证同样的输入产生同样的合并结果(对 replay/时间旅行至关重要)。
2. 更新 versions_seen
把每个执行过的节点的 versions_seen 更新为它执行时见到的通道版本—— 这就是"防止同一份更新反复触发同一节点"的机制(第 10 讲的触发判定依赖它)。
3. consume:被触发的通道用过即焚
对被本步任务触发的通道调 consume()(第 3 讲)—— 比如 EphemeralValue 触发通道、NamedBarrierValue 屏障,被消费后清空,准备下一轮。
4. update + bump version
把每个通道本步收到的所有写作为一个序列调 channel.update(vals)(第 3/4 讲的 reducer 落点)。 若 update 返回 True(通道变了),则 bump 它的 channel_versions——版本涨了,下一步可能触发订阅者。
5. 空 update + finish
- 对本步没收到写的通道也调一次
update(空序列)(第 3 讲:"每步都调,哪怕空"), 让需要"感知新超步"的通道有机会更新内部状态。 - 若判断可能是最后一步,调
channel.finish(),让...AfterFinish通道暴露暂存值(第 5 讲)。
黄金法则的代码证明
现在我们能精确解释第 1 讲的"步内不可变、步间才可见":
- 执行阶段(Execute):节点的写只进
task.writes缓冲(第 13 讲),通道本体一点没动。 所以同一超步内并行的节点,谁也读不到对方的写——步内不可变。 - Update 阶段(after_tick):所有缓冲的写在这里一次性合并进通道、bump 版本。
- 下一超步(Plan):因为版本涨了,订阅者才被触发、才能
get()到新值——步间才可见。
sequenceDiagram
participant N1 as 节点A(本步)
participant N2 as 节点B(本步)
participant Buf as task.writes 缓冲
participant Ch as 通道(本体)
N1->>Buf: 写 x=1 (不动通道)
N2->>Buf: 写 y=2 (不动通道)
Note over Ch: 本步通道不可变, A/B 互不可见
Buf->>Ch: after_tick: apply_writes 批量合并
Note over Ch: bump 版本, 下一步才可见
使用场景
- 理解"为什么我在节点里改了 state 立刻读还是旧值":因为写在缓冲里,本步通道没变。 要在同一步内读到自己刚写的值,需要
local_read(fresh=True)(第 13 讲的特例)。 - 理解 durability 模式:
_put_checkpoint的时机(每步 vs 退出时)决定持久化开销, 第 12/17 讲会讲 sync/async/exit 三模式。 - 排查"写丢了/被覆盖":检查该通道的 reducer——
LastValue会让后写覆盖先写, 聚合字段才会保留全部。
动手实验
- 在
apply_writes打断点,打印每个通道update前后的值和版本号,观察 bump。 - 写一个节点
return {"x": 1}后立刻print(state["x"]),确认读到的是旧值(缓冲未合并)。 - 用
stream_mode="values"跑,确认每个 chunk 对应一次after_tick后的完整快照。
阅读作业
- 精读
_algo.py的apply_writes(232–345 行),逐行对照本讲流程图。 - 精读
_loop.py的after_tick(676–714 行)与_put_checkpoint(1064–1199 行,第 15 讲再深入)。
小结
after_tick= Update:apply_writes批量合并写、更新 versions_seen、consume、bump 版本、finish、存 checkpoint。- 黄金法则的真相:写先进
task.writes缓冲,Update 阶段才落通道并 bump 版本,下一步才可见。 updated_channels驱动下一步 Plan 的反向索引优化。
下一讲:PregelRunner 如何并行执行这些任务。