From 23f901635e0d82531eba5d278e47fbb52cdea426 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 12:49:14 +0100 Subject: [PATCH] Revert "Remove etcd cache" This reverts commit 0f4b273d3a2165dc6a90b7fe3374a2c053b79d9e. --- etcd/cache.go | 94 ++++++++++++++++++++++++++++++++++++++++++++++++ etcd/embed.go | 25 +++++++------ server/server.go | 2 +- 3 files changed, 109 insertions(+), 12 deletions(-) create mode 100644 etcd/cache.go diff --git a/etcd/cache.go b/etcd/cache.go new file mode 100644 index 000000000..8005c87ef --- /dev/null +++ b/etcd/cache.go @@ -0,0 +1,94 @@ +// Copyright 2017 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package etcd + +import ( + "context" + "sync" + "time" + + "github.com/pilosa/pilosa/v2/topology" +) + +// EtcdWithCache is a wrapper around the Etcd type which will return a +// cached value when the number of requests come in below a configured +// frequency. It also breaks the cache after a configured TTL. +type EtcdWithCache struct { + *Etcd + + peerMetadataMu sync.RWMutex + peerMetadata map[string][]byte + + peersMu sync.Mutex // peer-list cache updates + + nodes []*topology.Node // unmarshalled Node data + nodesTTL int // seconds + nodesLastRequest time.Time // last time requested +} + +// NewEtcdWithCache returns a new instance of Cache. +func NewEtcdWithCache(opt Options, replicas int) *EtcdWithCache { + return &EtcdWithCache{ + Etcd: NewEtcd(opt, replicas), + + nodesTTL: 6, + + peerMetadata: make(map[string][]byte), + } +} + +// Metadata is a cache wrapper around the Metadator.Metadata method. +func (c *EtcdWithCache) Metadata(ctx context.Context, peerID string) ([]byte, error) { + c.peerMetadataMu.RLock() + v, ok := c.peerMetadata[peerID] + c.peerMetadataMu.RUnlock() + if ok { + return v, nil + } + v, err := c.Etcd.Metadata(ctx, peerID) + if err == nil { + c.peerMetadataMu.Lock() + c.peerMetadata[peerID] = v + c.peerMetadataMu.Unlock() + } + return v, err +} + +// Nodes caches the result of the underlying implementation's node list. +func (c *EtcdWithCache) Nodes() []*topology.Node { + c.peersMu.Lock() + defer c.peersMu.Unlock() + + now := time.Now() + if now.Sub(c.nodesLastRequest) > (time.Duration(c.nodesTTL) * time.Second) { + c.nodes = c.Etcd.Nodes() + c.nodesLastRequest = now + } + return c.nodes +} + +// SetNodes implements the Noder interface as NOP +// (because we can't force to set nodes for etcd). +func (c *EtcdWithCache) SetNodes(nodes []*topology.Node) {} + +// AppendNode implements the Noder interface as NOP +// (because resizer is responsible for adding new nodes). +func (c *EtcdWithCache) AppendNode(node *topology.Node) {} + +// RemoveNode implements the Noder interface as NOP +// (because resizer is responsible for removing existing nodes) +func (c *EtcdWithCache) RemoveNode(nodeID string) bool { + return false +} diff --git a/etcd/embed.go b/etcd/embed.go index cacac1cf7..53a96c9a0 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -37,6 +37,7 @@ import ( "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" "go.etcd.io/etcd/etcdserver/api/v3client" + "go.etcd.io/etcd/mvcc" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -106,12 +107,6 @@ func NewEtcd(opt Options, replicas int) *Etcd { func (e *Etcd) Close() error { _ = testhook.Closed(pilosa.NewAuditor(), e, nil) - if e.cli != nil { - if err := e.cli.Close(); err != nil { - log.Printf("Error closing etcd client: %v", err) - } - } - if e.e != nil { if e.resizeCancel != nil { e.resizeCancel() @@ -123,6 +118,12 @@ func (e *Etcd) Close() error { <-e.e.Server.StopNotify() } + if e.cli != nil { + if err := e.cli.Close(); err != nil { + log.Printf("Error closing etcd client: %v", err) + } + } + return nil } @@ -239,7 +240,9 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e } func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) + kv := e.e.Server.KV() + + resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) if err != nil { return disco.NodeStateUnknown, err } @@ -247,20 +250,20 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateResizing, nil } - resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) + resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return disco.NodeStateUnknown, err } - if len(resp.Kvs) > 1 { + if len(resp.KVs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(resp.Kvs) == 0 { + if len(resp.KVs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } - return disco.NodeState(resp.Kvs[0].Value), nil + return disco.NodeState(resp.KVs[0].Value), nil } func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { diff --git a/server/server.go b/server/server.go index 1317a8ebb..c07ba17a4 100644 --- a/server/server.go +++ b/server/server.go @@ -390,7 +390,7 @@ func (m *Command) SetupServer() error { m.Config.Etcd.Dir = filepath.Join(path, pilosa.DefaultDiscoDir) } - e := petcd.NewEtcd(m.Config.Etcd, m.Config.Cluster.ReplicaN) + e := petcd.NewEtcdWithCache(m.Config.Etcd, m.Config.Cluster.ReplicaN) discoOpt := pilosa.OptServerDisCo(e, e, e, e, e, e, e) serverOptions := []pilosa.ServerOption{