Skip to content
监听机制

监听机制

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 发生后再发送事件