监听机制
etcd Watch 机制完整指南:Watch 创建与取消、前缀/范围监听、事件类型详解、从指定 revision 恢复、重连策略,以及与轮询的对比。
核心概念
Watch 是 etcd 最重要的特性之一:客户端创建一个 Watcher,服务端持续推送匹配 Key 的变更事件。
Client Server
│ │
│ ── Watch(key="/svc/") ──►│
│ │
│ (另一个客户端: │
│ Put /svc/host1 "A")│
│ │
│ ◄─ Event(PUT, /svc/host1, "A", rev=5) ──│
│ │
│ (另一个客户端: │
│ Put /svc/host1 "B")│
│ │
│ ◄─ Event(PUT, /svc/host1, "B", rev=6) ──│关键保证:
- 基于 revision 的顺序通知,不丢事件
- 支持从指定 revision 恢复(只要 revision 未被 compact)
- 基于 gRPC 双向流,实时推送,延迟极低
基本用法
创建与接收
import (
"context"
"log"
clientv3 "go.etcd.io/etcd/client/v3"
)
// 监听单个 Key(阻塞当前 goroutine)
watchCh := cli.Watch(context.Background(), "/config/db/host")
for resp := range watchCh {
if resp.Err() != nil {
log.Printf("Watch 错误: %v", resp.Err())
return
}
for _, event := range resp.Events {
fmt.Printf("Type=%s Key=%s Value=%s Rev=%d\n",
event.Type, event.Kv.Key, event.Kv.Value, event.Kv.ModRevision)
}
}前缀监听(最常用)
// 监听 /service/ 下所有变更
watchCh := cli.Watch(ctx, "/service/", clientv3.WithPrefix())
for resp := range watchCh {
for _, event := range resp.Events {
switch event.Type {
case mvccpb.PUT:
if event.IsCreate() {
fmt.Printf("新增服务: %s = %s\n", event.Kv.Key, event.Kv.Value)
} else {
fmt.Printf("更新服务: %s = %s\n", event.Kv.Key, event.Kv.Value)
}
case mvccpb.DELETE:
fmt.Printf("服务下线: %s\n", event.Kv.Key)
}
}
}Event 类型与判断
type Event struct {
Type mvccpb.Event_EventType // PUT 或 DELETE
Kv *mvccpb.KeyValue // 当前值
PrevKv *mvccpb.KeyValue // 旧值(需 WithPrevKV)
}
// Event 类型判断
event.IsCreate() // Type == PUT && CreateRevision == ModRevision(首次创建)
event.IsModify() // Type == PUT && CreateRevision != ModRevision(更新)
event.Type == mvccpb.PUT // 写入事件(包含创建和修改)
event.Type == mvccpb.DELETE // 删除事件Watch 选项速查
| Option | 说明 | 示例 |
|---|---|---|
WithPrefix() |
前缀匹配监听 | Watch(ctx, "/svc/", WithPrefix()) |
WithRange(end) |
范围监听 [key, end) | Watch(ctx, "/a", WithRange("/z")) |
WithRev(n) |
从指定 revision 开始监听 | Watch(ctx, "/foo", WithRev(100)) |
WithPrevKV() |
返回变更前的旧值 | Watch(ctx, "/foo", WithPrevKV()) |
WithProgressNotify() |
定期收到空响应(检测连接活性) | Watch(ctx, "/foo", WithProgressNotify()) |
WithFilterDelete() |
过滤掉 DELETE 事件 | Watch(ctx, "/foo", WithFilterDelete()) |
WithFilterPut() |
过滤掉 PUT 事件 | — |
WithCreatedNotify() |
创建 Watcher 时收到一个空响应标记创建完成 | — |
WithFragment() |
允许服务端分包推送大响应 | — |
多个 Watch 组合
// 用 NewWatcher 创建独立的 Watcher 实例
watcher := clientv3.NewWatcher(cli)
// 同时监听多个 Key 范围
go func() {
ch := watcher.Watch(ctx, "/config/", clientv3.WithPrefix())
for resp := range ch {
// 处理配置变更
}
}()
go func() {
ch := watcher.Watch(ctx, "/service/", clientv3.WithPrefix())
for resp := range ch {
// 处理服务变更
}
}()
// 或者一次创建多个 watch
watchCtx, cancel := context.WithCancel(context.Background())
defer cancel()
ch1 := watcher.Watch(watchCtx, "/config/", clientv3.WithPrefix())
ch2 := watcher.Watch(watchCtx, "/service/", clientv3.WithPrefix())Watch 响应结构
type WatchResponse struct {
Header *ResponseHeader
Events []*Event // 事件列表
CompactRevision int64 // 当 >= 0 时,表示历史已被压缩到此 revision
Created bool // 是否刚创建 Watcher
Canceled bool // 是否被取消
CancelReason string // 取消原因
}
// 典型的处理循环
for resp := range watchCh {
// 1. 先检查错误
if resp.Err() != nil {
log.Printf("watch error: %v", resp.Err())
break
}
// 2. 检查是否被取消
if resp.Canceled {
log.Printf("watch canceled: %s", resp.CancelReason)
break
}
// 3. 处理事件
for _, ev := range resp.Events {
handleEvent(ev)
}
}从 Revision 恢复 Watch
这是 Watch 能力最重要的使用模式:先批量加载当前状态,再 Watch 增量变更,保证不丢事件。
// 1. 获取当前所有服务
resp, err := cli.Get(ctx, "/service/", clientv3.WithPrefix())
if err != nil {
log.Fatal(err)
}
// 记录当前所有服务
for _, kv := range resp.Kvs {
addService(kv)
}
// 2. 从当前 revision+1 开始监听后续变更
// 保证不会漏掉 Get 之后的任何事件
watchRev := resp.Header.Revision + 1
go func() {
watchCh := cli.Watch(context.Background(),
"/service/",
clientv3.WithPrefix(),
clientv3.WithRev(watchRev),
)
for resp := range watchCh {
if resp.Err() != nil {
log.Printf("watch error: %v", resp.Err())
return
}
for _, ev := range resp.Events {
switch ev.Type {
case mvccpb.PUT:
addService(ev.Kv)
case mvccpb.DELETE:
removeService(ev.Kv.Key)
}
}
}
}()Compaction 错误的处理
func watchWithCompactHandling(cli *clientv3.Client, key string) {
var watchRev int64 = 0
for {
opts := []clientv3.OpOption{clientv3.WithPrefix()}
if watchRev > 0 {
opts = append(opts, clientv3.WithRev(watchRev))
}
watchCh := cli.Watch(context.Background(), key, opts...)
for resp := range watchCh {
// 🔑 关键:检查 Compact 错误
if resp.CompactRevision > 0 {
log.Printf("Revision %d 已被压缩,需要全量重新同步", resp.CompactRevision)
// 重新 Get 全量数据
getResp, err := cli.Get(ctx, key, clientv3.WithPrefix())
if err != nil {
log.Fatal(err)
}
rebuildFromFullSync(getResp.Kvs)
watchRev = getResp.Header.Revision + 1
break // 跳出内层循环,重新创建 Watcher
}
if resp.Err() != nil {
log.Printf("Watch 错误: %v", resp.Err())
watchRev = 0 // 从头开始 watch
break
}
for _, ev := range resp.Events {
handleEvent(ev)
}
}
}
}重连与重试策略
etcd 的 Watch 在连接断开后会自动重连,但需要业务层配合处理:
func watchForever(cli *clientv3.Client, key string) {
for {
// 获取当前最新 revision
getResp, err := cli.Get(ctx, key, clientv3.WithPrefix())
rev := int64(0)
if err == nil {
rev = getResp.Header.Revision + 1
}
// 创建 Watcher
opts := []clientv3.OpOption{clientv3.WithPrefix()}
if rev > 0 {
opts = append(opts, clientv3.WithRev(rev))
}
watchCh := cli.Watch(ctx, key, opts...)
for resp := range watchCh {
err := resp.Err()
if err == rpctypes.ErrGRPCCompacted {
// 被 compact 了,需要重新全量读取
break
}
if err != nil {
log.Printf("Watch 错误,%v 后重试", err)
time.Sleep(time.Second)
break
}
for _, ev := range resp.Events {
handleEvent(ev)
}
}
}
}性能与注意事项
⚡ 性能提示:
| 场景 | 建议 |
|---|---|
| 大量并发 watch | 单个客户端建议 < 1000 个 watcher;超过时考虑按 range 合并 watch |
| 高频写入 | 一个 prefix 下如果有大量写入,考虑拆分为更细的 key range |
| 大 cluster | etcd 服务端为每个 watcher 维护一个推送队列,watcher 消费慢会占用内存 |
| Progress Notify | 建议开启 WithProgressNotify(),定期验证连接活性 |
Watch vs 轮询
| 维度 | Watch | 轮询(Polling) |
|---|---|---|
| 延迟 | 毫秒级实时推送 | 取决于轮询间隔(秒级) |
| 资源开销 | 维持长连接,事件驱动 | 每次轮询一次全文扫描 |
| 事件可靠性 | 基于 revision,不丢事件 | 轮询间隙可能丢失中间状态 |
| 适用场景 | 变更通知、配置热更新 | 定期检查(如健康检查) |
💡 最佳实践:优先使用 Watch 监听变更;健康检查等定期场景用轮询配合
WithSerializable()。
常见陷阱
| 陷阱 | 说明 |
|---|---|
| 🚨 Compaction 后 Watch 断连 | 从已压缩 revision watch 会收到 ErrCompacted,必须全量重新同步 |
| 🚨 闭包捕获错误 | goroutine 中 watch 回调引用的外部变量需要注意线程安全 |
| 🚨 忘记 WithPrefix | Watch(ctx, "/foo") 只监听 /foo 这一个 Key,不包含子 Key |
| 🚨 Watch channel 不消费导致内存泄漏 | 服务端为每个 watcher 维护推送队列,消费慢会占用 etcd 内存 |
| 🚨 revision+1 超出当前 revision | WithRev(revision+1) 可以安全指定 future revision——etcd 会阻塞等到该 revision 发生后再发送事件 |