From 16eff6de8c2b6ff867492cfe29dd1744b5a40cf6 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 8 Aug 2018 14:41:54 -0500 Subject: [PATCH] prevent anti entropy and cluster resize from running simultaneously --- cluster.go | 33 +++++++++++++++++++++++++++++++++ holder.go | 5 ++++- server.go | 22 +++++++++++++++++++++- 3 files changed, 58 insertions(+), 2 deletions(-) diff --git a/cluster.go b/cluster.go index 309846e0b..689706d75 100644 --- a/cluster.go +++ b/cluster.go @@ -204,6 +204,8 @@ type cluster struct { // nolint: maligned joining chan struct{} joined bool + abortAntiEntropyCh chan struct{} + mu sync.RWMutex jobs map[int64]*resizeJob currentJob *resizeJob @@ -235,6 +237,33 @@ func newCluster() *cluster { } } +// initializeAntiEntropy is called by the anti entropy routine when it starts. +// If the AE channel is created without a routine reading from it, cluster will +// block indefinitely when calling abortAntiEntropy(). +func (c *cluster) initializeAntiEntropy() { + c.mu.Lock() + c.abortAntiEntropyCh = make(chan struct{}) + c.mu.Unlock() +} + +// abortAntiEntropyQ checks whether the cluster wants to abort the anti entropy +// process (so that it can resize). It does not block. +func (c *cluster) abortAntiEntropyQ() bool { + select { + case <-c.abortAntiEntropyCh: + return true + default: + return false + } +} + +// abortAntiEntropy blocks until the anti-entropy routine calls abortAntiEntropyQ +func (c *cluster) abortAntiEntropy() { + if c.abortAntiEntropyCh != nil { + c.abortAntiEntropyCh <- struct{}{} + } +} + func (c *cluster) coordinatorNode() *Node { c.mu.RLock() defer c.mu.RUnlock() @@ -405,6 +434,10 @@ func (c *cluster) unprotectedSetState(state string) { c.state = state + if state == ClusterStateResizing { + c.abortAntiEntropy() + } + // TODO: consider NOT running cleanup on an active node that has // been removed. // It's safe to do a cleanup after state changes back to normal. diff --git a/holder.go b/holder.go index c2e8f5a6e..cf895e802 100644 --- a/holder.go +++ b/holder.go @@ -577,8 +577,11 @@ type holderSyncer struct { Closing <-chan struct{} } -// IsClosing returns true if the syncer has been marked to close. +// IsClosing returns true if the syncer has been asked to close. func (s *holderSyncer) IsClosing() bool { + if s.Cluster.abortAntiEntropyQ() { + return true + } select { case <-s.Closing: return true diff --git a/server.go b/server.go index 2cc4accdc..d120fcc09 100644 --- a/server.go +++ b/server.go @@ -417,6 +417,8 @@ func (s *Server) monitorAntiEntropy() { if s.antiEntropyInterval == 0 { return // anti entropy disabled } + s.cluster.initializeAntiEntropy() + ticker := time.NewTicker(s.antiEntropyInterval) defer ticker.Stop() @@ -428,11 +430,17 @@ func (s *Server) monitorAntiEntropy() { select { case <-s.closing: return + case <-s.cluster.abortAntiEntropyCh: // receive here so we don't block resizing + continue case <-ticker.C: s.holder.Stats.Count("AntiEntropy", 1, 1.0) } t := time.Now() - + if s.cluster.State() == ClusterStateResizing { + continue // don't launch anti-entropy during resize. + // the cluster sets its state to resizing and *then* sends to + // abortAntiEntropyCh before starting to resize + } // Sync holders. s.logger.Printf("holder sync beginning") if err := s.syncer.SyncHolder(); err != nil { @@ -444,6 +452,18 @@ func (s *Server) monitorAntiEntropy() { s.logger.Printf("holder sync complete") dif := time.Since(t) s.holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0) + + // Drain tick channel since we just finished anti-entropy. If the AE + // process took a long time, we don't want them to pile up on each + // other. + for { + select { + case <-ticker.C: + continue + default: + } + break + } } }