mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 11:57:51 +00:00
finish implementing snap := ClusterSnapshot()
This commit is contained in:
parent
1473e11a27
commit
2f66501160
5 changed files with 138 additions and 63 deletions
26
api.go
26
api.go
|
|
@ -497,7 +497,10 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string,
|
|||
qcx := api.Txf().NewQcx()
|
||||
defer qcx.Abort()
|
||||
|
||||
nodes := api.cluster.shardNodes(indexName, shard)
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
|
||||
|
||||
nodes := snap.ShardNodes(indexName, shard)
|
||||
errCh := make(chan error, len(nodes))
|
||||
for _, node := range nodes {
|
||||
node := node
|
||||
|
|
@ -619,8 +622,11 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
|
|||
return errors.Wrap(err, "validating api method")
|
||||
}
|
||||
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
|
||||
|
||||
// Validate that this handler owns the shard.
|
||||
if !api.cluster.ownsShard(api.Node().ID, indexName, shard) {
|
||||
if !snap.OwnsShard(api.Node().ID, indexName, shard) {
|
||||
api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName)
|
||||
return ErrClusterDoesNotOwnShard
|
||||
}
|
||||
|
|
@ -668,7 +674,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
|
|||
}
|
||||
|
||||
if index.Keys() {
|
||||
if store := index.TranslateStore(api.cluster.idPartition(indexName, columnID)); store == nil {
|
||||
if store := index.TranslateStore(snap.IDToShardPartition(indexName, columnID)); store == nil {
|
||||
return errors.Wrap(err, "partition does not exist")
|
||||
} else if colStr, err = store.TranslateID(columnID); err != nil {
|
||||
return errors.Wrap(err, "translating column")
|
||||
|
|
@ -702,7 +708,10 @@ func (api *API) ShardNodes(ctx context.Context, indexName string, shard uint64)
|
|||
return nil, errors.Wrap(err, "validating api method")
|
||||
}
|
||||
|
||||
return api.cluster.shardNodes(indexName, shard), nil
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
|
||||
|
||||
return snap.ShardNodes(indexName, shard), nil
|
||||
}
|
||||
|
||||
// FragmentBlockData is an endpoint for internal usage. It is not guaranteed to
|
||||
|
|
@ -1683,8 +1692,10 @@ 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)
|
||||
// Validate that this handler owns the shard.
|
||||
if !api.cluster.ownsShard(api.Node().ID, indexName, shard) {
|
||||
if !snap.OwnsShard(api.Node().ID, indexName, shard) {
|
||||
api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName)
|
||||
return ErrClusterDoesNotOwnShard
|
||||
}
|
||||
|
|
@ -2003,7 +2014,10 @@ func (api *API) CreateFieldKeys(ctx context.Context, index, field string, keys .
|
|||
|
||||
// PrimaryReplicaNodeURL returns the URL of the cluster's primary replica.
|
||||
func (api *API) PrimaryReplicaNodeURL() url.URL {
|
||||
node := api.cluster.PrimaryReplicaNode()
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN)
|
||||
|
||||
node := snap.PrimaryReplicaNode(api.Node().ID)
|
||||
if node == nil {
|
||||
return url.URL{}
|
||||
}
|
||||
|
|
|
|||
79
cluster.go
79
cluster.go
|
|
@ -666,9 +666,12 @@ func (c *cluster) fragsByHost(idx *Index) fragsByHost {
|
|||
// by creating every combination of field/view specified in `fieldViews` up
|
||||
// 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)
|
||||
|
||||
t := make(fragsByHost)
|
||||
_ = availableShards.ForEach(func(i uint64) error {
|
||||
nodes := c.shardNodes(idx, i)
|
||||
nodes := snap.ShardNodes(idx, i)
|
||||
for _, n := range nodes {
|
||||
// for each field/view combination:
|
||||
for field, views := range fieldViews {
|
||||
|
|
@ -828,9 +831,13 @@ func (c *cluster) translationNodes(to *cluster) (map[string][]*translationResize
|
|||
m[n.ID] = nil
|
||||
}
|
||||
|
||||
// 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)
|
||||
|
||||
for pid := 0; pid < c.partitionN; pid++ {
|
||||
fNodes := c.partitionNodes(pid)
|
||||
tNodes := to.partitionNodes(pid)
|
||||
fNodes := fSnap.PartitionNodes(pid)
|
||||
tNodes := toSnap.PartitionNodes(pid)
|
||||
|
||||
// For `to` cluster, we include all nodes containing a
|
||||
// replica for the partition. The source for each replica
|
||||
|
|
@ -888,9 +895,12 @@ func (c *cluster) shardDistributionByIndex(indexName string) map[string]map[stri
|
|||
c.mu.RLock()
|
||||
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)
|
||||
|
||||
for _, shard := range available {
|
||||
p := c.shardToShardPartition(indexName, shard)
|
||||
nodes := c.partitionNodes(p)
|
||||
p := snap.ShardToShardPartition(indexName, shard)
|
||||
nodes := snap.PartitionNodes(p)
|
||||
dist[nodes[0].ID]["primary-shards"] = append(dist[nodes[0].ID]["primary-shards"], shard)
|
||||
for k := 1; k < len(nodes); k++ {
|
||||
dist[nodes[k].ID]["replica-shards"] = append(dist[nodes[k].ID]["replica-shards"], shard)
|
||||
|
|
@ -931,11 +941,6 @@ func keyToKeyPartition(index, key string, partitionN int) int {
|
|||
return int(h.Sum64() % uint64(partitionN))
|
||||
}
|
||||
|
||||
// idPartition returns the partition that an id belongs to.
|
||||
func (c *cluster) idPartition(index string, id uint64) int {
|
||||
return shardToShardPartition(index, id/ShardWidth, c.partitionN)
|
||||
}
|
||||
|
||||
// ShardNodes returns a list of nodes that own a fragment. Safe for concurrent use.
|
||||
func (c *cluster) ShardNodes(index string, shard uint64) []*topology.Node {
|
||||
c.mu.RLock()
|
||||
|
|
@ -960,13 +965,6 @@ func (c *cluster) keyNodes(index, key string) []*topology.Node {
|
|||
return c.partitionNodes(c.Topology.KeyPartition(index, key))
|
||||
}
|
||||
|
||||
// ownsShard returns true if a host owns a fragment.
|
||||
func (c *cluster) ownsShard(nodeID string, index string, shard uint64) bool {
|
||||
c.mu.RLock()
|
||||
defer c.mu.RUnlock()
|
||||
return topology.Nodes(c.shardNodes(index, shard)).ContainsID(nodeID)
|
||||
}
|
||||
|
||||
// partitionNodes returns a list of nodes that own a partition. unprotected.
|
||||
func (c *cluster) partitionNodes(partitionID int) []*topology.Node {
|
||||
// Default replica count to between one and the number of nodes.
|
||||
|
|
@ -1466,10 +1464,14 @@ func (c *cluster) unprotectedGenerateResizeJobByAction(nodeAction nodeAction) (*
|
|||
j.IDs[node.ID] = true
|
||||
continue
|
||||
}
|
||||
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
|
||||
|
||||
instr := &ResizeInstruction{
|
||||
JobID: j.ID,
|
||||
Node: toCluster.unprotectedNodeByID(node.ID),
|
||||
Coordinator: c.unprotectedCoordinatorNode(),
|
||||
Coordinator: snap.PrimaryFieldTranslationNode(),
|
||||
Sources: fragmentSourcesByNode[node.ID],
|
||||
TranslationSources: translationSourcesByNode[node.ID],
|
||||
NodeStatus: c.nodeStatus(), // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node.
|
||||
|
|
@ -1490,7 +1492,10 @@ func (c *cluster) completeCurrentJob(state string) error {
|
|||
}
|
||||
|
||||
func (c *cluster) unprotectedCompleteCurrentJob(state string) error {
|
||||
if !c.unprotectedIsCoordinator() {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
|
||||
// TODO: this needs to become: IsPrimaryFieldTranslationNode(c.Node.ID)
|
||||
if !snap.IsCoordinatorNode(c.Node.ID) {
|
||||
return ErrNodeNotCoordinator
|
||||
}
|
||||
if c.currentJob == nil {
|
||||
|
|
@ -2419,16 +2424,19 @@ func (c *cluster) unprotectedPrimaryReplicaNode() *topology.Node {
|
|||
// the case where the local node is not coordinator, then this method will forward the translation
|
||||
// request to the coordinator.
|
||||
func (c *cluster) translateFieldKeys(ctx context.Context, field *Field, keys []string, writable bool) (ids []uint64, err error) {
|
||||
coordinator := c.coordinatorNode()
|
||||
if coordinator == nil {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
|
||||
|
||||
primary := snap.PrimaryFieldTranslationNode()
|
||||
if primary == nil {
|
||||
return nil, errors.Errorf("translating field(%s/%s) keys(%v) - cannot find coordinator node", field.Index(), field.Name(), keys)
|
||||
}
|
||||
|
||||
if c.Node.ID == coordinator.ID {
|
||||
if c.Node.ID == primary.ID {
|
||||
ids, err = field.TranslateStore().TranslateKeys(keys, writable)
|
||||
} else {
|
||||
// If it's writable, then forward the request to the coordinator.
|
||||
ids, err = c.InternalClient.TranslateKeysNode(ctx, &coordinator.URI, field.Index(), field.Name(), keys, writable)
|
||||
ids, err = c.InternalClient.TranslateKeysNode(ctx, &primary.URI, field.Index(), field.Name(), keys, writable)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
|
|
@ -2588,15 +2596,18 @@ func (c *cluster) translateFieldIDs(field *Field, ids map[uint64]struct{}) (map[
|
|||
}
|
||||
|
||||
func (c *cluster) translateFieldListIDs(field *Field, ids []uint64) (keys []string, err error) {
|
||||
coordinator := c.coordinatorNode()
|
||||
if coordinator == nil {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
|
||||
|
||||
primary := snap.PrimaryFieldTranslationNode()
|
||||
if primary == nil {
|
||||
return nil, errors.Errorf("translating field(%s/%s) ids(%v) - cannot find coordinator node", field.Index(), field.Name(), ids)
|
||||
}
|
||||
|
||||
if c.Node.ID == coordinator.ID {
|
||||
if c.Node.ID == primary.ID {
|
||||
keys, err = field.TranslateStore().TranslateIDs(ids)
|
||||
} else {
|
||||
keys, err = c.InternalClient.TranslateIDsNode(context.Background(), &coordinator.URI, field.Index(), field.Name(), ids)
|
||||
keys, err = c.InternalClient.TranslateIDsNode(context.Background(), &primary.URI, field.Index(), field.Name(), ids)
|
||||
}
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "translating field(%s/%s) ids(%v)", field.Index(), field.Name(), ids)
|
||||
|
|
@ -2669,10 +2680,13 @@ func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, ke
|
|||
return nil, ErrIndexNotFound
|
||||
}
|
||||
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
|
||||
|
||||
// Split keys by partition.
|
||||
keysByPartition := make(map[int][]string, c.partitionN)
|
||||
for key := range keySet {
|
||||
partitionID := c.Topology.KeyPartition(indexName, key)
|
||||
partitionID := snap.KeyToKeyPartition(indexName, key)
|
||||
keysByPartition[partitionID] = append(keysByPartition[partitionID], key)
|
||||
}
|
||||
|
||||
|
|
@ -2686,7 +2700,7 @@ func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, ke
|
|||
g.Go(func() (err error) {
|
||||
var ids []uint64
|
||||
|
||||
primary := c.primaryPartitionNode(partitionID)
|
||||
primary := snap.PrimaryPartitionNode(partitionID)
|
||||
if primary == nil {
|
||||
return errors.Errorf("translating index(%s) keys(%v) on partition(%d) - cannot find primary node", indexName, keys, partitionID)
|
||||
}
|
||||
|
|
@ -2949,10 +2963,13 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS
|
|||
return nil, newNotFoundError(ErrIndexNotFound, indexName)
|
||||
}
|
||||
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
|
||||
|
||||
// Split ids by partition.
|
||||
idsByPartition := make(map[int][]uint64, c.partitionN)
|
||||
for id := range idSet {
|
||||
partitionID := c.idPartition(indexName, id)
|
||||
partitionID := snap.IDToShardPartition(indexName, id)
|
||||
idsByPartition[partitionID] = append(idsByPartition[partitionID], id)
|
||||
}
|
||||
|
||||
|
|
@ -2966,7 +2983,7 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS
|
|||
g.Go(func() (err error) {
|
||||
var keys []string
|
||||
|
||||
primary := c.primaryPartitionNode(partitionID)
|
||||
primary := snap.PrimaryPartitionNode(partitionID)
|
||||
if primary == nil {
|
||||
return errors.Errorf("translating index(%s) ids(%v) on partition(%d) - cannot find primary node", indexName, ids, partitionID)
|
||||
}
|
||||
|
|
|
|||
29
executor.go
29
executor.go
|
|
@ -4712,8 +4712,11 @@ 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)
|
||||
|
||||
ret := false
|
||||
for _, node := range e.Cluster.shardNodes(index, shard) {
|
||||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
|
|
@ -5070,7 +5073,10 @@ func (e *executor) executeSetBitField(ctx context.Context, qcx *Qcx, index strin
|
|||
shard := colID / ShardWidth
|
||||
ret := false
|
||||
|
||||
for _, node := range e.Cluster.shardNodes(index, shard) {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN)
|
||||
|
||||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
|
|
@ -5113,7 +5119,10 @@ func (e *executor) executeSetValueField(ctx context.Context, qcx *Qcx, index str
|
|||
shard := colID / ShardWidth
|
||||
ret := false
|
||||
|
||||
for _, node := range e.Cluster.shardNodes(index, shard) {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN)
|
||||
|
||||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
|
|
@ -5157,10 +5166,12 @@ func (e *executor) executeClearValueField(ctx context.Context, qcx *Qcx, index s
|
|||
shard := colID / ShardWidth
|
||||
ret := false
|
||||
|
||||
for _, node := range e.Cluster.shardNodes(index, shard) {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN)
|
||||
|
||||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
// Update locally if host matches.
|
||||
if node.ID == e.Node.ID {
|
||||
|
||||
idx := e.Holder.Index(index)
|
||||
tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard})
|
||||
if err != nil {
|
||||
|
|
@ -5441,9 +5452,15 @@ func (e *executor) remoteExec(ctx context.Context, node *topology.Node, index st
|
|||
func (e *executor) shardsByNode(nodes []*topology.Node, index string, shards []uint64) (map[*topology.Node][]uint64, error) {
|
||||
m := make(map[*topology.Node][]uint64)
|
||||
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
// 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)
|
||||
|
||||
loop:
|
||||
for _, shard := range shards {
|
||||
for _, node := range e.Cluster.ShardNodes(index, shard) {
|
||||
for _, node := range snap.ShardNodes(index, shard) {
|
||||
if topology.Nodes(nodes).Contains(node) {
|
||||
m[node] = append(m[node], shard)
|
||||
continue loop
|
||||
|
|
|
|||
55
holder.go
55
holder.go
|
|
@ -1343,6 +1343,10 @@ func (s *holderSyncer) SyncHolder() error {
|
|||
s.mu.Lock() // only allow one instance of SyncHolder to be running at a time
|
||||
defer s.mu.Unlock()
|
||||
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)
|
||||
|
||||
// Iterate over schema in sorted order.
|
||||
for _, di := range s.Holder.Schema() {
|
||||
// Verify syncer has not closed.
|
||||
|
|
@ -1377,7 +1381,7 @@ func (s *holderSyncer) SyncHolder() error {
|
|||
itr.Seek(0)
|
||||
for shard, eof := itr.Next(); !eof; shard, eof = itr.Next() {
|
||||
// Ignore shards that this host doesn't own.
|
||||
if !s.Cluster.ownsShard(s.Node.ID, di.Name, shard) {
|
||||
if !snap.OwnsShard(s.Node.ID, di.Name, shard) {
|
||||
continue
|
||||
}
|
||||
|
||||
|
|
@ -1539,16 +1543,19 @@ func (s *holderSyncer) resetTranslationSync() error {
|
|||
return errors.Wrap(err, "stop translation sync")
|
||||
}
|
||||
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN)
|
||||
|
||||
// Set read-only flag for all translation stores.
|
||||
s.setTranslateReadOnlyFlags()
|
||||
s.setTranslateReadOnlyFlags(snap)
|
||||
|
||||
// Connect to each node that has a primary for which we are a replica.
|
||||
if err := s.initializeIndexTranslateReplication(); err != nil {
|
||||
if err := s.initializeIndexTranslateReplication(snap); err != nil {
|
||||
return errors.Wrap(err, "initialize index translate replication")
|
||||
}
|
||||
|
||||
// Connect to coordinator to stream field data.
|
||||
if err := s.initializeFieldTranslateReplication(); err != nil {
|
||||
if err := s.initializeFieldTranslateReplication(snap); err != nil {
|
||||
return errors.Wrap(err, "initialize field translate replication")
|
||||
}
|
||||
return nil
|
||||
|
|
@ -1619,9 +1626,10 @@ func (s *holderSyncer) stopTranslationSync() error {
|
|||
// setTranslateReadOnlyFlags updates all translation stores to enable or disable
|
||||
// writing new translation keys. Index stores are writable if the node owns the
|
||||
// partition. Field stores are writable if the node is the coordinator.
|
||||
func (s *holderSyncer) setTranslateReadOnlyFlags() {
|
||||
func (s *holderSyncer) setTranslateReadOnlyFlags(snap *topology.ClusterSnapshot) {
|
||||
s.Cluster.mu.RLock()
|
||||
isCoordinator := s.Cluster.unprotectedIsCoordinator()
|
||||
// TODO: this needs to become: IsPrimaryFieldTranslationNode(s.Cluster.Node.ID) {
|
||||
isPrimaryFieldTranslator := snap.IsCoordinatorNode(s.Cluster.Node.ID)
|
||||
|
||||
for _, index := range s.Holder.Indexes() {
|
||||
// There is a race condition here:
|
||||
|
|
@ -1642,8 +1650,8 @@ func (s *holderSyncer) setTranslateReadOnlyFlags() {
|
|||
//
|
||||
// Update: there was another path down to Index.Close(), so
|
||||
// we shrink to lock to be inside index.TranslateStore() now.
|
||||
for partitionID := 0; partitionID < s.Cluster.partitionN; partitionID++ {
|
||||
primary := s.Cluster.unprotectedPrimaryPartitionNode(partitionID)
|
||||
for partitionID := 0; partitionID < snap.PartitionN; partitionID++ {
|
||||
primary := snap.PrimaryPartitionNode(partitionID)
|
||||
isPrimary := primary != nil && s.Node.ID == primary.ID
|
||||
|
||||
if ts := index.TranslateStore(partitionID); ts != nil {
|
||||
|
|
@ -1652,7 +1660,7 @@ func (s *holderSyncer) setTranslateReadOnlyFlags() {
|
|||
}
|
||||
|
||||
for _, field := range index.Fields() {
|
||||
field.TranslateStore().SetReadOnly(!isCoordinator)
|
||||
field.TranslateStore().SetReadOnly(!isPrimaryFieldTranslator)
|
||||
}
|
||||
}
|
||||
s.Cluster.mu.RUnlock()
|
||||
|
|
@ -1660,8 +1668,8 @@ func (s *holderSyncer) setTranslateReadOnlyFlags() {
|
|||
|
||||
// initializeIndexTranslateReplication connects to each node that is the
|
||||
// primary for a partition that we are a replica of.
|
||||
func (s *holderSyncer) initializeIndexTranslateReplication() error {
|
||||
for _, node := range s.Cluster.Nodes() {
|
||||
func (s *holderSyncer) initializeIndexTranslateReplication(snap *topology.ClusterSnapshot) error {
|
||||
for _, node := range snap.Nodes {
|
||||
// Skip local node.
|
||||
if node.ID == s.Node.ID {
|
||||
continue
|
||||
|
|
@ -1673,8 +1681,8 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error {
|
|||
if !index.Keys() {
|
||||
continue
|
||||
}
|
||||
for partitionID := 0; partitionID < s.Cluster.partitionN; partitionID++ {
|
||||
partitionNodes := s.Cluster.partitionNodes(partitionID)
|
||||
for partitionID := 0; partitionID < snap.PartitionN; partitionID++ {
|
||||
partitionNodes := snap.PartitionNodes(partitionID)
|
||||
isPrimary := partitionNodes[0].ID == node.ID // remote is primary?
|
||||
isReplica := topology.Nodes(partitionNodes[1:]).ContainsID(s.Node.ID) // local is replica?
|
||||
if !isPrimary || !isReplica {
|
||||
|
|
@ -1713,9 +1721,10 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error {
|
|||
}
|
||||
|
||||
// initializeFieldTranslateReplication connects the coordinator to stream field data.
|
||||
func (s *holderSyncer) initializeFieldTranslateReplication() error {
|
||||
func (s *holderSyncer) initializeFieldTranslateReplication(snap *topology.ClusterSnapshot) error {
|
||||
// Skip if coordinator.
|
||||
if s.Cluster.isCoordinator() {
|
||||
// TODO: this needs to become: IsPrimaryFieldTranslationNode(s.Cluster.Node.ID) {
|
||||
if !snap.IsCoordinatorNode(s.Cluster.Node.ID) {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -1737,9 +1746,9 @@ func (s *holderSyncer) initializeFieldTranslateReplication() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// Connect to coordinator and begin streaming.
|
||||
coordinator := s.Cluster.coordinatorNode()
|
||||
rd, err := s.Holder.OpenTranslateReader(context.Background(), coordinator.URI.String(), m)
|
||||
// Connect to primary and begin streaming.
|
||||
primary := snap.PrimaryFieldTranslationNode()
|
||||
rd, err := s.Holder.OpenTranslateReader(context.Background(), primary.URI.String(), m)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -1754,6 +1763,9 @@ func (s *holderSyncer) initializeFieldTranslateReplication() error {
|
|||
}
|
||||
|
||||
func (s *holderSyncer) readIndexTranslateReader(rd TranslateEntryReader) {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN)
|
||||
|
||||
for {
|
||||
var entry TranslateEntry
|
||||
if err := rd.ReadEntry(&entry); err != nil {
|
||||
|
|
@ -1769,7 +1781,7 @@ func (s *holderSyncer) readIndexTranslateReader(rd TranslateEntryReader) {
|
|||
}
|
||||
|
||||
// Apply replication to store.
|
||||
store := idx.TranslateStore(s.Cluster.Topology.KeyPartition(entry.Index, entry.Key))
|
||||
store := idx.TranslateStore(snap.KeyToKeyPartition(entry.Index, entry.Key))
|
||||
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
|
||||
s.Holder.Logger.Printf("cannot force set index translation data: %d=%q", entry.ID, entry.Key)
|
||||
return
|
||||
|
|
@ -1825,6 +1837,9 @@ func (c *holderCleaner) IsClosing() bool {
|
|||
// CleanHolder compares the holder with the cluster state and removes
|
||||
// any unnecessary fragments and files.
|
||||
func (c *holderCleaner) CleanHolder() error {
|
||||
// Create a snapshot of the cluster to use for node/partition calculations.
|
||||
snap := topology.NewClusterSnapshot(c.Cluster.noder, c.Cluster.Hasher, c.Cluster.ReplicaN)
|
||||
|
||||
for _, index := range c.Holder.Indexes() {
|
||||
// Verify cleaner has not closed.
|
||||
if c.IsClosing() {
|
||||
|
|
@ -1832,7 +1847,7 @@ func (c *holderCleaner) CleanHolder() error {
|
|||
}
|
||||
|
||||
// Get the fragments that node is responsible for (based on hash(index, node)).
|
||||
containedShards := c.Cluster.containsShards(index.Name(), index.AvailableShards(includeRemote), c.Node)
|
||||
containedShards := snap.ContainsShards(index.Name(), index.AvailableShards(includeRemote), c.Node)
|
||||
|
||||
// Get the fragments registered in memory.
|
||||
for _, field := range index.Fields() {
|
||||
|
|
|
|||
|
|
@ -138,6 +138,18 @@ func (c *ClusterSnapshot) IsPrimaryFieldTranslationNode(nodeID string) bool {
|
|||
return c.PrimaryFieldTranslationNode().ID == nodeID
|
||||
}
|
||||
|
||||
// IsCoordinatorNode returns true if nodeID represents the coordinator
|
||||
// node responsible for field translation. TODO: this is temporary until
|
||||
// we transition over to using primary
|
||||
func (c *ClusterSnapshot) IsCoordinatorNode(nodeID string) bool {
|
||||
for i := range c.Nodes {
|
||||
if c.Nodes[i].ID == nodeID && c.Nodes[i].IsCoordinator {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// PrimaryPartitionNode returns the primary node of the given partition.
|
||||
func (c *ClusterSnapshot) PrimaryPartitionNode(partition int) *Node {
|
||||
if nodes := c.PartitionNodes(partition); len(nodes) > 0 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue