6. Watch API 实时监听机制

6.Watch API 实时监听机制

在分布式系统中,如何高效、实时地感知数据的变化,是一个核心挑战。传统的“轮询”(Polling)方式,即客户端定期向服务器查询数据是否有更新,不仅会给服务器带来巨大的压力,而且延迟较高,无法满足实时性的要求。etcd 的 Watch API 提供了一种基于事件流的异步监听机制,完美地解决了这个问题。

从本质上讲,etcd 的 Watch 是一种长连接。客户端与 etcd 服务器建立一个 gRPC 双向流(Bidirectional Stream),客户端通过这个流发送监听请求,服务器则通过同一个流将数据变更事件持续地推送回客户端。这种“订阅-发布”模型,相比轮询,具有极高的效率和极低的延迟。

核心概念:事件与修订号

理解 Watch 机制,首先要理解两个核心概念:事件(Event)和修订号(Revision)。

etcd 的数据模型是多版本并发控制(MVCC)的,这意味着每一次对键值对的修改,都不会覆盖旧数据,而是生成一个新的版本。每一次修改操作(如 Put、Delete)都会被分配一个全局单调递增的修订号。这个修订号是 etcd 的逻辑时钟,它定义了所有操作的全局顺序。

当一个键被修改时,etcd 会生成一个事件。根据 api.md 中的定义,一个事件包含以下关键信息:

message Event {
  enum EventType {
    PUT = 0;
    DELETE = 1;
  }
  EventType type = 1;         // 事件类型:PUT 或 DELETE
  KeyValue kv = 2;            // 事件发生后的键值对
  KeyValue prev_kv = 3;       // 事件发生前的键值对(可选)
}
  • type: 表示是创建/更新(PUT)还是删除(DELETE)操作。
  • kv: 包含了事件发生后该键的最新状态。对于 PUT 事件,如果 kv.version 为 1,说明这是一个新创建的键。对于 DELETE 事件,kv 包含了被删除的键及其删除时的修订号。
  • prev_kv: 这是一个非常有用的字段,它记录了事件发生前该键的状态。通过它,我们可以知道什么数据被覆盖了。不过,为了节省带宽,这个字段默认是不填充的,需要客户端在创建 Watch 时显式开启。

Watch 的工作流程

Watch 的工作流程可以概括为以下几步:

  1. 建立流:客户端与 etcd 服务器建立一个 gRPC 双向流。
  2. 发送创建请求:客户端通过流发送一个 WatchCreateRequest 消息,指定要监听的键(或键范围)、起始修订号等参数。
  3. 服务器确认:服务器收到请求后,会返回一个 WatchResponse,其中 created 字段为 true,并分配一个唯一的 watch_id。这个 ID 用于在同一个流中区分不同的 Watch 实例。
  4. 持续推送事件:从指定的修订号开始,一旦该键(或键范围)发生任何变更,服务器就会将对应的 Event 封装在 WatchResponse 中,通过流推送给客户端。
  5. 处理事件:客户端持续从流中读取 WatchResponse,解析其中的事件列表,并执行相应的业务逻辑。

这种机制的优势在于,客户端只需发起一次连接请求,后续的事件推送都由服务器主动完成,避免了反复的请求-响应开销,实现了真正的实时监听。

事件流处理模型

etcd 的 Watch 机制在设计上保证了事件流的可靠性和有序性,这对于构建正确的分布式应用至关重要。

可靠性与有序性保证

根据 api_guarantees.md 的描述,etcd 对 Watch 事件做出了以下关键保证:

  • 有序性(Ordered):所有事件都严格按照修订号的顺序发布。一个事件绝不会在比它更早的事件已经被发布之后才出现。
  • 唯一性(Unique):同一个事件绝不会在 Watch 中出现两次。
  • 可靠性(Reliable):事件序列绝不会丢失任何可用历史窗口内的子序列。如果事件 a, b, c 按时间顺序发生,Watch 收到了 a 和 c,那么只要 b 还在历史窗口内,它就保证会收到 b。
  • 原子性(Atomic):一个事件列表保证涵盖完整的修订号。同一个修订号内对多个键的更新,不会被拆分到多个事件列表中。
  • 可恢复性(Resumable):如果 Watch 连接中断,客户端可以通过在最后一个收到的事件修订号之后建立新的 Watch 来恢复监听,只要该修订号仍在历史窗口内。

这些保证意味着,客户端可以放心地按照接收顺序处理事件,并且在断线重连后能够准确地从上次中断的地方继续,不会错过任何数据变更。

多路复用(Multiplexing)

一个 gRPC 流可以承载多个 Watch。客户端可以在同一个流上发送多个 WatchCreateRequest 来创建不同的 Watch。服务器会为每个 Watch 分配唯一的 watch_id,并将所有事件通过这个共享的流发送回来,每个事件都附带其对应的 watch_id。

这种多路复用机制极大地减少了客户端与服务器之间的连接数和内存开销,尤其是在需要监听大量不同键的场景下,优势非常明显。

过滤器与历史监听

在创建 Watch 时,客户端可以通过 WatchCreateRequest 提供的字段进行精细控制:

  • start_revision:这是 Watch 最强大的功能之一。它允许客户端从一个历史修订号开始监听。这在什么场景下使用呢?假设你的应用因为网络问题或重启而与 etcd 断开连接,在此期间 etcd 可能已经发生了多次数据变更。重新连接后,你可以先读取当前最新的修订号,然后创建一个 start_revision 为上次记录的修订号 + 1 的 Watch。这样,etcd 就会把错过的所有事件按顺序重新推送给你,确保状态同步万无一失。
  • filters:允许在服务器端过滤掉不感兴趣的事件类型。例如,如果你只关心删除事件,可以设置 FilterType.NODELETE,这样服务器就不会发送任何 DELETE 类型的事件,从而节省网络带宽和客户端的处理开销。
  • prev_kv:如前所述,开启此选项后,PUT 和 DELETE 事件都会包含变更前的键值对,对于审计日志或实现回滚逻辑非常有用。

历史版本监听

“历史版本监听”是 Watch API 的一个核心应用场景,它使得 etcd 不仅仅是一个实时的配置中心,更是一个可靠的、可追溯的事件总线。

为什么需要监听历史版本?

考虑以下场景:

  1. 应用重启/网络闪断:应用在处理完一个事件后,与 etcd 的连接意外中断。在此期间,etcd 中的数据又发生了变化。当应用重新连接时,它需要知道错过了哪些事件,以便将本地状态恢复到最新。
  2. 初始化状态同步:一个新的服务实例启动,它需要获取当前所有相关的配置,并且还需要知道这些配置是如何演变到当前状态的,以便做出正确的初始决策。
  3. 审计与调试:需要完整地追溯某个键的所有变更历史。

如何实现历史监听?

实现历史监听的关键就是 WatchCreateRequest 中的 start_revision 字段。

让我们通过一个具体的例子来理解。假设我们执行了以下操作:

# 修订号 2
etcdctl put key1 value1

# 修订号 3
etcdctl put key2 value2

# 修订号 4
etcdctl delete key1

现在,我们的应用因为某些原因,在修订号 3 之后就断开了连接。当它在修订号 5 时重新启动,它需要同步从 4 开始的所有事件。

应用可以这样做:

  1. 首先,通过任何一次读取操作(例如 etcdctl get mykey -w=json)获取当前的修订号,假设是 5。
  2. 然后,它创建一个 Watch,指定 start_revision = 4。
message WatchCreateRequest {
  bytes key = "key1"; // 或者使用前缀监听
  int64 start_revision = 4;
}

etcd 服务器收到这个请求后,会检查历史记录,发现修订号 4 有一个对 key1 的删除操作。它会立即生成一个 DELETE 事件发送给客户端。这样,客户端就成功地“回放”了错过的事件,将本地状态与 etcd 保持一致。

这个机制是实现分布式锁、服务发现等高级功能的基础,因为它保证了状态的最终一致性和可恢复性。

Watch 进度通知

在长时间运行的 Watch 中,客户端可能会关心:“我的 Watch 还活着吗?我当前的数据有多新?” 特别是在网络不稳定的情况下,客户端可能无法确定是连接已经断开,还是仅仅因为长时间没有事件发生。

进度通知的作用

WatchProgressRequest 和 progress_notify 选项就是为了解决这个问题。它允许客户端主动或被动地获取 Watch 的“心跳”和当前进度。

两种使用方式

  1. 被动通知(progress_notify): 客户端在创建 Watch 时,可以将 progress_notify 字段设置为 true。

    message WatchCreateRequest {
      // ... other fields
      bool progress_notify = true;
    }
    

    设置后,如果 etcd 服务器在一段时间内没有产生任何新的事件,它会主动向客户端发送一个特殊的 WatchResponse。这个响应中 events 列表为空,但 header.revision 字段会包含服务器当前的最新修订号。 这对于客户端来说是一个信号:虽然没有数据变更,但 Watch 连接正常,并且我可以确认我的缓存数据至少更新到了这个修订号。

  2. 主动查询: 客户端也可以在 Watch 运行过程中,随时通过同一个流主动发送一个 WatchProgressRequest 消息。

    message WatchProgressRequest {
      // no fields
    }
    

    服务器收到这个请求后,会立即(或在下一个合适的时机)发送一个包含当前修订号的进度通知响应。这允许客户端在需要时(例如,在执行一次强一致性读之前)精确地了解 Watch 的最新进度。

进度通知机制是构建健壮的、可感知状态的客户端的关键,它使得客户端能够实现更精细的重连策略和状态同步逻辑。

Watch 实践案例

理论知识最终要落实到实践中。这里我们结合 etcdctl 命令行工具,演示几个常见的 Watch 使用场景。

场景一:实时监控单个配置项

假设我们有一个名为 config/app_setting 的配置项,我们希望在它被修改时立即得到通知。

终端 1:启动监听

etcdctl watch config/app_setting

终端 2:修改配置

etcdctl put config/app_setting "new_value"

终端 1 的输出:

PUT
config/app_setting
new_value

场景二:监听前缀(服务发现)

在微服务架构中,服务实例通常会将自己的地址注册到 etcd 的一个统一前缀下,例如 services/user-service/。网关需要监听这个前缀下的所有变化,以动态更新路由表。

终端 1:监听前缀

# --prefix 选项告诉 etcd 监听所有以 services/user-service/ 开头的键
etcdctl watch --prefix services/user-service/

终端 2:注册两个实例

etcdctl put services/user-service/192.168.1.10 "metadata"
etcdctl put services/user-service/192.168.1.11 "metadata"

终端 1 的输出:

PUT
services/user-service/192.168.1.10
metadata
PUT
services/user-service/192.168.1.11
metadata

终端 2:下线一个实例

etcdctl delete services/user-service/192.168.1.10

终端 1 的输出:

DELETE
services/user-service/192.168.1.10

场景三:断线重连与历史回溯

这是一个更复杂的场景,演示了如何利用 --rev 选项来恢复错过的事件。

步骤 1:准备数据

# 假设当前是修订号 10
etcdctl put mykey v1      # rev 11
etcdctl put mykey v2      # rev 12
etcdctl put mykey v3      # rev 13

步骤 2:模拟客户端断开连接前的状态

客户端记录下当前的修订号,假设是 12(即它已经收到了 v1 和 v2 的更新)。

步骤 3:在客户端断开期间,数据发生变化

# 在另一个终端执行
etcdctl put mykey v4      # rev 14
etcdctl put mykey v5      # rev 15

步骤 4:客户端重连并同步

客户端重新启动,它知道上次处理的修订号是 12。为了避免错过事件,它从 13 开始监听。

# 终端 1 (客户端)
etcdctl watch --rev=13 mykey

终端 1 的输出(立即回放历史并进入实时监听):

PUT
mykey
v3      # 回放 rev 13 的事件
PUT
mykey
v4      # 回放 rev 14 的事件
PUT
mykey
v5      # 回放 rev 15 的事件
# 此时空闲,等待新的事件

这个例子清晰地展示了 --rev 选项如何保证 Watch 的可恢复性,这是构建高可用分布式系统时不可或缺的特性。


本章我们深入探讨了 etcd 的 Watch API,从其底层的工作原理、事件流的处理模型,到如何监听历史版本和获取进度通知,并通过实践案例展示了其强大的功能。Watch 机制是 etcd 实现实时协调的核心,掌握它对于理解和使用 etcd 至关重要。下一章,我们将转向另一个核心 API——Lease API,它主要用于解决分布式环境下的节点存活检测问题。