mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 09:05:55 +00:00
add deleted (rebalanced) shards to remoteAvailableShards
This commit is contained in:
parent
6bd81b87eb
commit
0e2bb550db
3 changed files with 18 additions and 5 deletions
15
holder.go
15
holder.go
|
|
@ -1233,6 +1233,15 @@ func (c *holderCleaner) CleanHolder() error {
|
|||
|
||||
// Get the fragments registered in memory.
|
||||
for _, field := range index.Fields() {
|
||||
// deletedShards is used to track which shards for the field
|
||||
// were deleted. Any shards that get deleted from this node
|
||||
// get added to remoteAvailableShards. This is done because
|
||||
// the CleanHolder process is cleaning up shards which got
|
||||
// moved to other nodes. Because those shards still exist
|
||||
// (just no longer on this particular node), this node still
|
||||
// needs to consider each of them as an available shard in
|
||||
// the cluster.
|
||||
var deletedShards []uint64
|
||||
for _, view := range field.views() {
|
||||
for _, fragment := range view.allFragments() {
|
||||
fragShard := fragment.shard
|
||||
|
|
@ -1244,6 +1253,12 @@ func (c *holderCleaner) CleanHolder() error {
|
|||
if err := view.deleteFragment(fragShard); err != nil {
|
||||
return errors.Wrap(err, "deleting fragment")
|
||||
}
|
||||
deletedShards = append(deletedShards, fragShard)
|
||||
}
|
||||
}
|
||||
if len(deletedShards) > 0 {
|
||||
if err := field.AddRemoteAvailableShards(roaring.NewBitmap(deletedShards...)); err != nil {
|
||||
return errors.Wrap(err, "adding remote available shards")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -138,7 +138,7 @@ func TestClusterResize_EmptyNodes(t *testing.T) {
|
|||
}
|
||||
|
||||
// Ensure that adding a node correctly resizes the cluster.
|
||||
func TestXClusterResize_AddNode(t *testing.T) {
|
||||
func TestClusterResize_AddNode(t *testing.T) {
|
||||
t.Run("NoData", func(t *testing.T) {
|
||||
clus := test.MustRunCluster(t, 2)
|
||||
defer clus.Close()
|
||||
|
|
@ -387,8 +387,6 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
} else if res != exp {
|
||||
t.Fatalf("unexpected result: %s", res)
|
||||
} else {
|
||||
fmt.Println(res)
|
||||
}
|
||||
|
||||
// Configure node1
|
||||
|
|
|
|||
4
view.go
4
view.go
|
|
@ -278,11 +278,11 @@ func (v *view) CreateFragmentIfNotExists(shard uint64) (*fragment, error) {
|
|||
frag.RowAttrStore = v.rowAttrStore
|
||||
|
||||
v.fragments[shard] = frag
|
||||
v.notifyIfNew(shard)
|
||||
v.notifyIfNewShard(shard)
|
||||
return frag, nil
|
||||
}
|
||||
|
||||
func (v *view) notifyIfNew(shard uint64) {
|
||||
func (v *view) notifyIfNewShard(shard uint64) {
|
||||
if v.remoteShardPresent(shard) { //checks the fields remoteShards bitmap to see if broadcast needed
|
||||
return
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue