From b304de65360bdef30d2548a66e6f1a7699a3088a Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 6 Aug 2018 15:57:22 -0500 Subject: [PATCH] Treat coordinator as primary translate store. Daisy-chain other nodes based on their position in the cluster. Deprecate the `primary-url` configuration option. --- api.go | 6 +-- cluster.go | 27 ++++++++-- ctl/server.go | 2 +- holder.go | 22 ++++++++ http/translator.go | 15 ++++-- server.go | 31 +++++------ server/config.go | 2 +- server/server.go | 7 ++- translate.go | 131 +++++++++++++++++++++++++++++++++++++++++++-- translate_test.go | 4 +- 10 files changed, 206 insertions(+), 41 deletions(-) diff --git a/api.go b/api.go index 847cc7ebd..fa28e0cb9 100644 --- a/api.go +++ b/api.go @@ -136,9 +136,9 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er } // Translate column attributes, if necessary. - if api.server.translateFile != nil { + if api.server.holder.translateFile != nil { for _, col := range resp.ColumnAttrSets { - v, err := api.server.translateFile.TranslateColumnToString(req.Index, col.ID) + v, err := api.server.holder.translateFile.TranslateColumnToString(req.Index, col.ID) if err != nil { return resp, err } @@ -764,7 +764,7 @@ func (api *API) ResizeAbort() error { // GetTranslateData provides a reader for key translation logs starting at offset. func (api *API) GetTranslateData(ctx context.Context, offset int64) (io.ReadCloser, error) { - rc, err := api.server.translateFile.Reader(ctx, offset) + rc, err := api.server.holder.translateFile.Reader(ctx, offset) if err != nil { return nil, errors.Wrap(err, "read from translate store") } diff --git a/cluster.go b/cluster.go index 715176b1c..32095fa24 100644 --- a/cluster.go +++ b/cluster.go @@ -169,7 +169,7 @@ type nodeAction struct { type cluster struct { // nolint: maligned id string Node *Node - Nodes []*Node // TODO phase this out? + Nodes []*Node // Hashing algorithm used to assign partitions to nodes. Hasher Hasher @@ -935,7 +935,6 @@ func (c *cluster) allNodesReady() (ret bool) { } func (c *cluster) handleNodeAction(nodeAction nodeAction) error { - c.mu.Lock() j, err := c.unprotectedGenerateResizeJob(nodeAction) c.mu.Unlock() @@ -959,7 +958,7 @@ func (c *cluster) handleNodeAction(nodeAction nodeAction) error { c.logger.Printf("wait for jobResult") jobResult := <-j.result - // Make sure j.Run() didn't return an error. + // Make sure j.run() didn't return an error. if eg.Wait() != nil { return errors.Wrap(err, "running job") } @@ -1779,6 +1778,10 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error { } } + // If the cluster membership has changed, reset the primary for + // translate store replication. + c.holder.setPrimaryTranslateStore(c.unprotectedPreviousNode()) + c.unprotectedSetState(cs.State) c.markAsJoined() @@ -1786,6 +1789,24 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error { return nil } +// unprotectedPreviousNode returns the node listed before the current node in c.Nodes. +// If there is only one node in the cluster, returns nil. +// If the current node is the first node in the list, returns the last node. +func (c *cluster) unprotectedPreviousNode() *Node { + if len(c.Nodes) <= 1 { + return nil + } + + pos := c.nodePositionByID(c.Node.ID) + if pos == -1 { + return nil + } else if pos == 0 { + return c.Nodes[len(c.Nodes)-1] + } else { + return c.Nodes[pos-1] + } +} + // setStatic is unprotected, but only called before the cluster has been started // (and therefore not concurrently). func (c *cluster) setStatic(hosts []string) error { diff --git a/ctl/server.go b/ctl/server.go index ad1ac5df5..90fdfb46e 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -44,7 +44,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.DurationVarP((*time.Duration)(&srv.Config.Cluster.LongQueryTime), "cluster.long-query-time", "", time.Minute, "Duration that will trigger log and stat messages for slow queries.") // Translation - flags.StringVarP(&srv.Config.Translation.PrimaryURL, "translation.primary-url", "", srv.Config.Translation.PrimaryURL, "URL for primary translation node for replication.") + flags.StringVarP(&srv.Config.Translation.PrimaryURL, "translation.primary-url", "", srv.Config.Translation.PrimaryURL, "DEPRECATED: URL for primary translation node for replication.") // Gossip flags.StringVarP(&srv.Config.Gossip.Port, "gossip.port", "", srv.Config.Gossip.Port, "Port to which pilosa should bind for internal state sharing.") diff --git a/holder.go b/holder.go index c2e8f5a6e..89b4dc368 100644 --- a/holder.go +++ b/holder.go @@ -46,6 +46,10 @@ type Holder struct { // Indexes by name. indexes map[string]*Index + // Key/ID translation + translateFile *TranslateFile + NewPrimaryTranslateStore func(*Node) TranslateStore + // opened channel is closed once Open() completes. opened chan struct{} @@ -77,6 +81,9 @@ func NewHolder() *Holder { opened: make(chan struct{}), + translateFile: NewTranslateFile(), + NewPrimaryTranslateStore: newNopTranslateStore, + broadcaster: NopBroadcaster, Stats: NopStatsClient, @@ -160,6 +167,13 @@ func (h *Holder) Close() error { return errors.Wrap(err, "closing index") } } + + if h.translateFile != nil { + if err := h.translateFile.Close(); err != nil { + return err + } + } + return nil } @@ -560,6 +574,14 @@ func (h *Holder) logStartup() error { return nil } +func (h *Holder) setPrimaryTranslateStore(node *Node) { + var nodeID string + if node != nil { + nodeID = node.ID + } + h.translateFile.SetPrimaryStore(nodeID, h.NewPrimaryTranslateStore(node)) +} + // holderSyncer is an active anti-entropy tool that compares the local holder // with a remote holder based on block checksums and resolves differences. type holderSyncer struct { diff --git a/http/translator.go b/http/translator.go index 80c73bf5e..715493c6b 100644 --- a/http/translator.go +++ b/http/translator.go @@ -19,12 +19,17 @@ var _ pilosa.TranslateStore = (*translateStore)(nil) // translateStore represents an implementation of pilosa.TranslateStore that // communicates over HTTP. This is used with the TranslateHandler. type translateStore struct { - URL string + node *pilosa.Node } -// NewTranslateStore returns a new instance of TranslateStore. -func NewTranslateStore(rawurl string) *translateStore { - return &translateStore{URL: rawurl} +// DEPRECATED: NewTranslateStore returns a new instance of translateStore. +func NewTranslateStore(rawurl string) *translateStore { // nolint: unparam + return &translateStore{node: nil} +} + +// NewNodeTranslateStore returns a new instance of TranslateStore based on node. +func NewNodeTranslateStore(node *pilosa.Node) pilosa.TranslateStore { + return &translateStore{node: node} } // TranslateColumnsToUint64 is not currently implemented. @@ -50,7 +55,7 @@ func (s *translateStore) TranslateRowToString(index, frame string, values uint64 // Reader returns a reader that can stream data from a remote store. func (s *translateStore) Reader(ctx context.Context, off int64) (io.ReadCloser, error) { // Generate remote URL. - u, err := url.Parse(s.URL) + u, err := url.Parse(s.node.URI.String()) if err != nil { return nil, err } diff --git a/server.go b/server.go index 2cc4accdc..b4a704b80 100644 --- a/server.go +++ b/server.go @@ -49,7 +49,6 @@ type Server struct { // nolint: maligned // Internal holder *Holder cluster *cluster - translateFile *TranslateFile diagnostics *diagnosticsCollector executor *executor hosts []string @@ -70,8 +69,6 @@ type Server struct { // nolint: maligned isCoordinator bool syncer holderSyncer - primaryTranslateStore TranslateStore - defaultClient InternalClient dataDir string } @@ -163,9 +160,17 @@ func OptServerInternalClient(c InternalClient) ServerOption { } } +// DEPRECATED func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption { return func(s *Server) error { - s.primaryTranslateStore = store + s.logger.Printf("DEPRECATED: OptServerPrimaryTranslateStore") + return nil + } +} + +func OptServerPrimaryTranslateStoreFunc(tf func(*Node) TranslateStore) ServerOption { + return func(s *Server) error { + s.holder.NewPrimaryTranslateStore = tf return nil } } @@ -261,6 +266,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { } s.holder.Path = path + s.holder.translateFile.Path = filepath.Join(path, ".keys") s.holder.Logger = s.logger s.holder.Stats.SetLogger(s.logger) @@ -268,11 +274,6 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.cluster.logger = s.logger s.cluster.holder = s.holder - // Initialize translation database. - s.translateFile = NewTranslateFile() - s.translateFile.Path = filepath.Join(path, ".keys") - s.translateFile.PrimaryTranslateStore = s.primaryTranslateStore - // Get or create NodeID. s.nodeID = s.loadNodeID() if s.isCoordinator { @@ -299,7 +300,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.executor.Holder = s.holder s.executor.Node = node s.executor.Cluster = s.cluster - s.executor.TranslateStore = s.translateFile + s.executor.TranslateStore = s.holder.translateFile s.executor.MaxWritesPerRequest = s.maxWritesPerRequest s.cluster.broadcaster = s s.cluster.maxWritesPerRequest = s.maxWritesPerRequest @@ -324,7 +325,7 @@ func (s *Server) Open() error { } // Initialize id-key storage. - if err := s.translateFile.Open(); err != nil { + if err := s.holder.translateFile.Open(); err != nil { return err } @@ -370,7 +371,6 @@ func (s *Server) Close() error { s.wg.Wait() var errh error - var errt error var errc error if s.cluster != nil { errc = s.cluster.close() @@ -378,17 +378,12 @@ func (s *Server) Close() error { if s.holder != nil { errh = s.holder.Close() } - if s.translateFile != nil { - errt = s.translateFile.Close() - } - // prefer to return holder error over translateFile error over cluster + // prefer to return holder error over cluster // error. This order is somewhat arbitrary. It would be better if we had // some way to combine all the errors, but probably not important enough to // warrant the extra complexity. if errh != nil { return errors.Wrap(errh, "closing holder") - } else if errt != nil { - return errors.Wrap(errt, "closing translateFile") } return errors.Wrap(errc, "closing cluster") } diff --git a/server/config.go b/server/config.go index 39a769ce2..8367eb835 100644 --- a/server/config.go +++ b/server/config.go @@ -71,7 +71,7 @@ type Config struct { // Gossip config is based around memberlist.Config. Gossip gossip.Config `toml:"gossip"` - // Translation config supports translation store replication. + // DEPRECATED: Translation config supports translation store replication. Translation struct { PrimaryURL string `toml:"primary-url"` } `toml:"translation"` diff --git a/server/server.go b/server/server.go index 1273a238f..d10fa3ac2 100644 --- a/server/server.go +++ b/server/server.go @@ -250,10 +250,9 @@ func (m *Command) SetupServer() error { c := http.GetHTTPClient(TLSConfig) - // Setup connection to primary store if this is a replica. - var primaryTranslateStore pilosa.TranslateStore + // Primary store configuration is handled automatically now. if m.Config.Translation.PrimaryURL != "" { - primaryTranslateStore = http.NewTranslateStore(m.Config.Translation.PrimaryURL) + m.logger.Printf("DEPRECATED: The primary-url configuration option is no longer used.") } // Set Coordinator. @@ -278,7 +277,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerStatsClient(statsClient), pilosa.OptServerURI(uri), pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)), - pilosa.OptServerPrimaryTranslateStore(primaryTranslateStore), + pilosa.OptServerPrimaryTranslateStoreFunc(http.NewNodeTranslateStore), pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts), pilosa.OptServerSerializer(proto.Serializer{}), coordinatorOpt, diff --git a/translate.go b/translate.go index 20696ef6c..a37eb5b74 100644 --- a/translate.go +++ b/translate.go @@ -8,6 +8,7 @@ import ( "errors" "fmt" "io" + "io/ioutil" "log" "os" "path/filepath" @@ -71,6 +72,10 @@ type TranslateFile struct { // If non-nil, data is streamed from a primary and this is a read-only store. PrimaryTranslateStore TranslateStore + primaryID string // unique ID used to identify the primary store + replicationClosing chan struct{} + primaryStoreEvents chan primaryStoreEvent + repWG sync.WaitGroup // Delay after attempting to connect to a primary that the store will retry. replicationRetryInterval time.Duration @@ -86,6 +91,9 @@ func NewTranslateFile() *TranslateFile { mapSize: defaultMapSize, + replicationClosing: make(chan struct{}), + primaryStoreEvents: make(chan primaryStoreEvent), + replicationRetryInterval: defaultReplicationRetryInterval, } } @@ -109,10 +117,65 @@ func (s *TranslateFile) Open() (err error) { return err } - // Stream from primary, if available. + // Listen to primaryStoreEvents channel. + s.wg.Add(1) + go func() { defer s.wg.Done(); s.monitorPrimaryStoreEvents() }() + + return nil +} + +// primaryStoreEvent is used to set/change the primary translate store. +// It contains a TranslateStore along with an associated string ID which +// is used to determine whether the primary needs to be changed from the +// current value. +type primaryStoreEvent struct { + id string + ts TranslateStore +} + +// SetPrimaryStore sets the translate files's primary translate store. +// The id value is used to determine whether the primary needs to be changed +// from the current value (i.e. calling this multiple times with the same +// input values will no-op on all subsequent calls). +func (s *TranslateFile) SetPrimaryStore(id string, ts TranslateStore) { + go func() { + s.primaryStoreEvents <- primaryStoreEvent{ + id: id, + ts: ts, + } + }() +} + +// handlePrimaryStoreEvent changes the PrimaryTranslateStore +// used for replication by TranslateFile. +func (s *TranslateFile) handlePrimaryStoreEvent(ev primaryStoreEvent) error { + s.mu.Lock() + defer s.mu.Unlock() + + if ev.id == s.primaryID { + return nil + } + + // Stop translate store replication. + log.Printf("stop monitor replication") + close(s.replicationClosing) + s.repWG.Wait() + + // Set the primary node for translate store replication. + log.Printf("set primary translate store to %s", ev.id) + s.primaryID = ev.id + if ev.id == "" { + s.PrimaryTranslateStore = nil + } else { + s.PrimaryTranslateStore = ev.ts + } + + // Start translate store replication. Stream from primary, if available. + log.Printf("start monitor replication") if s.PrimaryTranslateStore != nil { - s.wg.Add(1) - go func() { defer s.wg.Done(); s.monitorReplication() }() + s.replicationClosing = make(chan struct{}) + s.repWG.Add(1) + go func() { defer s.repWG.Done(); s.monitorReplication() }() } return nil @@ -259,7 +322,13 @@ func (s *TranslateFile) replayEntries() error { func (s *TranslateFile) monitorReplication() { // Create context that will cancel on close. ctx, cancel := context.WithCancel(context.Background()) - go func() { <-s.closing; cancel() }() + go func() { + select { + case <-s.closing: + case <-s.replicationClosing: + } + cancel() + }() // Keep attempting to replicate until the store closes. for { @@ -270,12 +339,31 @@ func (s *TranslateFile) monitorReplication() { select { case <-s.closing: return + case <-s.replicationClosing: + return case <-time.After(s.replicationRetryInterval): log.Printf("pilosa: reconnecting to primary replica") } } } +// monitorPrimaryStoreEvents is executed in a separate goroutine and listens for changes +// to the primary store assignment. +func (s *TranslateFile) monitorPrimaryStoreEvents() { + log.Printf("monitor primary store events") + // Keep handling events until the store closes. + for { + select { + case <-s.closing: + return + case ev := <-s.primaryStoreEvents: + if err := s.handlePrimaryStoreEvent(ev); err != nil { + log.Printf("handle primary store event") + } + } + } +} + func (s *TranslateFile) replicate(ctx context.Context) error { off := s.size() @@ -989,3 +1077,38 @@ func uVarintSize(x uint64) (i int) { } return i + 1 } + +// nopTStore represents a TranslateStore that doesn't do anything. +var nopTStore TranslateStore = nopTranslateStore{} + +// newNopTranslateStore returns a translate store which does nothing. It returns a global +// object to avoid unnecessary allocations. +func newNopTranslateStore(*Node) TranslateStore { return nopTStore } + +// nopTranslateStore represents a no-op implementation of the TranslateStore interface. +type nopTranslateStore struct{} + +// TranslateColumnsToUint64 is a no-op implementation of the TranslateStore TranslateColumnsToUint64 method. +func (s nopTranslateStore) TranslateColumnsToUint64(index string, values []string) ([]uint64, error) { + return []uint64{}, nil +} + +// TranslateColumnToString is a no-op implementation of the TranslateStore TranslateColumnToString method. +func (s nopTranslateStore) TranslateColumnToString(index string, values uint64) (string, error) { + return "", nil +} + +// TranslateRowsToUint64 is a no-op implementation of the TranslateStore TranslateRowsToUint64 method. +func (s nopTranslateStore) TranslateRowsToUint64(index, field string, values []string) ([]uint64, error) { + return []uint64{}, nil +} + +// TranslateRowToString is a no-op implementation of the TranslateStore TranslateRowToString method. +func (s nopTranslateStore) TranslateRowToString(index, field string, values uint64) (string, error) { + return "", nil +} + +// Reader is a no-op implementation of the TranslateStore Reader method. +func (s nopTranslateStore) Reader(ctx context.Context, off int64) (io.ReadCloser, error) { + return ioutil.NopCloser(bytes.NewReader(nil)), nil +} diff --git a/translate_test.go b/translate_test.go index 9856bc85c..4feb482f1 100644 --- a/translate_test.go +++ b/translate_test.go @@ -383,7 +383,7 @@ func TestTranslateFile_PrimaryTranslateStore(t *testing.T) { // Create a replica that accepts writes from primary. replica := NewTranslateFile() - replica.PrimaryTranslateStore = primary + replica.SetPrimaryStore("primary", primary) if err := replica.Open(); err != nil { t.Fatal(err) } @@ -567,7 +567,7 @@ func (s *TranslateFile) Reopen() error { s.TranslateFile = pilosa.NewTranslateFile() s.lock.Unlock() s.Path = prev.Path - s.PrimaryTranslateStore = prev.PrimaryTranslateStore + s.SetPrimaryStore("restored-primary", prev.PrimaryTranslateStore) return s.Open() }