From dee46d442392d48b1f8676d3a3195ab66f6de207 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Tue, 26 Apr 2022 11:01:48 -0500 Subject: [PATCH] add PartitionToNodeAssignment as a new option We default to the jmp-hash method which we had previously, and allow a user to set the "modulus" option which uses a simple mod operation to ensure an even spread of partitions across nodes. I think that ideally we would have new indexes uses modulus and existing indexes use jmp-hash which implies supporting this configuration on a per-index basis. If we don't do per index, we should probably run the whole test suite both ways. --- api.go | 30 +++++++++++++++--------------- cluster.go | 24 +++++++++++++++--------- cluster_internal_test.go | 4 ++-- ctl/server.go | 1 + executor.go | 10 +++++----- fragment.go | 6 +++--- holder.go | 4 ++-- server.go | 15 +++++++++++---- server/config.go | 9 ++++++++- server/server.go | 1 + topology/noder.go | 2 +- topology/snapshot.go | 21 ++++++++++++++------- 12 files changed, 78 insertions(+), 49 deletions(-) diff --git a/api.go b/api.go index f172b853d..d2fcbbc83 100644 --- a/api.go +++ b/api.go @@ -286,7 +286,7 @@ func (api *API) DeleteIndex(ctx context.Context, indexName string) error { return errors.Wrap(err, "sending DeleteIndex message") } // Delete ids allocated for index if any present - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() if snap.IsPrimaryFieldTranslationNode(api.NodeID()) { if err := api.holder.ida.reset(indexName); err != nil { return errors.Wrap(err, "deleting id allocation for index") @@ -529,7 +529,7 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, defer qcx.Abort() // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() nodes := snap.ShardNodes(indexName, shard) errCh := make(chan error, len(nodes)) @@ -654,7 +654,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() // Validate that this handler owns the shard. if !snap.OwnsShard(api.NodeID(), indexName, shard) { @@ -740,7 +740,7 @@ func (api *API) ShardNodes(ctx context.Context, indexName string, shard uint64) } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() return snap.ShardNodes(indexName, shard), nil } @@ -755,7 +755,7 @@ func (api *API) PartitionNodes(ctx context.Context, partitionID int) ([]*topolog } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() return snap.PartitionNodes(partitionID), nil } @@ -862,7 +862,7 @@ func (api *API) TranslateData(ctx context.Context, indexName string, partition i } // Find the node that can service the request. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() nodes := snap.PartitionNodes(partition) var upNode *topology.Node for _, node := range nodes { @@ -944,7 +944,7 @@ func (api *API) NodeID() string { // PrimaryNode returns the primary node for the cluster. func (api *API) PrimaryNode() *topology.Node { // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() return snap.PrimaryFieldTranslationNode() } @@ -1867,7 +1867,7 @@ func (api *API) IngestOperations(ctx context.Context, qcx *Qcx, indexName string return errors.Wrap(err, "sharding input data") } // now that we have this, let's assign the shards to nodes - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() // oh hey an easy case: we're presumably the only node if len(snap.Nodes) == 1 { return api.ingestNodeOperationsForFields(ctx, qcx, index, knownFields, sharded) @@ -2095,7 +2095,7 @@ func (api *API) LongQueryTime() time.Duration { func (api *API) validateShardOwnership(indexName string, shard uint64) error { // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() // Validate that this handler owns the shard. if !snap.OwnsShard(api.NodeID(), indexName, shard) { api.server.logger.Errorf("node %s does not own shard %d of index %s", api.NodeID(), shard, indexName) @@ -2389,7 +2389,7 @@ func (api *API) MatchField(ctx context.Context, index, field string, like string // PrimaryReplicaNodeURL returns the URL of the cluster's primary replica. func (api *API) PrimaryReplicaNodeURL() url.URL { // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() node := snap.PrimaryReplicaNode(api.NodeID()) if node == nil { @@ -2500,7 +2500,7 @@ func (api *API) ReserveIDs(key IDAllocKey, session [32]byte, offset uint64, coun } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() if !snap.IsPrimaryFieldTranslationNode(api.NodeID()) { return nil, errors.New("cannot reserve IDs on a non-primary node") @@ -2515,7 +2515,7 @@ func (api *API) CommitIDs(key IDAllocKey, session [32]byte, count uint64) error } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() if !snap.IsPrimaryFieldTranslationNode(api.NodeID()) { return errors.New("cannot commit IDs on a non-primary node") @@ -2530,7 +2530,7 @@ func (api *API) ResetIDAlloc(index string) error { } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() if !snap.IsPrimaryFieldTranslationNode(api.NodeID()) { return errors.New("cannot reset IDs on a non-primary node") @@ -2588,7 +2588,7 @@ func (api *API) TranslateFieldDB(ctx context.Context, indexName, fieldName strin // RestoreShard func (api *API) RestoreShard(ctx context.Context, indexName string, shard uint64, rd io.Reader) error { - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() if !snap.OwnsShard(api.server.nodeID, indexName, shard) { return ErrClusterDoesNotOwnShard // TODO (twg)really just node doesn't own shard but leave for now } @@ -2785,7 +2785,7 @@ func (api *API) MutexCheck(ctx context.Context, qcx *Qcx, indexName string, fiel return nil, errors.New("can only check mutex state for mutex fields") } // request data from other nodes as well - snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + snap := api.cluster.NewSnapshot() eg, _ := errgroup.WithContext(ctx) myID := api.NodeID() results := make([]map[uint64]map[uint64][]uint64, len(snap.Nodes)) diff --git a/cluster.go b/cluster.go index 61c641d13..0fab2875c 100644 --- a/cluster.go +++ b/cluster.go @@ -106,6 +106,8 @@ type cluster struct { // nolint: maligned confirmDownRetries int confirmDownSleep time.Duration + + partitionAssigner string } // newCluster returns a new instance of Cluster with defaults. @@ -168,7 +170,7 @@ func (c *cluster) primaryNode() *topology.Node { // unprotectedPrimaryNode returns the primary node. func (c *cluster) unprotectedPrimaryNode() *topology.Node { // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() return snap.PrimaryFieldTranslationNode() } @@ -759,7 +761,7 @@ func (c *cluster) fragsByHost(idx *Index) fragsByHost { // for the given set of shards with data. func (c *cluster) fragCombos(idx string, availableShards *roaring.Bitmap, fieldViews viewsByField) fragsByHost { // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() t := make(fragsByHost) _ = availableShards.ForEach(func(i uint64) error { @@ -926,8 +928,8 @@ func (c *cluster) translationNodes(to *cluster) (map[string][]*translationResize } // Create a snapshot of the cluster to use for node/partition calculations. - fSnap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) - toSnap := topology.NewClusterSnapshot(to.noder, c.Hasher, to.ReplicaN) + fSnap := c.NewSnapshot() + toSnap := topology.NewClusterSnapshot(to.noder, c.Hasher, c.partitionAssigner, to.ReplicaN) for pid := 0; pid < c.partitionN; pid++ { fNodes := fSnap.PartitionNodes(pid) @@ -990,7 +992,7 @@ func (c *cluster) shardDistributionByIndex(indexName string) map[string]map[stri defer c.mu.RUnlock() // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() for _, shard := range available { p := snap.ShardToShardPartition(indexName, shard) @@ -1504,7 +1506,7 @@ func (c *cluster) translateFieldIDs(ctx context.Context, field *Field, ids map[u func (c *cluster) translateFieldListIDs(ctx context.Context, field *Field, ids []uint64) (keys []string, err error) { // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() primary := snap.PrimaryFieldTranslationNode() if primary == nil { @@ -1628,7 +1630,7 @@ func (c *cluster) findIndexKeys(ctx context.Context, indexName string, keys ...s } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() // Split keys by partition. keysByPartition := make(map[int][]string, c.partitionN) @@ -1737,7 +1739,7 @@ func (c *cluster) createIndexKeys(ctx context.Context, indexName string, keys .. } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() // Split keys by partition. keysByPartition := make(map[int][]string, c.partitionN) @@ -1858,7 +1860,7 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() // Split ids by partition. idsByPartition := make(map[int][]uint64, c.partitionN) @@ -1907,6 +1909,10 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS return idMap, nil } +func (c *cluster) NewSnapshot() *topology.ClusterSnapshot { + return topology.NewClusterSnapshot(c.noder, c.Hasher, c.partitionAssigner, c.ReplicaN) +} + // ClusterStatus describes the status of the cluster including its // state and node topology. type ClusterStatus struct { diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 9e123a96f..6b992644c 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -369,7 +369,7 @@ func TestCluster_Owners(t *testing.T) { cNodes := c.noder.Nodes() // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() // Verify nodes are distributed. if a := snap.PartitionNodes(0); !reflect.DeepEqual(a, []*topology.Node{cNodes[0], cNodes[1]}) { @@ -433,7 +433,7 @@ func TestCluster_ContainsShards(t *testing.T) { cNodes := c.noder.Nodes() // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + snap := c.NewSnapshot() shards := snap.ContainsShards("test", roaring.NewBitmap(0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10), cNodes[2]) diff --git a/ctl/server.go b/ctl/server.go index ff2f06862..ca84646e9 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -37,6 +37,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.IntVar(&srv.Config.Cluster.ReplicaN, "cluster.replicas", 1, "Number of hosts each piece of data should be stored on.") flags.DurationVar((*time.Duration)(&srv.Config.Cluster.LongQueryTime), "cluster.long-query-time", time.Duration(srv.Config.Cluster.LongQueryTime), "RENAMED TO 'long-query-time': Duration that will trigger log and stat messages for slow queries.") // negative duration indicates invalid value because 0 is meaningful flags.StringVar(&srv.Config.Cluster.Name, "cluster.name", srv.Config.Cluster.Name, "Human-readable name for the cluster.") + flags.StringVar(&srv.Config.Cluster.PartitionToNodeAssignment, "cluster.partition-to-node-assignment", srv.Config.Cluster.PartitionToNodeAssignment, "How to assign partitions to nodes. jmp-hash or modulus") // Translation flags.StringVar(&srv.Config.Translation.PrimaryURL, "translation.primary-url", srv.Config.Translation.PrimaryURL, "DEPRECATED: URL for primary translation node for replication.") diff --git a/executor.go b/executor.go index bde4830c1..54534079c 100644 --- a/executor.go +++ b/executor.go @@ -5178,7 +5178,7 @@ func (e *executor) executeClearBitField(ctx context.Context, qcx *Qcx, index str shard := colID / ShardWidth // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + snap := e.Cluster.NewSnapshot() ret := false for _, node := range snap.ShardNodes(index, shard) { @@ -5539,7 +5539,7 @@ func (e *executor) executeSetBitField(ctx context.Context, qcx *Qcx, index strin ret := false // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + snap := e.Cluster.NewSnapshot() for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. @@ -5585,7 +5585,7 @@ func (e *executor) executeSetValueField(ctx context.Context, qcx *Qcx, index str ret := false // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + snap := e.Cluster.NewSnapshot() for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. @@ -5632,7 +5632,7 @@ func (e *executor) executeClearValueField(ctx context.Context, qcx *Qcx, index s ret := false // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + snap := e.Cluster.NewSnapshot() for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. @@ -5699,7 +5699,7 @@ func (e *executor) shardsByNode(nodes []*topology.Node, index string, shards []u // We use e.Cluster.Nodes() here instead of e.Cluster.noder because we need // the node states in order to ensure that we don't include an unavailable // node in the map of nodes to which we distribute the query. - snap := topology.NewClusterSnapshot(topology.NewLocalNoder(e.Cluster.Nodes()), e.Cluster.Hasher, e.Cluster.ReplicaN) + snap := topology.NewClusterSnapshot(topology.NewLocalNoder(e.Cluster.Nodes()), e.Cluster.Hasher, e.Cluster.partitionAssigner, e.Cluster.ReplicaN) loop: for _, shard := range shards { diff --git a/fragment.go b/fragment.go index cff68282e..3097e2f28 100644 --- a/fragment.go +++ b/fragment.go @@ -3065,7 +3065,7 @@ func (s *fragmentSyncer) syncFragment() error { defer span.Finish() // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + snap := s.Cluster.NewSnapshot() // Determine replica set. nodes := snap.ShardNodes(s.Fragment.index(), s.Fragment.shard) @@ -3185,7 +3185,7 @@ func (s *fragmentSyncer) syncBlockFromPrimary(id int) error { f := s.Fragment // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + snap := s.Cluster.NewSnapshot() // Determine replica set. Return early if this is not // the primary node. @@ -3237,7 +3237,7 @@ func (s *fragmentSyncer) syncBlock(id int) error { f := s.Fragment // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + snap := s.Cluster.NewSnapshot() // Read pairs from each remote block. var uris []*pnet.URI diff --git a/holder.go b/holder.go index d1d17b7c3..c8100bc17 100644 --- a/holder.go +++ b/holder.go @@ -1235,7 +1235,7 @@ func (s *holderSyncer) SyncHolder() error { ti := time.Now() // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + snap := s.Cluster.NewSnapshot() schema, err := s.Holder.Schema() if err != nil { @@ -1351,7 +1351,7 @@ func (s *holderSyncer) resetTranslationSync() error { } // Create a snapshot of the cluster to use for node/partition calculations. - snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + snap := s.Cluster.NewSnapshot() // Set read-only flag for all translation stores. s.setTranslateReadOnlyFlags(snap) diff --git a/server.go b/server.go index baa3fb463..bbce8c55e 100644 --- a/server.go +++ b/server.go @@ -427,6 +427,13 @@ func OptServerLookupDB(dsn string) ServerOption { } } +func OptServerPartitionAssigner(p string) ServerOption { + return func(s *Server) error { + s.cluster.partitionAssigner = p + return nil + } +} + // NewServer returns a new instance of Server. func NewServer(opts ...ServerOption) (*Server, error) { cluster := newCluster() @@ -1295,7 +1302,7 @@ func (s *Server) monitorRuntime() { } func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, remote bool) (*Transaction, error) { - snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN) + snap := srv.cluster.NewSnapshot() node := srv.node() if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 { return nil, ErrNodeNotPrimary @@ -1342,7 +1349,7 @@ func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time } func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool) (*Transaction, error) { - snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN) + snap := srv.cluster.NewSnapshot() node := srv.node() if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 { return nil, ErrNodeNotPrimary @@ -1372,7 +1379,7 @@ func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool } func (srv *Server) Transactions(ctx context.Context) (map[string]*Transaction, error) { - snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN) + snap := srv.cluster.NewSnapshot() node := srv.node() if !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 { return nil, ErrNodeNotPrimary @@ -1382,7 +1389,7 @@ func (srv *Server) Transactions(ctx context.Context) (map[string]*Transaction, e } func (srv *Server) GetTransaction(ctx context.Context, id string, remote bool) (*Transaction, error) { - snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN) + snap := srv.cluster.NewSnapshot() node := srv.node() if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 { diff --git a/server/config.go b/server/config.go index 15fd7df0f..5b1b5b239 100644 --- a/server/config.go +++ b/server/config.go @@ -129,7 +129,8 @@ type Config struct { ReplicaN int `toml:"replicas"` Name string `toml:"name"` // This LongQueryTime is deprecated but still exists for backward compatibility - LongQueryTime toml.Duration `toml:"long-query-time"` + LongQueryTime toml.Duration `toml:"long-query-time"` + PartitionToNodeAssignment string `toml:"partition-to-node-assignment"` } `toml:"cluster"` // Etcd config is based on embedded etcd. @@ -311,6 +312,11 @@ func (c *Config) validate() error { return nil } +const ( + PartitionToNodeJmp string = "jmp-hash" + PartitionToNodeModulus string = "modulus" +) + // NewConfig returns an instance of Config with default options. func NewConfig() *Config { c := &Config{ @@ -345,6 +351,7 @@ func NewConfig() *Config { c.Cluster.Name = "cluster0" c.Cluster.ReplicaN = 1 c.Cluster.LongQueryTime = toml.Duration(-time.Minute) //TODO remove this once cluster.longQueryTime is fully deprecated + c.Cluster.PartitionToNodeAssignment = PartitionToNodeJmp // AntiEntropy config. c.AntiEntropy.Interval = toml.Duration(0) diff --git a/server/server.go b/server/server.go index 5e368cfab..bedf68f33 100644 --- a/server/server.go +++ b/server/server.go @@ -485,6 +485,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerRBFConfig(m.Config.RBFConfig), pilosa.OptServerMaxQueryMemory(m.Config.MaxQueryMemory), pilosa.OptServerQueryHistoryLength(m.Config.QueryHistoryLength), + pilosa.OptServerPartitionAssigner(m.Config.Cluster.PartitionToNodeAssignment), discoOpt, } diff --git a/topology/noder.go b/topology/noder.go index 63c8c5e69..9e9ab3942 100644 --- a/topology/noder.go +++ b/topology/noder.go @@ -60,7 +60,7 @@ func (n *localNoder) Nodes() []*Node { // PrimaryNodeID implements the Noder interface. func (n *localNoder) PrimaryNodeID(hasher Hasher) string { - snap := NewClusterSnapshot(NewLocalNoder(n.nodes), hasher, 1) + snap := NewClusterSnapshot(NewLocalNoder(n.nodes), hasher, "jmp-hash", 1) primaryNode := snap.PrimaryFieldTranslationNode() if primaryNode == nil { return "" diff --git a/topology/snapshot.go b/topology/snapshot.go index be394aac8..313a11882 100644 --- a/topology/snapshot.go +++ b/topology/snapshot.go @@ -31,10 +31,12 @@ type ClusterSnapshot struct { // The number of replicas a partition has. ReplicaN int + + PartitionAssignment string } // NewClusterSnapshot returns a new instance of ClusterSnapshot. -func NewClusterSnapshot(noder Noder, hasher Hasher, replicas int) *ClusterSnapshot { +func NewClusterSnapshot(noder Noder, hasher Hasher, partitionAssignment string, replicas int) *ClusterSnapshot { nodes := noder.Nodes() // Make sure replica count doesn't exceed the number of nodes. @@ -46,10 +48,11 @@ func NewClusterSnapshot(noder Noder, hasher Hasher, replicas int) *ClusterSnapsh } return &ClusterSnapshot{ - Nodes: nodes, - Hasher: hasher, - PartitionN: DefaultPartitionN, - ReplicaN: replicas, + Nodes: nodes, + Hasher: hasher, + PartitionN: DefaultPartitionN, + ReplicaN: replicas, + PartitionAssignment: partitionAssignment, } } @@ -160,7 +163,11 @@ func (c *ClusterSnapshot) IsPrimary(nodeID string, partition int) bool { // PrimaryNodeIndex returns the index (position in the cluster) of the primary // node for the given partition. func (c *ClusterSnapshot) PrimaryNodeIndex(partition int) int { - return partition % len(c.Nodes) + if c.PartitionAssignment == "modulus" { + return partition % len(c.Nodes) + } else { + return c.Hasher.Hash(uint64(partition), len(c.Nodes)) + } } // NonPrimaryReplicas returns the list of node IDs which are replicas for the @@ -277,7 +284,7 @@ func NodePositionByID(nodes []*Node, nodeID string) int { // and a hasher. The order of the node IDs provided does not matter because this // function will re-order them in a deterministic way. func PrimaryNodeID(nodeIDs []string, hasher Hasher) string { - snap := NewClusterSnapshot(NewIDNoder(nodeIDs), hasher, 1) + snap := NewClusterSnapshot(NewIDNoder(nodeIDs), hasher, "jmp-hash", 1) primaryNode := snap.PrimaryFieldTranslationNode() if primaryNode == nil { return ""