From a94c745a7f137c58e54fe690fc71a82415a04ca9 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 29 Mar 2021 10:01:21 -0500 Subject: [PATCH] drop "batched" flag from Add operation The "batched" flag creates a complexity which is that the return value of Add might or might not be meaningful, but it doesn't really buy us very much. If we are concerned about the ops log size of writing single ops as 21-byte arrays of 1 op rather than as 13-byte ops, we can make the AddN code smarter about how it writes ops. And probably should. Along with this, change Remove to use the batched operation form, which writes a more meaningful ops log, and return a meaningful value for changes made. Otherwise, it ends up writing potentially thousands of ops to the ops log without reporting any OpN, because the number of ops written isn't the same as the number of changes those ops made. This could result in files growing by megabytes without OpN changing. There was a comment here about a test failing with RemoveN. I can't prove it, but I strongly suspect that this was actually a result of that test case hitting a particular bug that we eventually fixed, and which we might have fixed sooner if we'd realized why using RemoveN made that test fail. --- bluegreentx.go | 6 +++--- bolt.go | 20 ++++---------------- bolt_test.go | 28 ++++++++++++++-------------- catcher.go | 4 ++-- dbshard_internal_test.go | 3 +-- fragment.go | 5 ++--- fragment_internal_test.go | 4 ++-- rbf.go | 20 +++++--------------- rrtx.go | 28 +++++----------------------- stattx.go | 4 ++-- tx.go | 23 +---------------------- 11 files changed, 41 insertions(+), 104 deletions(-) diff --git a/bluegreentx.go b/bluegreentx.go index e0d9ddee9..394f4e296 100644 --- a/bluegreentx.go +++ b/bluegreentx.go @@ -539,7 +539,7 @@ func (c *blueGreenTx) isIn(index, field, view string, shard uint64, ckey uint64) return } -func (c *blueGreenTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { +func (c *blueGreenTx) Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { c.checker.see(index, field, view, shard) //vv("blueGreenTx) Add(index=%v, field=%v, view=%v, shard=%v", index, field, view, shard) defer func() { @@ -554,10 +554,10 @@ func (c *blueGreenTx) Add(index, field, view string, shard uint64, batched bool, a2 := make([]uint64, len(a)) copy(a2, a) - ach, errA := c.a.Add(index, field, view, shard, batched, a...) + ach, errA := c.a.Add(index, field, view, shard, a...) _, _ = ach, errA - bch, errB := c.b.Add(index, field, view, shard, batched, a2...) + bch, errB := c.b.Add(index, field, view, shard, a2...) if !c.o.blueGreenOff { diff --git a/bolt.go b/bolt.go index 3b37a2bec..fcf07e296 100644 --- a/bolt.go +++ b/bolt.go @@ -625,28 +625,16 @@ func (tx *BoltTx) RemoveContainer(index, field, view string, shard uint64, ckey } // Add sets all the a bits hot in the specified fragment. -func (tx *BoltTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { - return tx.addOrRemove(index, field, view, shard, batched, false, a...) +func (tx *BoltTx) Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { + return tx.addOrRemove(index, field, view, shard, false, a...) } // Remove clears all the specified a bits in the chosen fragment. func (tx *BoltTx) Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { - const batched = false - const remove = true - return tx.addOrRemove(index, field, view, shard, batched, remove, a...) + return tx.addOrRemove(index, field, view, shard, true, a...) } -func (tx *BoltTx) addOrRemove(index, field, view string, shard uint64, batched, remove bool, a ...uint64) (changeCount int, err error) { - - // pure hack to match RoaringTx - defer func() { - if !remove && !batched { - if changeCount > 0 { - changeCount = 1 - } - } - }() - +func (tx *BoltTx) addOrRemove(index, field, view string, shard uint64, remove bool, a ...uint64) (changeCount int, err error) { if len(a) == 0 { return 0, nil } diff --git a/bolt_test.go b/bolt_test.go index 3c41ddc41..7a8c32f99 100644 --- a/bolt_test.go +++ b/bolt_test.go @@ -55,7 +55,7 @@ func BoltMustSetBitvalue(dbwrap *BoltWrapper, index, field, view string, shard u tx, _ := dbwrap.NewTx(writable, index, Txo{}) // add a bit - changed, err := tx.Add(index, field, view, shard, doBatched, putme) + changed, err := tx.Add(index, field, view, shard, putme) if changed != 1 { panic("should have 1 bit changed") } @@ -122,7 +122,7 @@ func TestBolt_DeleteFragment(t *testing.T) { views := []string{"v1", "v2"} for _, view := range views { for _, v := range bits { - changed, err := tx.Add(index, field, view, shard, doBatched, v) + changed, err := tx.Add(index, field, view, shard, v) if changed <= 0 { panic("should have changed") } @@ -238,7 +238,7 @@ func TestBolt_SetBitmap(t *testing.T) { index, field, view, shard := "i", "f", "v", uint64(0) tx, _ := dbwrap.NewTx(writable, index, Txo{}) bitvalue := uint64(0) - changed, err := tx.Add(index, field, view, shard, doBatched, bitvalue) + changed, err := tx.Add(index, field, view, shard, bitvalue) if changed <= 0 { panic("should have changed") } @@ -280,14 +280,14 @@ func TestBolt_OffsetRange(t *testing.T) { tx, _ := dbwrap.NewTx(writable, index, Txo{}) bitvalue := uint64(1 << 20) - changed, err := tx.Add(index, field, view, shard, doBatched, bitvalue) + changed, err := tx.Add(index, field, view, shard, bitvalue) if changed <= 0 { panic("should have changed") } panicOn(err) bitvalue2 := uint64(1<<20 + 1) - changed, err = tx.Add(index, field, view, shard, doBatched, bitvalue2) + changed, err = tx.Add(index, field, view, shard, bitvalue2) if changed <= 0 { panic("should have changed") } @@ -373,7 +373,7 @@ func TestBolt_Count_dense_containers(t *testing.T) { expected := 0 for i := uint64(0); i < (1<<16)+2; i += 2 { - changed, err := tx.Add(index, field, view, shard, doBatched, i) + changed, err := tx.Add(index, field, view, shard, i) panicOn(err) if changed <= 0 { panic("wat? should have changed") @@ -419,7 +419,7 @@ func TestBolt_ContainerIterator_on_one_bit(t *testing.T) { bitvalue := uint64(42) // add a bit - changed, err := tx.Add(index, field, view, shard, doBatched, bitvalue) + changed, err := tx.Add(index, field, view, shard, bitvalue) if changed <= 0 { panic("should have changed") } @@ -477,7 +477,7 @@ func TestBolt_ContainerIterator_on_one_bit_fail_to_find(t *testing.T) { searchme := putme + 1 // add a bit - changed, err := tx.Add(index, field, view, shard, doBatched, putme) + changed, err := tx.Add(index, field, view, shard, putme) if changed <= 0 { panic("should have changed") } @@ -532,7 +532,7 @@ func TestBolt_ContainerIterator_empty_iteration_loop(t *testing.T) { searchme := uint64(1 << 17) // in the next container, key:2 // add a bit - changed, err := tx.Add(index, field, view, shard, doBatched, putme) + changed, err := tx.Add(index, field, view, shard, putme) if changed <= 0 { panic("should have changed") } @@ -581,7 +581,7 @@ func TestBolt_ForEach_on_one_bit(t *testing.T) { bitvalue := uint64(42) // add a bit - changed, err := tx.Add(index, field, view, shard, doBatched, bitvalue) + changed, err := tx.Add(index, field, view, shard, bitvalue) if changed <= 0 { panic("should have changed") } @@ -1131,7 +1131,7 @@ func TestBolt_DeleteIndex(t *testing.T) { bitvalue := uint64(777) bits := []uint64{0, 3, 1 << 16, 1<<16 + 3, 8 << 16} for _, v := range bits { - changed, err := tx.Add(index, field, view, shard, doBatched, v) + changed, err := tx.Add(index, field, view, shard, v) if changed <= 0 { panic("should have changed") } @@ -1139,7 +1139,7 @@ func TestBolt_DeleteIndex(t *testing.T) { } index2 := "i2" // should not be deleted, even though it shares a prefix with 'i' - changed, err := tx.Add(index2, field, view, shard, doBatched, bitvalue) + changed, err := tx.Add(index2, field, view, shard, bitvalue) if changed <= 0 { panic("should have changed") } @@ -1195,7 +1195,7 @@ func TestBolt_DeleteIndex_over100k(t *testing.T) { //limit := uint64(101) for v := uint64(1); v < limit; v++ { // shift by << 16 to get into a different shard - changed, err := tx.Add(index, field, view, shard, doBatched, v<<16) + changed, err := tx.Add(index, field, view, shard, v<<16) if changed <= 0 { panic("should have changed") } @@ -1207,7 +1207,7 @@ func TestBolt_DeleteIndex_over100k(t *testing.T) { } index2 := "i2" // should not be deleted, even though it shares a prefix with 'i' - changed, err := tx.Add(index2, field, view, shard, doBatched, bitvalue) + changed, err := tx.Add(index2, field, view, shard, bitvalue) if changed <= 0 { panic("should have changed") } diff --git a/catcher.go b/catcher.go index 09f24ce93..bdde74934 100644 --- a/catcher.go +++ b/catcher.go @@ -151,7 +151,7 @@ func (c *catcherTx) IsDone() bool { return c.b.IsDone() } -func (c *catcherTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { +func (c *catcherTx) Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { defer func() { if r := recover(); r != nil { @@ -159,7 +159,7 @@ func (c *catcherTx) Add(index, field, view string, shard uint64, batched bool, a panic(r) } }() - return c.b.Add(index, field, view, shard, batched, a...) + return c.b.Add(index, field, view, shard, a...) } func (c *catcherTx) Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { diff --git a/dbshard_internal_test.go b/dbshard_internal_test.go index 18858485c..c966e4a5b 100644 --- a/dbshard_internal_test.go +++ b/dbshard_internal_test.go @@ -300,7 +300,6 @@ func makeRBFtestDB(path string, h *Holder, shard uint64) { func makeTxTestDBWithViewsShards(holder *Holder, idx *Index, exp *FieldView2Shards) { // TODO(jea): need date time quantum views!! - batched := false for field, viewmap := range exp.m { for view, shset := range viewmap { @@ -310,7 +309,7 @@ func makeTxTestDBWithViewsShards(holder *Holder, idx *Index, exp *FieldView2Shar // simply write 1 bit to each shard to force its creation. bits := []uint64{(shard << shardwidth.Exponent) + 1} tx := idx.holder.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard}) - changeCount, err := tx.Add(idx.name, field, view, shard, batched, bits...) + changeCount, err := tx.Add(idx.name, field, view, shard, bits...) panicOn(err) if changeCount != len(bits) { panic(fmt.Sprintf("writing field '%v', view '%v' shard '%v', expected changeCount to equal len bits = %v but was %v", field, view, shard, len(bits), changeCount)) diff --git a/fragment.go b/fragment.go index bd29a0246..7f5990df4 100644 --- a/fragment.go +++ b/fragment.go @@ -708,7 +708,7 @@ func (f *fragment) unprotectedSetBit(tx Tx, rowID, columnID uint64) (changed boo // Write to storage. changeCount := 0 - changeCount, err = tx.Add(f.index(), f.field(), f.view(), f.shard, doBatched, pos) + changeCount, err = tx.Add(f.index(), f.field(), f.view(), f.shard, pos) changed = changeCount > 0 if err != nil { return false, errors.Wrap(err, "writing") @@ -2465,8 +2465,7 @@ func (f *fragment) importPositions(tx Tx, set, clear []uint64, rowSet map[uint64 f.stats.Count(MetricImportingN, int64(len(set)), 1) // TODO benchmark Add/RemoveN behavior with sorted/unsorted positions - // Note: AddN() avoids writing to the op-log. While Add() does. - changedN, err := tx.Add(f.index(), f.field(), f.view(), f.shard, !doBatched, set...) + changedN, err := tx.Add(f.index(), f.field(), f.view(), f.shard, set...) if err != nil { return errors.Wrap(err, "adding positions") } diff --git a/fragment_internal_test.go b/fragment_internal_test.go index fa45aace4..baa5120c3 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -5426,7 +5426,7 @@ func TestRemapCache(t *testing.T) { }() // create a container - _, err := tx.Add(index, field, view, shard, !doBatched, 65537) + _, err := tx.Add(index, field, view, shard, 65537) if err != nil { t.Fatalf("storage add: %v", err) } @@ -5439,7 +5439,7 @@ func TestRemapCache(t *testing.T) { _ = f.mustRow(tx, 0) // add a bit that isn't in that container, so that container doesn't // change - _, err = tx.Add(index, field, view, shard, !doBatched, 2) + _, err = tx.Add(index, field, view, shard, 2) if err != nil { t.Fatalf("storage add: %v", err) } diff --git a/rbf.go b/rbf.go index 4fa2f8900..c9fb65c0d 100644 --- a/rbf.go +++ b/rbf.go @@ -30,6 +30,7 @@ import ( rbfcfg "github.com/pilosa/pilosa/v2/rbf/cfg" "github.com/pilosa/pilosa/v2/roaring" txkey "github.com/pilosa/pilosa/v2/short_txkey" + //txkey "github.com/pilosa/pilosa/v2/txkey" "github.com/pkg/errors" ) @@ -241,27 +242,16 @@ func (tx *RBFTx) RemoveContainer(index, field, view string, shard uint64, key ui } // Add sets all the a bits hot in the specified fragment. -func (tx *RBFTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { - return tx.addOrRemove(index, field, view, shard, batched, false, a...) +func (tx *RBFTx) Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { + return tx.addOrRemove(index, field, view, shard, false, a...) } // Remove clears all the specified a bits in the chosen fragment. func (tx *RBFTx) Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { - const batched = false - const remove = true - return tx.addOrRemove(index, field, view, shard, batched, remove, a...) + return tx.addOrRemove(index, field, view, shard, true, a...) } -func (tx *RBFTx) addOrRemove(index, field, view string, shard uint64, batched, remove bool, a ...uint64) (changeCount int, err error) { - // pure hack to match RoaringTx - defer func() { - if !remove && !batched { - if changeCount > 0 { - changeCount = 1 - } - } - }() - +func (tx *RBFTx) addOrRemove(index, field, view string, shard uint64, remove bool, a ...uint64) (changeCount int, err error) { if len(a) == 0 { return 0, nil } diff --git a/rrtx.go b/rrtx.go index ca89b6d3d..75cc46034 100644 --- a/rrtx.go +++ b/rrtx.go @@ -200,25 +200,16 @@ func (tx *RoaringTx) RemoveContainer(index, field, view string, shard uint64, ke return nil } -func (tx *RoaringTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { +func (tx *RoaringTx) Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { //vv("RoaringTx.Add(index='%v', shard='%v') stack=\n%v", index, shard, stack()) b, err := tx.bitmap(index, field, view, shard) if err != nil { return 0, err } - if !batched { - changed, err := b.Add(a...) - if changed { - return 1, err - } - return 0, err - } - // Note: do not replace b.AddN() with b.DirectAddN(). - // DirectAddN() does not do op-log operations inside roaring - // This creates a problem because RoaringTx needs the op-log - // to know when to flush the fragment to disk. - count, err := b.AddN(a...) // AddN does oplog batches. needed to keep op-log up to date. + // DirectAddN() does not do op-log operations inside roaring, so the + // on-disk representation no longer matches the in-memory operations. + count, err := b.AddN(a...) return count, err } @@ -227,16 +218,7 @@ func (tx *RoaringTx) Remove(index, field, view string, shard uint64, a ...uint64 if err != nil { return 0, err } - changed, err := b.Remove(a...) // green TestFragment_Bug_Q2DoubleDelete - if changed { - return 1, err - } else { - return 0, err - } - - // Note: don't replace b.Remove(a...) with b.RemoveN(a...) or - // with b.DirectRemoveN(a...). If you do, you'll see - // TestFragment_Bug_Q2DoubleDelete go red. + return b.RemoveN(a...) } func (tx *RoaringTx) Contains(index, field, view string, shard uint64, v uint64) (exists bool, err error) { diff --git a/stattx.go b/stattx.go index 7009de9ac..1d815b32e 100644 --- a/stattx.go +++ b/stattx.go @@ -445,7 +445,7 @@ func (c *statTx) IsDone() (done bool) { return c.b.IsDone() } -func (c *statTx) Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) { +func (c *statTx) Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { me := kAdd t0 := time.Now() @@ -459,7 +459,7 @@ func (c *statTx) Add(index, field, view string, shard uint64, batched bool, a .. panic(r) } }() - return c.b.Add(index, field, view, shard, batched, a...) + return c.b.Add(index, field, view, shard, a...) } func (c *statTx) Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) { diff --git a/tx.go b/tx.go index f33648cb0..958491a39 100644 --- a/tx.go +++ b/tx.go @@ -22,11 +22,6 @@ import ( //txkey "github.com/pilosa/pilosa/v2/txkey" ) -// batch operations want Tx.Add(batched=doBatch), while bit-at-a-time want Tx.Add(batched=!doBatched) -// Used in Tx.Add() to get consistency between RoaringTx and other Tx implementations on -// the changeCount returned. -const doBatched = false // must be false, do not change this without adjusting the Add() implementations correspondingly. - // writable initializes Tx that update, use !writable for read-only. const writable = true @@ -136,23 +131,7 @@ type Tx interface { // in the specified fragment. RemoveContainer(index, field, view string, shard uint64, ckey uint64) error - // Add adds the 'a' bits to the specified fragment. - // - // Using batched=true allows efficient bulk-import. - // - // Notes on the RoaringTx implementation: - // If the batched flag is true, then the roaring.Bitmap.AddN() is used, which does oplog batches. - // If the batched flag is false, then the roaring.Bitmap.Add() is used, which does simple opTypeAdd single adds. - // - // Beware: if batched is false, then changeCount will only ever be 0 or 1, - // because it calls roaring.Add(). - // If batched is true, we call roaring.DirectAddN() and then changeCount - // will be accurate if the changeCount is greater than 0. - // - // Hence: only ever call Add(batched=false) if changeCount is expected to be 0 or 1. - // Or, must use Add(batched=true) if changeCount can be > 1. - // - Add(index, field, view string, shard uint64, batched bool, a ...uint64) (changeCount int, err error) + Add(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error) // Remove removes the 'a' values from the Bitmap for the fragment. Remove(index, field, view string, shard uint64, a ...uint64) (changeCount int, err error)