From d1360e372116e519a90e2ba7f032ec180f72ad26 Mon Sep 17 00:00:00 2001 From: Ben Lockhart Date: Fri, 9 May 2025 09:10:36 -0400 Subject: [PATCH 1/2] Add logging to binlog watcher actions Signed-off-by: Ben Lockhart --- go/vt/vttablet/tabletserver/binlog_watcher.go | 1 + go/vt/vttablet/tabletserver/state_manager.go | 2 ++ go/vt/vttablet/tabletserver/vstreamer/engine.go | 4 ++++ 3 files changed, 7 insertions(+) diff --git a/go/vt/vttablet/tabletserver/binlog_watcher.go b/go/vt/vttablet/tabletserver/binlog_watcher.go index 80ac1194c7e..220debd1e76 100644 --- a/go/vt/vttablet/tabletserver/binlog_watcher.go +++ b/go/vt/vttablet/tabletserver/binlog_watcher.go @@ -92,6 +92,7 @@ func (blw *BinlogWatcher) process(ctx context.Context) { } for { + log.Info("Binlog watcher: streaming") // VStreamer will reload the schema when it encounters a DDL. err := blw.vs.Stream(ctx, "current", nil, filter, throttlerapp.BinlogWatcherName, func(events []*binlogdatapb.VEvent) error { return nil diff --git a/go/vt/vttablet/tabletserver/state_manager.go b/go/vt/vttablet/tabletserver/state_manager.go index 0ccd0e42735..7bb5ff52e04 100644 --- a/go/vt/vttablet/tabletserver/state_manager.go +++ b/go/vt/vttablet/tabletserver/state_manager.go @@ -447,7 +447,9 @@ func (sm *stateManager) verifyTargetLocked(ctx context.Context, target *querypb. } func (sm *stateManager) servePrimary() error { + log.Info("servePrimary: closing binlog watcher") sm.watcher.Close() + log.Info("servePrimary: binlog watcher closed") if err := sm.connect(topodatapb.TabletType_PRIMARY, true); err != nil { return err diff --git a/go/vt/vttablet/tabletserver/vstreamer/engine.go b/go/vt/vttablet/tabletserver/vstreamer/engine.go index eba6e736a21..6f089305ad1 100644 --- a/go/vt/vttablet/tabletserver/vstreamer/engine.go +++ b/go/vt/vttablet/tabletserver/vstreamer/engine.go @@ -270,10 +270,14 @@ func (vse *Engine) Stream(ctx context.Context, startPos string, tablePKs []*binl // Remove stream from map and decrement wg when it ends. defer func() { + log.Info("VStreamer engine: locking for close") vse.mu.Lock() defer vse.mu.Unlock() + log.Info("VStreamer engine: deleting streamers") delete(vse.streamers, idx) + log.Info("VStreamer engine: decrementing wait group") vse.wg.Done() + log.Info("VStreamer engine: unlocking") }() // No lock is held while streaming, but wg is incremented. From 8cfa31620e7c52cd2ff19b099e052b2c4b602882 Mon Sep 17 00:00:00 2001 From: Ben Lockhart Date: Tue, 3 Jun 2025 16:37:19 -0400 Subject: [PATCH 2/2] remove most new log lines, leave 'watcher closed' in Signed-off-by: Ben Lockhart --- go/vt/vttablet/tabletserver/binlog_watcher.go | 1 - go/vt/vttablet/tabletserver/state_manager.go | 1 - go/vt/vttablet/tabletserver/vstreamer/engine.go | 4 ---- 3 files changed, 6 deletions(-) diff --git a/go/vt/vttablet/tabletserver/binlog_watcher.go b/go/vt/vttablet/tabletserver/binlog_watcher.go index 220debd1e76..80ac1194c7e 100644 --- a/go/vt/vttablet/tabletserver/binlog_watcher.go +++ b/go/vt/vttablet/tabletserver/binlog_watcher.go @@ -92,7 +92,6 @@ func (blw *BinlogWatcher) process(ctx context.Context) { } for { - log.Info("Binlog watcher: streaming") // VStreamer will reload the schema when it encounters a DDL. err := blw.vs.Stream(ctx, "current", nil, filter, throttlerapp.BinlogWatcherName, func(events []*binlogdatapb.VEvent) error { return nil diff --git a/go/vt/vttablet/tabletserver/state_manager.go b/go/vt/vttablet/tabletserver/state_manager.go index 7bb5ff52e04..f53fce59288 100644 --- a/go/vt/vttablet/tabletserver/state_manager.go +++ b/go/vt/vttablet/tabletserver/state_manager.go @@ -447,7 +447,6 @@ func (sm *stateManager) verifyTargetLocked(ctx context.Context, target *querypb. } func (sm *stateManager) servePrimary() error { - log.Info("servePrimary: closing binlog watcher") sm.watcher.Close() log.Info("servePrimary: binlog watcher closed") diff --git a/go/vt/vttablet/tabletserver/vstreamer/engine.go b/go/vt/vttablet/tabletserver/vstreamer/engine.go index 6f089305ad1..eba6e736a21 100644 --- a/go/vt/vttablet/tabletserver/vstreamer/engine.go +++ b/go/vt/vttablet/tabletserver/vstreamer/engine.go @@ -270,14 +270,10 @@ func (vse *Engine) Stream(ctx context.Context, startPos string, tablePKs []*binl // Remove stream from map and decrement wg when it ends. defer func() { - log.Info("VStreamer engine: locking for close") vse.mu.Lock() defer vse.mu.Unlock() - log.Info("VStreamer engine: deleting streamers") delete(vse.streamers, idx) - log.Info("VStreamer engine: decrementing wait group") vse.wg.Done() - log.Info("VStreamer engine: unlocking") }() // No lock is held while streaming, but wg is incremented.