发布日期

第 11 讲 · 超步循环(下):Update 阶段

学习目标

  • 看懂 after_tick()apply_writes() 如何把缓冲的写合并进通道。
  • 理解版本 bump、consume、finish 的精确时机。
  • 彻底理解第 1 讲的黄金法则"步内不可变、步间才可见"。

after_tick():Update 阶段主函数

_loop.pyafter_tick()(676–714 行)。它在一个超步的所有节点执行完后被调用,做四件事:

  1. 收集所有任务的写:task.writes(执行阶段缓冲在这里,第 13 讲)。
  2. apply_writes(checkpoint, channels, tasks, ...) → 真正改通道,返回 updated_channels
  3. 若输出通道有更新 → emit values 流。
  4. 清空 pending writes,调 _put_checkpoint({"source": "loop"}) 存档(并 step += 1)。
  5. 检查 interrupt_after(第 18 讲)。

updated_channels(本步真正变化的通道集合)会传给下一步的 prepare_next_tasks, 驱动 trigger_to_nodes 反向索引优化(第 10 讲)。

apply_writes:Update 的核心算法

_algo.pyapply_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 会让后写覆盖先写, 聚合字段才会保留全部。

动手实验

  1. apply_writes 打断点,打印每个通道 update 前后的值和版本号,观察 bump。
  2. 写一个节点 return {"x": 1} 后立刻 print(state["x"]),确认读到的是旧值(缓冲未合并)。
  3. stream_mode="values" 跑,确认每个 chunk 对应一次 after_tick 后的完整快照。

阅读作业

  • 精读 _algo.pyapply_writes(232–345 行),逐行对照本讲流程图。
  • 精读 _loop.pyafter_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 如何并行执行这些任务。