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.
This commit is contained in:
Matthew Jaffee 2022-04-26 11:01:48 -05:00 • committed by Matthew Jaffee
parent 24af93e7d8
commit dee46d4423
12 changed files with 78 additions and 49 deletions

30
api.go
View file

@ -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))

View file

@ -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 {

View file

@ -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])

View file

@ -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.")

View file

@ -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 {

View file

@ -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

View file

@ -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)

View file

@ -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 {

View file

@ -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)

View file

@ -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,
}

View file

@ -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 ""

View file

@ -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 ""