wal: send watch events only when needed
This commit is contained in:
parent
1e1ba89a3f
commit
0f9a5f9c4b
|
@ -480,6 +480,7 @@ func (w *WalManager) Watch(ctx context.Context, revision int64) <-chan *WatchEle
|
||||||
defer close(walCh)
|
defer close(walCh)
|
||||||
for wresp := range wch {
|
for wresp := range wch {
|
||||||
we := &WatchElement{ChangeGroupsRevisions: make(changeGroupsRevisions)}
|
we := &WatchElement{ChangeGroupsRevisions: make(changeGroupsRevisions)}
|
||||||
|
send := false
|
||||||
|
|
||||||
if wresp.Canceled {
|
if wresp.Canceled {
|
||||||
err := wresp.Err()
|
err := wresp.Err()
|
||||||
|
@ -501,6 +502,7 @@ func (w *WalManager) Watch(ctx context.Context, revision int64) <-chan *WatchEle
|
||||||
|
|
||||||
switch {
|
switch {
|
||||||
case strings.HasPrefix(key, etcdWalsDir+"/"):
|
case strings.HasPrefix(key, etcdWalsDir+"/"):
|
||||||
|
send = true
|
||||||
switch ev.Type {
|
switch ev.Type {
|
||||||
case mvccpb.PUT:
|
case mvccpb.PUT:
|
||||||
var walData *WalData
|
var walData *WalData
|
||||||
|
@ -514,6 +516,7 @@ func (w *WalManager) Watch(ctx context.Context, revision int64) <-chan *WatchEle
|
||||||
}
|
}
|
||||||
|
|
||||||
case strings.HasPrefix(key, etcdChangeGroupsDir+"/"):
|
case strings.HasPrefix(key, etcdChangeGroupsDir+"/"):
|
||||||
|
send = true
|
||||||
switch ev.Type {
|
switch ev.Type {
|
||||||
case mvccpb.PUT:
|
case mvccpb.PUT:
|
||||||
changeGroup := path.Base(string(ev.Kv.Key))
|
changeGroup := path.Base(string(ev.Kv.Key))
|
||||||
|
@ -523,13 +526,18 @@ func (w *WalManager) Watch(ctx context.Context, revision int64) <-chan *WatchEle
|
||||||
we.ChangeGroupsRevisions[changeGroup] = 0
|
we.ChangeGroupsRevisions[changeGroup] = 0
|
||||||
}
|
}
|
||||||
|
|
||||||
|
case key == etcdPingKey:
|
||||||
|
send = true
|
||||||
|
|
||||||
default:
|
default:
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if send {
|
||||||
walCh <- we
|
walCh <- we
|
||||||
}
|
}
|
||||||
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
return walCh
|
return walCh
|
||||||
|
|
Loading…
Reference in New Issue