watch
# 1.概述
# 1.1.watch机制
etcd提供了watch机制,用于避免感知数据变化时客户端的反复轮询。客户端会watch一系列key,这些key更新时,etcd就会通知客户端。watch机制也是kubernetes控制器的工作基础,各类控制器的workload就是围绕监听、比较资源状态进行协调工作,确保资源的最终一致性。
# 1.2.轮询与推送
etcd v2中的client通过HTTP/1.1协议长连接实现轮询server,获取最新的数据变化事件,方式会造成大量空负载。因此,etcd v3采用HTTP/2的gRPC协议、双向流的watch API设计,实现连接多路复用。HTTP/2协议的消息被分解独立的帧,交错发送,帧会标识关联流,一个数据流对应一个请求或响应包。形式上,一个client/TCP连接支持多gRPC stream,一个gRPC stream支持多个watcher。
# 1.3.事件存储
etcd v2使用简单的环形数组存储历史事件版本,eventQueue的容量固定是1000,超出的事件版本会覆盖,造成事件丢失。此时,client不得不触发大量的expensive查询操作获取最新的事件及版本号。type EventHistory struct { Queue eventQueue StartIndex uint64 LastIndex uint64 rwl sync.RWMutex } func newEventHistory(capacity int) *EventHistory { return &EventHistory{ Queue: eventQueue{ Capacity: capacity, Events: make([]*Event, capacity), // 1000 }, } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15etcdv3采用MVCC机制,将一个key的历史修改版本保存在boltdb,boltdb是一个基于磁盘文件的持久化存储。因此重启后历史事件不会丢失,可以通过配置压缩策略控制保存的历史版本数。
# 1.4.事件推送机制
发起
watch key请求时,etcd的gRPCWatchServer会创建一个serverWatchStream,负责接收client的gRPC Stream的create/cancel watcher请求,调用MVCC模块的WatchStream子模块分配一个watch id,将watcher注册到MVCC的WatchableKV模块,以便将MVCC模块接收的Watch事件转发给client。
WatchableKV模块会运行syncWatchersLoop和syncVictimsLoop,用于同步数据、记录失败event,通过异步任务执行推送重试。当涉及数据修改操作时,请求经过KVServer、Raft模块最终Apply到状态机时,MVCC的PUT事务会将本次修改的mvccpb.KeyValue保存在一个changes数组。PUT事务执行结束后,KeyValue会转换成Event事件,回调watchableStore.notify,将监听过此key及synced watcherGroup中的watcher匹配出来,将大等于监听版本的事件发送到watcher的channel。sendLoop监听到channel消息后,读出消息立即推给client,这就是最新事件推送阶段。
watcher的channel buffer默认容量为1024,由于事件堆积造成buffer满时,etcd会将此watcher从synced watcherGroup删除,将watcher和事件列表保存到victim的watcherBatch结构,通过异步机制重试保证事件的可靠性。syncVictimsLoop会遍历victim watcherBatch数据结构,尝试将堆积的事件再次推送给watcher channel,推送失败时会再次加入victim watcherBatch数据结构等待下次重试。推送成功时,回检watcher监听版本号,最小版本号超出server当前版本号时会加入synced watcher,进行最新事件通知机制。否则加入到unsynced watcherGroup,syncWatchersLoop会遍历unsynced watcherGroup每个watcher,批量推送历史事件。当某个watcher历史事件同步跟上后,又会加入synced watcherGroup开始最新事件同步。
快速定位[key,watcher]的实质就是上面的区间树
# 2.Client分析
# 2.1.核心数据
client对etcd客户端进行抽象,封装了监听回调、节点管理、grpc连接代理等模块。type Client struct { ... Watcher // 监听回调模块 ... conn *grpc.ClientConn //grpc连接代理,用于向server建立etcd长连接 ... }1
2
3
4
5
6
7watcher是etcd客户端的监听回调模块,内置用于和grpc服务端建立长连接的remote,维护streams字段,通过ctxKey映射多笔和服务端通信的长连接代理对象watchGrpcStream,同一笔watchGrpcStream一般会复用。type watcher struct { remote pb.WatchClient // 用于和grpc服务端建立长连接 ... mu sync.RWMutexg ... streams map[string]*watchGrpcStream // 映射ctxKey和长连接代理对象 }1
2
3
4
5
6
7watchGrpcStream是client与server间长连接的抽象,同时也是处理create/cancel watch请求以及服务端watch事件回调的中枢模块。type watchGrpcStream struct { owner *watcher // 所属watcher模块 remote pb.WatchClient // 建立长连接客户端 substreams map[int64]*watcherStream // key为watchId,value对应watcherStream子处理流 reqc chan watchStreamRequest // 接受watch请求 respc chan *pb.WatchResponse // 服务端响应或watch事件回调 }1
2
3
4
5
6
7watcherStream是某个特定watch的处理流抽象,用于监听及事件回调的缓冲。type watcherStream struct { initReq watchRequest // 应用放创建watch传递参数 outc chan WatchResponse // 将watch回调事件推到更上层的endpointManager使用的chan recvc chan *WatchResponse // 接收watch回调事件的chan buf []*WatchResponse // 缓存watch回调事件的缓冲 }1
2
3
4
5
6
# 2.2.watch
当发起一个
watch请求时,会进行请求组装、相关长连接代理对象的初始化、子处理流对象的初始化,并将watch对应的标识与子处理流对象绑定,实现key的监听及事件回调数据获取。func (w *watcher) Watch(ctx context.Context, key string, opts ...OpOption) WatchChan { ... // 请求相关 wr := &watchRequest{ ctx: ctx, createdNotify: ow.createdNotify, key: string(ow.key), end: string(ow.end), rev: ow.rev, progressNotify: ow.progressNotify, fragment: ow.fragment, filters: filters, prevKV: ow.prevKV, retc: make(chan chan WatchResponse, 1), } ok := false // 使用的ctxKey,用于关联长连接代理对象 ctxKey := streamKeyFromCtx(ctx) // 查找或分配适当的grpc流 w.mu.Lock() ... // streams是一个map,保存所有由ctx值键控的活动grpc流 wgs := w.streams[ctxKey] if wgs == nil { // 找不到ctxKey关联grpc流则新建 wgs = w.newWatcherGrpcStream(ctx) w.streams[ctxKey] = wgs } donec := wgs.donec reqc := wgs.reqc w.mu.Unlock() // 初始化closeChan closeCh := make(chan WatchResponse, 1) select { // 提交创建watcher请求 case reqc <- wr: ok = true ... } // 请求提交,等待响应 if ok { select { // 创建watcher成功后,返回推送监听事件的chan case ret := <-wr.retc: return ret case <-ctx.Done(): case <-donec: if wgs.closeErr != nil { closeCh <- WatchResponse{Canceled: true, closeErr: wgs.closeErr} break } // 重试 return w.Watch(ctx, key, opts...) } } close(closeCh) return closeCh }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
# 2.3.newWatcherGrpcStream
newWatcherGrpcStream会构造一个watcherGrpcStream实例,调用watcherGrpcStream.run向server建立长连接,通过持续轮询处理来自上层的请求及更底层的响应。func (w *watcher) newWatcherGrpcStream(inctx context.Context) *watchGrpcStream { ctx, cancel := context.WithCancel(&valCtx{inctx}) // 构造watchGrpcStream wgs := &watchGrpcStream{ owner: w, remote: w.remote, callOpts: w.callOpts, ctx: ctx, ctxKey: streamKeyFromCtx(inctx), cancel: cancel, substreams: make(map[int64]*watcherStream), respc: make(chan *pb.WatchResponse), reqc: make(chan watchStreamRequest), donec: make(chan struct{}), errc: make(chan error, 1), closingc: make(chan *watcherStream), resumec: make(chan struct{}), lg: w.lg, } // 启动协程处理监听key的watch各种事件 go wgs.run() return wgs }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 2.4.run
run方法会通过newWatchClient方法构造服务端的通信长连接,利用轮询处理来自更上层的请求和更底层的响应。针对上层,持续从reqc管道取出上层写入的watch请求,调用serveSubstream方法异步启动服务于当前watch的subStream。针对下层,持续从respc中接收来自etcd服务端的响应事件,调用watch.dispatchEvent方法,根据事件归属的watch将其分配给所属的watch subStream。func (w *watchGrpcStream) run() { ... // 启用底层grpc连接 if wc, closeErr = w.newWatchClient(); closeErr != nil { return } ... var cur *pb.WatchResponse backoff := time.Millisecond for { select { // 取出上层watch请求或progress请求 case req := <-w.reqc: switch wreq := req.(type) { // watch请求 case *watchRequest: // 初始化输出管道 outc := make(chan WatchResponse, 1) // 每个wr会建立虚拟stream(共用底层连接) ws := &watcherStream{ initReq: *wreq, id: InvalidWatchID, outc: outc, // 无缓冲管道为了避免重复事件 recvc: make(chan *WatchResponse), } // 用于通知上层子流结束 ws.donec = make(chan struct{}) w.wg.Add(1) // 启动子流 go w.serveSubstream(ws, w.resumec) // resuming和substreams就好像预备党员和党员的关系,初始化的ws只是预备役,还没收到事件响应 w.resuming = append(w.resuming, ws) if len(w.resuming) == 1 { // 保证wr顺序发送,wr收到响应才法下一个,否则排队 if err := wc.Send(ws.initReq.toPB()); err != nil { ... } } // 进度请求(强制推进watch流进度、检查watch流有效) case *progressRequest: // 发送进度请求 if err := wc.Send(wreq.toPB()); err != nil { ... } } // 处理grpc流接收到的事件 case pbresp := <-w.respc: // cur用于支持Fragment分片,由于一些响应较大,etcd server可能会分片发送,cur记录的分片就可以合成完整响应 if cur == nil || pbresp.Created || pbresp.Canceled { // 更新事件 cur = pbresp } else if cur != nil && cur.WatchId == pbresp.WatchId { // 合并事件 cur.Events = append(cur.Events, pbresp.Events...) // 更新分片标识 cur.Fragment = pbresp.Fragment } switch { // 处理创建事件 case pbresp.Created: if ws := w.resuming[0]; ws != nil { // 预备役的wr收到响应转为正式役 w.addSubstream(pbresp, ws) // 分发事件,其实就是根据watchId关联ws处理响应 w.dispatchEvent(pbresp) // 清空resuming第一个预备役ws w.resuming[0] = nil } // 顺序发送下一个预备役ws if ws := w.nextResume(); ws != nil { if err := wc.Send(ws.initReq.toPB()); err != nil { ... } } // 重置当前事件 cur = nil // 如果是取消事件 case pbresp.Canceled && pbresp.CompactRevision == 0: // 从cancelSet移除watchId delete(cancelSet, pbresp.WatchId) // 从正式役移除ws if ws, ok := w.substreams[pbresp.WatchId]; ok { // 关闭子流 close(ws.recvc) // 标记为关闭状态 closing[ws] = struct{}{} } // 重置当前事件 cur = nil // 如果是分片事件 case cur.Fragment: // 跳过,直到分片全发完再接收 continue // 其他情况 default: // 分发事件 ok := w.dispatchEvent(cur) // 重置当前事件 cur = nil if ok { break } // 分发失败,发送取消请求 if _, ok := cancelSet[pbresp.WatchId]; ok { break } // cancelSet去重,标记取消 cancelSet[pbresp.WatchId] = struct{}{} // 构造取消请求 cr := &pb.WatchRequest_CancelRequest{ CancelRequest: &pb.WatchCancelRequest{ WatchId: pbresp.WatchId, }, } req := &pb.WatchRequest{RequestUnion: cr} ... // 发送取消请求 if err := wc.Send(req); err != nil { ... } } // 处理流接收错误 case err := <-w.errc: ... // 处理子流关闭错误 case ws := <-w.closingc: ... } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
# 2.5.dispatchEvent
dispatchEvent将事件响应分发给watchId关联的观察者流,主要就是包装响应、传递给对应watchSubStream。func (w *watchGrpcStream) dispatchEvent(pbresp *pb.WatchResponse) bool { // 构造wr响应,其实就是事件对象、事件类型 wr := &WatchResponse{ Header: *pbresp.Header, Events: events, CompactRevision: pbresp.CompactRevision, Created: pbresp.Created, Canceled: pbresp.Canceled, cancelReason: pbresp.CancelReason, } ... // 通知给上游接收管道 return w.unicastResponse(wr, pbresp.WatchId) } func (w *watchGrpcStream) unicastResponse(wr *WatchResponse, watchId int64) bool { // 通过watchId关联subStream ws, ok := w.substreams[watchId] if !ok { return false } select { // 响应送入subStream case ws.recvc <- wr: case <-ws.donec: return false } return true }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
# 2.6.newWatchClient
newWatchClient用于开启一个grpc连接,然后启动协程监听服务端的事件响应。func (w *watchGrpcStream) newWatchClient() (pb.Watch_WatchClient, error) { // 新建或重启,正式役的ws回退到预备役 close(w.resumec) w.resumec = make(chan struct{}) w.joinSubstreams() for _, ws := range w.substreams { ws.id = InvalidWatchID w.resuming = append(w.resuming, ws) } // strip out nils, if any var resuming []*watcherStream for _, ws := range w.resuming { if ws != nil { resuming = append(resuming, ws) } } w.resuming = resuming w.substreams = make(map[int64]*watcherStream) // connect to grpc stream while accepting watcher cancelation stopc := make(chan struct{}) donec := w.waitCancelSubstreams(stopc) wc, err := w.openWatchClient() ... // 启动协程接收服务端响应,响应会转发给watchGrpcStream的respc管道 go w.serveWatchClient(wc) return wc, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
# 2.7.serveSubstream
某个
watch对应的subStream通过recv接收到create类型的watchResponse,代表server端完成对watch的创建请求,此时watchGrpcStream会将从watchClient获得的响应从respc转移到subStream的recvc管道,serveSubstream又会从recvc管道取出事件,根据事件将watch channel投递到watcherStream.initReq.retc中,供client.watch方法获取返回,用于交互client端交互事件响应。func (w *watchGrpcStream) serveSubstream(ws *watcherStream, resumec chan struct{}) { ... emptyWr := &WatchResponse{} for { curWr := emptyWr outc := ws.outc // outc就是watch channel // ws.buf是缓冲队列,这里处理队头数据(历史) if len(ws.buf) > 0 { curWr = ws.buf[0] } else { outc = nil } select { // 发送队列数据 case outc <- *curWr: ... ws.buf[0] = nil ws.buf = ws.buf[1:] // 处理当前事件 case wr, ok := <-ws.recvc: ... // 安装事件 if wr.Created { if ws.initReq.retc != nil { // outc通过retc发出去,会被Watch()本身的循环接收到 ws.initReq.retc <- ws.outc // 避免重试时的重新发送 ws.initReq.retc = nil // 设置WithCreatedNotify,wr会立刻发出,不会进入缓冲队列 if ws.initReq.createdNotify { ws.outc <- *wr } // recv为0 if ws.initReq.rev == 0 { nextRev = wr.Header.Revision } } } else { nextRev = wr.Header.Revision + 1 } // 更新当前的watch进度 if len(wr.Events) > 0 { nextRev = wr.Events[len(wr.Events)-1].Kv.ModRevision + 1 } ws.initReq.rev = nextRev ... // 缓冲事件 ws.buf = append(ws.buf, wr) case <-w.ctx.Done(): return case <-ws.initReq.ctx.Done(): return case <-resumec: resuming = true return } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
# 2.8.总结
--- watcher 初始化长连接抽象对象watchGrpcStream,通过watchGrpcStream提交请求,获取watch channel --- watchGrpcStream 总览大局,负责用户侧发来的watchRequest,负责watchClient侧发来的watchResponse,同时转发到适当的watchSubStream --- watchClient 基于send发送watchRequest,循环监听etcd服务端事件响应,将响应转发至watchGrpcStream的respc --- watcherStream 循环处理WatchResponse通过respc传递事件,基于用于侧watch channel(基于watchGrpcStream.retc传递到用户侧)发送响应,维护缓冲队列(来不及处理的消息)1
2
3
4
5
6
7
8
9
10
11
# 3.Server分析
# 3.1.watchableStore
watchableStore是MVCC和watch的关键模块,负责注册、管理及触发watcher的功能,每个watchableStore会组合store的字段和方法,同时维护两个watcherGroup管理watcher的同步进度。type watchableStore struct { *store // 存储模块 victims []watcherBatch // watcher channel阻塞时会暂时记录到恢复者 victimc chan struct{} // 新的watcherBatch实例加入victims,会向该管道发送消息 unsynced watcherGroup // 记录未同步的watcher synced watcherGroup // 已完成同步的watcher } type watcher struct { key []byte // 监听起始值 end []byte // 监听结束值 victim bool // watcher channel是否阻塞 compacted bool // 是否压缩 ... // 最小的revision minRev int64 id WatchID ... ch chan<- WatchResponse }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20watchableStore初始化时会启动syncWatchersLoop和syncVictimsLoop同步未同步的观察者以及发送失败的待恢复者。func newWatchableStore(log,Backend,lease,AuthStore,ConsistentIndexGetter,StoreConfig) *watchableStore { s := &watchableStore{ store: NewStore(lg, b, le, ig, cfg), victimc: make(chan struct{}, 1), unsynced: newWatcherGroup(), synced: newWatcherGroup(), stopc: make(chan struct{}), } s.store.ReadView = &readView{s} s.store.WriteView = &writeView{s} ... s.wg.Add(2) // 启动syncWatchersLoop go s.syncWatchersLoop() // 启动syncVictimsLoop go s.syncVictimsLoop() return s }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 3.2.syncWatchersLoop
syncWatchersLoop会每隔100ms调用syncWatchers同步未完成的watcher,保证watchableStore维护的观察者状态是最新的。func (s *watchableStore) syncWatchersLoop() { defer s.wg.Done() for { s.mu.RLock() st := time.Now() lastUnsyncedWatchers := s.unsynced.size() s.mu.RUnlock() unsyncedWatchers := 0 // 存在未同步的watcher,调用syncWatchers()同步 if lastUnsyncedWatchers > 0 { unsyncedWatchers = s.syncWatchers() } ... } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16syncWatchers会从未同步的观察者中选一组,迭代获得集合中观察者跟踪的最小recv。完成上述动作后,会使用最小修订获取键值并发送事件给这些观察者,同步成功的观察者会从未同步集合移除并加入同步组,这也是是syncWatchersLoop模块的核心。func (s *watchableStore) syncWatchers() int { s.mu.Lock() defer s.mu.Unlock() // 不存在未同步观察者 if s.unsynced.size() == 0 { return 0 } // 涉及到键值读,加锁 s.store.revMu.RLock() defer s.store.revMu.RUnlock() curRev := s.store.currentRev compactionRev := s.store.compactMainRev // 根据当前版本从未同步组选一些watcher及范围内最小版本 wg, minRev := s.unsynced.choose(maxWatchersPerSync, curRev, compactionRev) minBytes, maxBytes := newRevBytes(), newRevBytes() revToBytes(revision{main: minRev}, minBytes) revToBytes(revision{main: curRev + 1}, maxBytes) // 从boltdb取出当前版本范围内的数据及转换为事件 tx := s.store.b.ReadTx() tx.RLock() // 从区间树返回键及查到的值,键是revision,值是实际键值数据 revs, vs := tx.UnsafeRange(keyBucketName, minBytes, maxBytes, 0) var evs []mvccpb.Event if s.store != nil && s.store.lg != nil { evs = kvsToEvents(s.store.lg, wg, revs, vs) } tx.RUnlock() var victims watcherBatch // 观察者映射到匹配事件,其实就是map wb := newWatcherBatch(wg, evs) for w := range wg.watchers { ... // 事件发送到watcher channel if w.send(WatchResponse{WatchID: w.id, Events: eb.evs, Revision: curRev}) { pendingEventsGauge.Add(float64(len(eb.evs))) } ... // 发送失败加入victims if w.victim { victims[w] = eb } else { ... // 加入已同步观察者 s.synced.add(w) } // 从未同步观察者删除 s.unsynced.delete(w) } // 加入victims组 s.addVictim(victims) ... // 返回未同步观察者数量 return s.unsynced.size() }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
# 3.3.syncVictimsLoop
接下来就是
syncVictimsLoop,负责重试发送阻塞的watcher,将事件重新发送到watcher channel,同时根据发送状态将watcher移交到同步组。func (s *watchableStore) syncVictimsLoop() { defer s.wg.Done() for { // 尝试处理所有阻塞watcher for s.moveVictims() != 0 { // try to update all victim watchers } ... // 阻塞watcher不为空时间隔10ms再次触发 } }1
2
3
4
5
6
7
8
9
10
11
12moveVictims方法是阻塞watcher同步核心,用于触发阻塞watcher的重试以及转移。func (s *watchableStore) moveVictims() (moved int) { s.mu.Lock() ... for _, wb := range victims { // 尝试再次发送 for w, eb := range wb { // event事件包装为watchRresponse,写入watcher channel if w.send(WatchResponse{WatchID: w.id, Events: eb.evs, Revision: rev}) { pendingEventsGauge.Add(float64(len(eb.evs))) } else { ... // 失败继续放回Victim组 newVictim[w] = eb continue } moved++ } s.mu.Lock() s.store.revMu.RLock() // 将victim分配到unsync/sync组 curRev := s.store.currentRev for w, eb := range wb { // 发送失败的保留在victim组 if newVictim != nil && newVictim[w] != nil { // couldn't send watch response; stays victim continue } w.victim = false ... // 观察者要求的rev小于数据最新的curv,watcher放入unsync组 if w.minRev <= curRev { s.unsynced.add(w) // 否则watcher放入sync组 } else { slowWatcherGauge.Dec() s.synced.add(w) } } s.store.revMu.RUnlock() s.mu.Unlock() } ... return moved }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
# 3.4.watchServer
etcd server启动时会注册多个grpc server,涉及到watch操作的是watch server,用户处理客户端的watch请求,并根据注册的watcher发送事件响应。func Server(*EtcdServer,*tls.Config,grpc.UnaryServerInterceptor,...grpc.ServerOption) *grpc.Server { grpcServer := grpc.NewServer(append(opts, gopts...)...) ... // 注册 WatchServer pb.RegisterWatchServer(grpcServer, NewWatchServer(s)) return grpcServer }1
2
3
4
5
6
7
8watchServer只用于处理watch请求,内部关联实现mvcc.WatchableKV接口的watchableStore,用于监听存储中数据变更。type WatchServer interface { // 处理watch的接口抽象 Watch(Watch_WatchServer) error } type watchServer struct { ... watchable mvcc.WatchableKV // 这里其实就是watchableStore,用于监听数据变更 ... } // 核心数据结构 type serverWatchStream struct { ... watchable mvcc.WatchableKV // KV存储,其实就是watchableStore实现 ... gRPCStream pb.Watch_WatchServer // 连接客户端的stream watchStream mvcc.WatchStream // key变动的消息管道 ctrlStream chan *pb.WatchResponse // 响应客户端请求的消息管道 ... progress map[mvcc.WatchID]bool // 进度请求,作为心跳探测stream有效 prevKV map[mvcc.WatchID]bool // 范围监听需要通知的watcher,a/b中的b变化,a也要通知 fragment map[mvcc.WatchID]bool // 分片标识,大数据拆分发送,对应客户端接收时的分片事件判断 ... }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
# 3.5.watch
etcd grpcWatchServer收到watch请求后,创建serverWatchStream,负责接收client的gRPC stream的create/cancel watcher请求——recvLoop模块,将从MVCC模块接收的watch事件转发给client——sendLoop模块。func (ws *watchServer) Watch(stream pb.Watch_WatchServer) (err error) { // 基于watchServer初始化serverWatchStream sws := serverWatchStream{ ... } sws.wg.Add(1) // 启动sendLoop,转发watch事件 go func() { sws.sendLoop() sws.wg.Done() }() ... // 启动recvLoop go func() { if rerr := sws.recvLoop(); rerr != nil { ... } }() ... sws.close() return err }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25注意
并非每次调用
watch就会创建一个serverWatchStream,而是每个grpc Stream创建一个,创建出的serverWatchStream会接管这个grpc Stream上所有的watcher。此外,一个连接可以创建多个grpc Stream,这也是etcd v3的watch相比etcd v2在资源上没有大幅降低的原因。
# 3.6.sendLoop
存储中出现
更新或删除事件时,事件会被送到watchStream持有的channel,sendLoop又会通过select监听多个channel中的数据,将接收到的数据封装成pb.WatchResponse结构并通过grpc Stream发送给客户端。func (sws *serverWatchStream) sendLoop() { ... for { select { // key变化的消息管道 case wresp, ok := <-sws.watchStream.Chan(): if !ok { return } // 整理事件 evs := wresp.Events events := make([]*mvccpb.Event, len(evs)) sws.mu.RLock() needPrevKV := sws.prevKV[wresp.WatchID] sws.mu.RUnlock() for i := range evs { events[i] = &evs[i] ... } ... // 事件响应发送给客户端 if !fragmented && !ok { // 完整发送 serr = sws.gRPCStream.Send(wr) } else { // 分片发送 serr = sendFragments(wr, sws.maxRequestBytes, sws.gRPCStream.Send) } ... // 响应客户端请求的消息管道 case c, ok := <-sws.ctrlStream: ... // 请求结果返回给客户端 if err := sws.gRPCStream.Send(c); err != nil { ... return } // track id creation wid := mvcc.WatchID(c.WatchId) // 请求取消,从ids列表删除 if c.Canceled && wid != clientv3.InvalidWatchID { delete(ids, wid) continue } // 创建事件,把watchId注册到ids列表,将缓存的event都发送到client if c.Created { ids[wid] = struct{}{} for _, v := range pending[wid] { // 发送事件 if err := sws.gRPCStream.Send(v); err != nil { ... return } } // 发送成功的从pending列表删除缓存的watcher event delete(pending, wid) } // 进度请求 case <-progressTicker.C: sws.mu.Lock() for id, ok := range sws.progress { if ok { // 定时发送进度响应,类似心跳包 sws.watchStream.RequestProgress(id) } sws.progress[id] = true } sws.mu.Unlock() // serverWatchStream关闭 case <-sws.closec: return } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76注意
1.
<-sws.watchStream.Chan()会不断获取数据操作引发的变更事件,将事件包装为watchResponse后发送到client。由于watcher是在recvLoop创建。然后才会通过消息异步将watchId注册到sendLoop的ids列表,因此需要将watchId未注册的event缓存起来。2.
recvLoop收到create/cancel watch请求后调用watchStream的create/cancel watcher,结果通过sws.ctrlStream管道送到sendLoop,以维护sendLoop中活跃的watchId,以及针对本次watch请求反馈给客户端的响应3.
<-progressTicker.C会不断处理进度请求的响应,这种请求用于探测grpc Stream连接状态,类似心跳检测
# 3.7.recvLoop
recvLoop处理客户端的watch请求,根据请求类型创建或取消watcher,将请求结果放入sws.ctrlStream管道,借助sendLoop将请求结果响应给client。此外,watch动作最终借助watchableStore.Watch实现。func (sws *serverWatchStream) recvLoop() error { for { // 获取client通过grpc Stream发起的请求 req, err := sws.gRPCStream.Recv() ... // 执行不同请求动作 switch uv := req.RequestUnion.(type) { case *pb.WatchRequest_CreateRequest: ... creq := uv.CreateRequest ... // watch权限检查 err := sws.isWatchPermitted(creq) if err != nil { ... // 无权限watch响应 wr := &pb.WatchResponse{ Header: sws.newResponseHeader(sws.watchStream.Rev()), WatchId: clientv3.InvalidWatchID, Canceled: true, Created: true, CancelReason: cancelReason, } select { // 写入请求结果管道,等待sendLoop处理 case sws.ctrlStream <- wr: continue case <-sws.closec: return nil } } ... // 创建本次请求的watcher id, err := sws.watchStream.Watch(mvcc.WatchID(creq.WatchId), creq.Key, creq.RangeEnd, rev, filters...) ... // 响应watcher创建结果 wr := &pb.WatchResponse{ Header: sws.newResponseHeader(wsrev), WatchId: int64(id), Created: true, Canceled: err != nil, } ... select { // 请求结果管道 case sws.ctrlStream <- wr: case <-sws.closec: return nil } // 取消请求 case *pb.WatchRequest_CancelRequest: if uv.CancelRequest != nil { id := uv.CancelRequest.WatchId // 取消watcher err := sws.watchStream.Cancel(mvcc.WatchID(id)) if err == nil { // 返回取消结果 sws.ctrlStream <- &pb.WatchResponse{...} ... } } // 进度请求 case *pb.WatchRequest_ProgressRequest: ... default: // 等待下一轮处理 continue } } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
# 3.8.watchStream.watch
watcher的创建通过调用sws.watchStream.Watch外包给了MVCC的watchStream,watchStream为watcher生成最新的watchId,调用MVCC存储的watch()创建watcher实例、更新revision,对应的结果包装为watchResponse后通过消息发送到sendLoop,实现watchId注册到活跃ids列表。type watchStream struct { watchable watchable // 记录关联的watchableStore ch chan WatchResponse // event事件写入管道 ... watchers map[WatchID]*watcher // 记录watchId与watcher关系 }1
2
3
4
5
6watchStream通过调用MVCC.watch创建watcher,并维护分配的watchId与watcher的关系,这里的MVCC模块其实就是watchableStore,会在watchServer初始化时调用watchable.NewWatchStream()初始化。func (ws *watchStream) Watch(id WatchID, key, end []byte, startRev int64, fcs ...FilterFunc) (WatchID, error) { ... ws.mu.Lock() defer ws.mu.Unlock() if ws.closed { return -1, ErrEmptyWatcherRange } // watch还未分配,分配一个自增ID if id == clientv3.AutoWatchID { for ws.watchers[ws.nextID] != nil { ws.nextID++ } id = ws.nextID ws.nextID++ } // 创建watcher,更新recv w, c := ws.watchable.watch(key, end, startRev, id, ws.ch, fcs...) ws.cancels[id] = c // 维护watchId-->watcher映射 ws.watchers[id] = w return id, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25ws.watchable.watch()创建实际的watcher实例,根据watcher的同步情况将watcher加入watchableStore.watcherGroup,watcherGroup是一个用于管理watcher的集合。// 又走到watchableStore的实现了 func (s *watchableStore) watch(key,end []byte,startRev int64,id WatchID,ch chan<- WatchResponse,fcs ...FilterFunc) (*watcher, cancelFunc) { // 初始化watcher wa := &watcher{ key: key, end: end, minRev: startRev, id: id, ch: ch, fcs: fcs, } s.mu.Lock() s.revMu.RLock() // 检查watcher同步情况 synced := startRev > s.store.currentRev || startRev == 0 // 创建场景、事件跟上最新的场景认为已同步 if synced { // 更新跟踪的版本 wa.minRev = s.store.currentRev + 1 if startRev > wa.minRev { wa.minRev = startRev } } // watcher加入已同步组 if synced { s.synced.add(wa) // watcher加入未同步组 } else { slowWatcherGauge.Inc() s.unsynced.add(wa) } s.revMu.RUnlock() s.mu.Unlock() watcherGauge.Inc() return wa, func() { s.cancelWatcher(wa) } }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39总结
1.
s.synced.add()、s.unsynced.add()会将watcher加入到watcherGroup中,即根据监听进度在watchableStore中记录已完成同步和未完成同步的watcher2.创建
watcher的请求响应推送到sws.ctrlStream channel,由sendLoop异步处理并返回给client3.数据变化或重试时,
syncWatchersLoop和syncVictimsLoop会将事件通过watcher.send()发送到watcher channel,对应上述分析的sendLoop监听watcher channel获取事件响应
# 4.总结
# 4.1.client
--- 核心流程 1.`client.Watch()`发起监听请求 2.根据`ctxKey`获取或初始化`WatcherGrpcStream`实例,`WatcherGrpcStream`是建立`server grpc stream`的长连接抽象 3.创建`WatcherGrpcStream`时调用`run()`方法,通过`newWatchClient`创建`wc`连接`grpc server`,同时启动协程监听服务端响应放入`watchGrpcStream`的`respc channel` 4.请求事件到来时,通过建立的`wc`将请求送到服务端,启动服务于当前`watch`的`serveSubstream`服务,监听服务端的响应 5.`watchGrpcStream`的`respc channel`收到响应事件时,按照`FIFO`分发给`watch`对应`watchStream recvc channel` 6.`serveSubstream`持续监听`recvc channel`事件,转发到`watcherStream.initReq.retc`,这个`channel`其实就是`client.watch()`请求时返回给客户端的`ret channel`,也是`watchRequest retc channel`,最终客户端会根据这个管道获取服务端针对`watch`的事件。1
2
3
4
5
6
7
注意
wc.send(req)—serveWatchClient接收resp写入watchGrpcStream.respc—dispatchEvent分发到watcherStream.recvc—serveSubstream写入watchRequest.retc—retc已经返回给客户端,可直接拿到响应
# 4.2.server
--- 核心流程 1.etcd服务端创建newWatchableStore开启group监听 2.调用mvcc中syncWatchersLoop将所有未通知的事件通知给所有的监听者; 3.对watcher通道阻塞时存入victim中数据,开启syncVictimsLoop; 4.watchServer响应客户端请求,发起watchStream及watcher实例新建,并将其添加至unsynced或synced组 5.client端通过grpc proxy向watcherServer发送watcher请求,recvLoop处理请求,sendLoop获取事件并响应 6.grpc proxy提供对同一个key的多次watch合并减少etcd server中重复watcher创建,以提高etcd server稳定性1
2
3
4
5
6
7
补充
由于本篇未分析
MVCC存储,所以未体现出数据变更时对watcher的通知,其实数据操作经过watchable存储后会生成新的版本,通过调用watchableStore.notify将事件发送到每个watcher channel,对上sendLoop对watcher channel的监听


