mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
Merge pull request #334 from travisturner/block-limits
Fix off-by-one maxRowID in block limits
This commit is contained in:
commit
f9cb7abccf
2 changed files with 110 additions and 4 deletions
|
|
@ -1776,8 +1776,8 @@ func (f *fragment) mergeBlock(id int, data []pairSet) (sets, clears []pairSet, e
|
|||
clears = make([]pairSet, len(data)+1)
|
||||
|
||||
// Limit upper row/column pair.
|
||||
maxRowID := uint64(id+1) * HashBlockSize
|
||||
maxColumnID := uint64(ShardWidth)
|
||||
maxRowID := (uint64(id+1) * HashBlockSize) - 1
|
||||
maxColumnID := uint64(ShardWidth) - 1
|
||||
|
||||
// Create buffered iterator for local block.
|
||||
itrs := make([]*bufIterator, 1, len(data)+1)
|
||||
|
|
@ -1857,8 +1857,8 @@ func (f *fragment) mergeBlock(id int, data []pairSet) (sets, clears []pairSet, e
|
|||
sets[i].rowIDs = append(sets[i].rowIDs, min.rowID)
|
||||
sets[i].columnIDs = append(sets[i].columnIDs, min.columnID)
|
||||
} else {
|
||||
clears[i].rowIDs = append(sets[i].rowIDs, min.rowID)
|
||||
clears[i].columnIDs = append(sets[i].columnIDs, min.columnID)
|
||||
clears[i].rowIDs = append(clears[i].rowIDs, min.rowID)
|
||||
clears[i].columnIDs = append(clears[i].columnIDs, min.columnID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
106
holder_test.go
106
holder_test.go
|
|
@ -490,6 +490,112 @@ func TestHolderSyncer_SyncHolder(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// Ensure holder can sync with a remote holder and respects
|
||||
// the row boundaries of the block.
|
||||
func TestHolderSyncer_BlockIteratorLimits(t *testing.T) {
|
||||
c := test.MustNewCluster(t, 3)
|
||||
c[0].Config.Cluster.ReplicaN = 3
|
||||
c[0].Config.AntiEntropy.Interval = 0
|
||||
c[1].Config.Cluster.ReplicaN = 3
|
||||
c[1].Config.AntiEntropy.Interval = 0
|
||||
err := c.Start()
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
defer c.Close()
|
||||
|
||||
_, err = c[0].API.CreateIndex(context.Background(), "i", pilosa.IndexOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("creating index i: %v", err)
|
||||
}
|
||||
_, err = c[0].API.CreateField(context.Background(), "i", "f", pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, pilosa.DefaultCacheSize))
|
||||
if err != nil {
|
||||
t.Fatalf("creating field f: %v", err)
|
||||
}
|
||||
|
||||
blockEdge := uint64(pilosa.HashBlockSize)
|
||||
|
||||
hldr0 := &test.Holder{Holder: c[0].Server.Holder()}
|
||||
hldr1 := &test.Holder{Holder: c[1].Server.Holder()}
|
||||
hldr2 := &test.Holder{Holder: c[2].Server.Holder()}
|
||||
|
||||
// Set data on the local holder.
|
||||
hldr0.SetBit("i", "f", blockEdge-1, 10)
|
||||
hldr0.SetBit("i", "f", blockEdge, 20)
|
||||
|
||||
// Set the same data on one of the replicas
|
||||
// so that we have a quorum.
|
||||
hldr1.SetBit("i", "f", blockEdge-1, 10)
|
||||
hldr1.SetBit("i", "f", blockEdge, 20)
|
||||
|
||||
// Leave the third replica empty to force a block merge.
|
||||
//
|
||||
|
||||
err = c[0].Server.SyncData()
|
||||
if err != nil {
|
||||
t.Fatalf("syncing node 0: %v", err)
|
||||
}
|
||||
|
||||
// Verify data is the same on all nodes.
|
||||
for i, hldr := range []*test.Holder{hldr0, hldr1, hldr2} {
|
||||
if a := hldr.Row("i", "f", blockEdge-1).Columns(); !reflect.DeepEqual(a, []uint64{10}) {
|
||||
t.Errorf("unexpected columns(%d/block 0): %+v", i, a)
|
||||
}
|
||||
if a := hldr.Row("i", "f", blockEdge).Columns(); !reflect.DeepEqual(a, []uint64{20}) {
|
||||
t.Errorf("unexpected columns(%d/block 1): %+v", i, a)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Ensure holder correctly handles clears during block sync.
|
||||
func TestHolderSyncer_Clears(t *testing.T) {
|
||||
c := test.MustNewCluster(t, 3)
|
||||
c[0].Config.Cluster.ReplicaN = 3
|
||||
c[0].Config.AntiEntropy.Interval = 0
|
||||
c[1].Config.Cluster.ReplicaN = 3
|
||||
c[1].Config.AntiEntropy.Interval = 0
|
||||
err := c.Start()
|
||||
if err != nil {
|
||||
t.Fatalf("starting cluster: %v", err)
|
||||
}
|
||||
defer c.Close()
|
||||
|
||||
_, err = c[0].API.CreateIndex(context.Background(), "i", pilosa.IndexOptions{})
|
||||
if err != nil {
|
||||
t.Fatalf("creating index i: %v", err)
|
||||
}
|
||||
_, err = c[0].API.CreateField(context.Background(), "i", "f", pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, pilosa.DefaultCacheSize))
|
||||
if err != nil {
|
||||
t.Fatalf("creating field f: %v", err)
|
||||
}
|
||||
|
||||
hldr0 := &test.Holder{Holder: c[0].Server.Holder()}
|
||||
hldr1 := &test.Holder{Holder: c[1].Server.Holder()}
|
||||
hldr2 := &test.Holder{Holder: c[2].Server.Holder()}
|
||||
|
||||
// Set data on the local holder that should be cleared
|
||||
// because it's the only instance of this value.
|
||||
hldr0.SetBit("i", "f", 0, 30)
|
||||
|
||||
// Set similar data on the replicas, but
|
||||
// different from what's on local. This should end
|
||||
// up being set on all replicas
|
||||
hldr1.SetBit("i", "f", 0, 20)
|
||||
hldr2.SetBit("i", "f", 0, 20)
|
||||
|
||||
err = c[0].Server.SyncData()
|
||||
if err != nil {
|
||||
t.Fatalf("syncing node 0: %v", err)
|
||||
}
|
||||
|
||||
// Verify data is the same on all nodes.
|
||||
for i, hldr := range []*test.Holder{hldr0, hldr1, hldr2} {
|
||||
if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{20}) {
|
||||
t.Errorf("unexpected columns(%d): %+v", i, a)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Ensure holder can sync time quantum views with a remote holder.
|
||||
func TestHolderSyncer_TimeQuantum(t *testing.T) {
|
||||
c := test.MustNewCluster(t, 2)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue