From d322cc7aef5852db85fb659e6aaef9666482ddd5 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 8 Mar 2022 13:00:08 -0600 Subject: [PATCH 1/8] batch size 65536 --- executor.go | 144 +++++++++++++++++++++++++++++++++++++++++++--------- 1 file changed, 121 insertions(+), 23 deletions(-) diff --git a/executor.go b/executor.go index 136753a65..f12c957d3 100644 --- a/executor.go +++ b/executor.go @@ -8309,33 +8309,23 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i func DeleteRows(ctx context.Context, src *Row, idx *Index, shard uint64) (bool, error) { return DeleteRowsWithFlow(ctx, src, idx, shard, false) } -func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, normalFlow bool) (bool, error) { +func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (bool, error) { var existenceFragment *fragment var deletedRowID uint64 var commitor Commitor = &NopCommitor{} - var err error - if len(src.segments) == 0 { //nothing to remove - return false, nil - } - columns := src.segments[0].data //should only be one segment - if columns.Count() == 0 { - return false, nil - } - - if idx.Keys() { - //store columns in exits field ToBeDelete row commited - if normalFlow { // normalFlow is the standard path, "not normal" is recoverory - existenceFragment = idx.Holder().fragment(idx.Name(), existenceFieldName, viewStandard, shard) - if existenceFragment == nil { - //no exists field - return false, errors.New("can't bulk delete without existence field") - } - deletedRowID, err = transactExistRow(ctx, idx, shard, existenceFragment, src) - } - commitor, err = deleteKeyTranslation(ctx, idx, shard, columns) - if err != nil { - return false, err + var err error //store columns in exits field ToBeDelete row commited + if normalFlow { // normalFlow is the standard path, "not normal" is recoverory + existenceFragment = idx.Holder().fragment(idx.Name(), existenceFieldName, viewStandard, shard) + if existenceFragment == nil { + //no exists field + return false, errors.New("can't bulk delete without existence field") } + src := NewRowFromBitmap(columns) + deletedRowID, err = transactExistRow(ctx, idx, shard, existenceFragment, src) + } + commitor, err = deleteKeyTranslation(ctx, idx, shard, columns) + if err != nil { + return false, err } writeTx := idx.Txf().NewTx(Txo{Write: writable, Index: idx, Shard: shard}) if err != nil { @@ -8420,6 +8410,114 @@ func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, return changed, nil } +func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (bool, error) { + var existenceFragment *fragment + var deletedRowID uint64 + var commitor Commitor = &NopCommitor{} + var err error + writeTx := idx.Txf().NewTx(Txo{Write: writable, Index: idx, Shard: shard}) + if err != nil { + return false, err + } + defer writeTx.Rollback() + changed := false + defer func() { + // if there is an error in the key commit, then rollback the delete + // write records before keys to remove possiblity of unmatch keys=records + err = writeTx.Commit() + if err != nil { + changed = false + commitor.Rollback() + return + } + }() + findExisting := roaring.NewBitmapBitmapFilter(columns, func(p uint64) error { return nil }) + resChan := make(chan countResults) + clearFragment := func(frag *fragment) (bool, error) { + posChan := make(chan uint64, 8192) + findExisting.SetCallback(func(pos uint64) error { + posChan <- pos + return nil + }) + go writeTx.RemoveChannel(frag.index(), frag.field(), frag.view(), frag.shard, posChan, resChan) + + err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) + close(posChan) + if err != nil { + return false, err + } + r := <-resChan + return r.changeCount > 0, r.err + } + + for _, field := range idx.Fields() { + for _, view := range field.views() { + + frag, ok := view.fragments[shard] + if !ok { + continue + } + c, err := clearFragment(frag) + if err != nil { + return false, err + } + if c { + changed = true + } + + } + } + close(resChan) + if existenceFragment != nil { //a string keys have been deleted and the deleteRow was created + if normalFlow { + existenceFragment.clearRow(writeTx, deletedRowID) + } else { + // this is if we are recovering from failure and cleaning up + rows, err := existenceFragment.rows(ctx, writeTx, 1) + if err != nil { + return false, err + } + for _, rowId := range rows { + existenceFragment.clearRow(writeTx, rowId) + } + } + } + return changed, nil +} + +func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, normalFlow bool) (bool, error) { + if len(src.segments) == 0 { //nothing to remove + return false, nil + } + columns := src.segments[0].data //should only be one segment + if columns.Count() == 0 { + return false, nil + } + bits := src.segments[0].data.Slice() + min := func(a, b int) int { + if a <= b { + return a + } + return b + } + var change bool + var err error + limit := 65536 + for i := 0; i < len(bits); i += limit { + batch := roaring.NewBitmap(bits[i:min(i+limit, len(bits))]...) + if idx.Keys() { + change, err = DeleteRowsWithFlowWithKeys(ctx, batch, idx, shard, normalFlow) + } else { + change, err = DeleteRowsWithOutKeysFlow(ctx, batch, idx, shard, normalFlow) + } + if err != nil { + return change, err + } + } + return change, err + +} + type Commitor interface { Rollback() Commit() error From dc7493ff75e98920211428859c8f770e5d84d458 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 8 Mar 2022 15:06:38 -0600 Subject: [PATCH 2/8] optimize container after remove --- rbf.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/rbf.go b/rbf.go index 5360958cb..1cdc59261 100644 --- a/rbf.go +++ b/rbf.go @@ -253,6 +253,8 @@ func (tx *RBFTx) RemoveChannel(index, field, view string, shard uint64, a chan u return } } else { + rc = roaring.Optimize(rc) + err = tx.tx.PutContainer(name, lastHi, rc) if err != nil { resChan <- countResults{0, errors.Wrap(err, "failed to put container")} From eff54f49a6fbbab4a48a4c5182d82b8e0783fb7a Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 8 Mar 2022 17:26:25 -0600 Subject: [PATCH 3/8] qcx finish/reset --- executor.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/executor.go b/executor.go index f12c957d3..e3f52d806 100644 --- a/executor.go +++ b/executor.go @@ -8303,6 +8303,8 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i return } + qcx.Finish() // have to release the tx in order for a snapshot to be able to occur + qcx.Reset() return DeleteRowsWithFlow(ctx, src, idx, shard, true) } From 1a8989a69b7e012765a7f27baefa7e11df830f25 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 8 Mar 2022 21:04:39 -0600 Subject: [PATCH 4/8] . --- executor.go | 2 +- txfactory.go | 3 --- 2 files changed, 1 insertion(+), 4 deletions(-) diff --git a/executor.go b/executor.go index e3f52d806..885200c4f 100644 --- a/executor.go +++ b/executor.go @@ -8303,7 +8303,7 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i return } - qcx.Finish() // have to release the tx in order for a snapshot to be able to occur + qcx.Abort() // have to release the tx in order for a snapshot to be able to occur qcx.Reset() return DeleteRowsWithFlow(ctx, src, idx, shard, true) } diff --git a/txfactory.go b/txfactory.go index 4fddf79a9..b96111c57 100644 --- a/txfactory.go +++ b/txfactory.go @@ -158,9 +158,6 @@ func (q *Qcx) Abort() { func (q *Qcx) Reset() { q.mu.Lock() defer q.mu.Unlock() - if !q.done { - vprint.PanicOn("must call Qcx.Abort() or Qcx.Finish() before calling Reset().") - } q.unprotected_reset() } From e209759283ab723ee86f9c3db7af7ca25987cce5 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 8 Mar 2022 21:16:11 -0600 Subject: [PATCH 5/8] make maxdelete an optional param --- executor.go | 2 +- rbf/cfg/cfg.go | 4 ++++ rbf/cfg/os.go | 3 +++ 3 files changed, 8 insertions(+), 1 deletion(-) diff --git a/executor.go b/executor.go index 885200c4f..d8422bd05 100644 --- a/executor.go +++ b/executor.go @@ -8504,7 +8504,7 @@ func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, } var change bool var err error - limit := 65536 + limit := idx.holder.cfg.RBFConfig.MaxDelete for i := 0; i < len(bits); i += limit { batch := roaring.NewBitmap(bits[i:min(i+limit, len(bits))]...) if idx.Keys() { diff --git a/rbf/cfg/cfg.go b/rbf/cfg/cfg.go index 184775e33..2b65a2b2b 100644 --- a/rbf/cfg/cfg.go +++ b/rbf/cfg/cfg.go @@ -41,6 +41,9 @@ type Config struct { // background checkpoints. It cannot be set from toml. The default is // to use stderr. Logger logger.Logger `toml:"-"` + + // The maximum number of bits to be deleted in a single transaction default(65536) + MaxDelete int `toml:"max-delete"` } func NewDefaultConfig() *Config { @@ -50,6 +53,7 @@ func NewDefaultConfig() *Config { MinWALCheckpointSize: DefaultMinWALCheckpointSize, MaxWALCheckpointSize: DefaultMaxWALCheckpointSize, FsyncEnabled: true, + MaxDelete: DefaultMaxDelete, // CI passed with 20. 50 was too big for CI, even on X-large instances. // For now we default to 0, which means use sync.Pool. diff --git a/rbf/cfg/os.go b/rbf/cfg/os.go index 5ca1f60d9..057b4320f 100644 --- a/rbf/cfg/os.go +++ b/rbf/cfg/os.go @@ -13,3 +13,6 @@ const DefaultMaxSize = 4 * (1 << 30) // size of the WAL. The size can be increased by updating the DB.MaxWALSize // and reopening the database. This setting mainly affects virtual space usage. const DefaultMaxWALSize = 4 * (1 << 30) + +// DefaultMaxDelete is the maximum number of bits that will be deleted in a single batch +const DefaultMaxDelete = 65536 From 82d3f54a2838b227b0c70caee8ce65d68ae9aaa3 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 10 Mar 2022 15:17:28 -0600 Subject: [PATCH 6/8] new delete flow to allow for qcx reset --- executor.go | 33 +++++++++++++++++---------------- 1 file changed, 17 insertions(+), 16 deletions(-) diff --git a/executor.go b/executor.go index d8422bd05..b5391393c 100644 --- a/executor.go +++ b/executor.go @@ -8243,9 +8243,23 @@ func (e *executor) executeDeleteRecords(ctx context.Context, qcx *Qcx, index str return false, errors.New("Delete() only accepts a single bitmap input") } + if len(c.Children) != 1 { + + } + + row, err := e.executeBitmapCall(ctx, qcx, index, c.Children[0], shards, opt) + qcx.Abort() + qcx.Reset() //release the qcx to allow for rbf checkpoint + // Execute calls in bulk on each remote node and merge. mapFn := func(ctx context.Context, shard uint64, mopt *mapOptions) (_ interface{}, err error) { - return e.executeDeleteRecordFromShard(ctx, qcx, index, c, shard) + + for i := range row.segments { + if row.segments[i].shard == shard { + return e.executeDeleteRecordFromShard(ctx, index, row.segments[i].data, shard) + } + } + return false, nil } // Merge returned results at coordinating node. @@ -8279,20 +8293,9 @@ func transactExistRow(ctx context.Context, idx *Index, shard uint64, frag *fragm } return rowID, tx.Commit() } -func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (changed bool, err error) { +func (e *executor) executeDeleteRecordFromShard(ctx context.Context, index string, columns *roaring.Bitmap, shard uint64) (changed bool, err error) { span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeDeleteRecordFromShard") defer span.Finish() - //need to build the bitmap in the call - child := c.Children[0] - src, er := e.executeBitmapCallShard(ctx, qcx, index, child, shard) - if er != nil { - err = er - return - } - if len(src.segments) == 0 { //nothing to remove - return - } - columns := src.segments[0].data //should only be one segment if columns.Count() == 0 { return } @@ -8302,9 +8305,7 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, qcx *Qcx, i err = newNotFoundError(ErrIndexNotFound, index) return } - - qcx.Abort() // have to release the tx in order for a snapshot to be able to occur - qcx.Reset() + src := NewRowFromBitmap(columns) return DeleteRowsWithFlow(ctx, src, idx, shard, true) } From e64dde0b4c6b7f44a3ad9df63122a520480dd4e4 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Fri, 11 Mar 2022 06:12:38 -0600 Subject: [PATCH 7/8] remove round trip --- executor.go | 27 ++++++++++----------------- 1 file changed, 10 insertions(+), 17 deletions(-) diff --git a/executor.go b/executor.go index b5391393c..5b8c2c092 100644 --- a/executor.go +++ b/executor.go @@ -8242,24 +8242,12 @@ func (e *executor) executeDeleteRecords(ctx context.Context, qcx *Qcx, index str } else if len(c.Children) > 1 { return false, errors.New("Delete() only accepts a single bitmap input") } - - if len(c.Children) != 1 { - - } - - row, err := e.executeBitmapCall(ctx, qcx, index, c.Children[0], shards, opt) qcx.Abort() qcx.Reset() //release the qcx to allow for rbf checkpoint // Execute calls in bulk on each remote node and merge. mapFn := func(ctx context.Context, shard uint64, mopt *mapOptions) (_ interface{}, err error) { - - for i := range row.segments { - if row.segments[i].shard == shard { - return e.executeDeleteRecordFromShard(ctx, index, row.segments[i].data, shard) - } - } - return false, nil + return e.executeDeleteRecordFromShard(ctx, index, c.Children[0], shard) } // Merge returned results at coordinating node. @@ -8293,18 +8281,23 @@ func transactExistRow(ctx context.Context, idx *Index, shard uint64, frag *fragm } return rowID, tx.Commit() } -func (e *executor) executeDeleteRecordFromShard(ctx context.Context, index string, columns *roaring.Bitmap, shard uint64) (changed bool, err error) { +func (e *executor) executeDeleteRecordFromShard(ctx context.Context, index string, bmCall *pql.Call, shard uint64) (changed bool, err error) { span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeDeleteRecordFromShard") defer span.Finish() - if columns.Count() == 0 { - return - } // Fetch index. idx := e.Holder.Index(index) if idx == nil { err = newNotFoundError(ErrIndexNotFound, index) return } + qcx := idx.Txf().NewQcx() + //bmCall is a bitmap + row, err := e.executeBitmapCallShard(ctx, qcx, index, bmCall, shard) + qcx.Abort() + columns := row.segments[0].data + if columns.Count() == 0 { + return + } src := NewRowFromBitmap(columns) return DeleteRowsWithFlow(ctx, src, idx, shard, true) } From f9ee18b8b073d00959f3c7dfe5528ab55bc029cd Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 14 Mar 2022 14:18:21 -0500 Subject: [PATCH 8/8] review suggestions --- executor.go | 107 +++++++++++++++++++++++----------------------------- 1 file changed, 47 insertions(+), 60 deletions(-) diff --git a/executor.go b/executor.go index 5b8c2c092..9d689c7bd 100644 --- a/executor.go +++ b/executor.go @@ -8294,6 +8294,9 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, index strin //bmCall is a bitmap row, err := e.executeBitmapCallShard(ctx, qcx, index, bmCall, shard) qcx.Abort() + if len(row.segments) == 0 { + return + } columns := row.segments[0].data if columns.Count() == 0 { return @@ -8305,6 +8308,26 @@ func (e *executor) executeDeleteRecordFromShard(ctx context.Context, index strin func DeleteRows(ctx context.Context, src *Row, idx *Index, shard uint64) (bool, error) { return DeleteRowsWithFlow(ctx, src, idx, shard, false) } +func clearFragment(writeTx Tx, columns *roaring.Bitmap, frag *fragment, resChan chan countResults) (changed bool, err error) { + posChan := make(chan uint64, 8192) + findExisting := roaring.NewBitmapBitmapFilter(columns, func(pos uint64) error { + posChan <- pos + return nil + }) + go writeTx.RemoveChannel(frag.index(), frag.field(), frag.view(), frag.shard, posChan, resChan) + + err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) + close(posChan) + if err != nil { + return false, err + } + r := <-resChan + + changed = r.changeCount > 0 + err = r.err + return +} + func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (bool, error) { var existenceFragment *fragment var deletedRowID uint64 @@ -8318,6 +8341,9 @@ func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, id } src := NewRowFromBitmap(columns) deletedRowID, err = transactExistRow(ctx, idx, shard, existenceFragment, src) + if err != nil { + return false, err + } } commitor, err = deleteKeyTranslation(ctx, idx, shard, columns) if err != nil { @@ -8352,24 +8378,7 @@ func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, id } }() - findExisting := roaring.NewBitmapBitmapFilter(columns, func(p uint64) error { return nil }) resChan := make(chan countResults) - clearFragment := func(frag *fragment) (bool, error) { - posChan := make(chan uint64, 8192) - findExisting.SetCallback(func(pos uint64) error { - posChan <- pos - return nil - }) - go writeTx.RemoveChannel(frag.index(), frag.field(), frag.view(), frag.shard, posChan, resChan) - - err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) - close(posChan) - if err != nil { - return false, err - } - r := <-resChan - return r.changeCount > 0, r.err - } for _, field := range idx.Fields() { for _, view := range field.views() { @@ -8378,7 +8387,7 @@ func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, id if !ok { continue } - c, err := clearFragment(frag) + c, err := clearFragment(writeTx, columns, frag, resChan) if err != nil { return false, err } @@ -8406,46 +8415,23 @@ func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, id return changed, nil } -func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (bool, error) { +func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (changed bool, err error) { var existenceFragment *fragment var deletedRowID uint64 var commitor Commitor = &NopCommitor{} - var err error writeTx := idx.Txf().NewTx(Txo{Write: writable, Index: idx, Shard: shard}) - if err != nil { - return false, err - } defer writeTx.Rollback() - changed := false defer func() { // if there is an error in the key commit, then rollback the delete // write records before keys to remove possiblity of unmatch keys=records - err = writeTx.Commit() + err := writeTx.Commit() if err != nil { changed = false commitor.Rollback() return } }() - findExisting := roaring.NewBitmapBitmapFilter(columns, func(p uint64) error { return nil }) resChan := make(chan countResults) - clearFragment := func(frag *fragment) (bool, error) { - posChan := make(chan uint64, 8192) - findExisting.SetCallback(func(pos uint64) error { - posChan <- pos - return nil - }) - go writeTx.RemoveChannel(frag.index(), frag.field(), frag.view(), frag.shard, posChan, resChan) - - err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) - close(posChan) - if err != nil { - return false, err - } - r := <-resChan - return r.changeCount > 0, r.err - } - for _, field := range idx.Fields() { for _, view := range field.views() { @@ -8453,7 +8439,7 @@ func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx if !ok { continue } - c, err := clearFragment(frag) + c, err := clearFragment(writeTx, columns, frag, resChan) if err != nil { return false, err } @@ -8464,24 +8450,27 @@ func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx } } close(resChan) - if existenceFragment != nil { //a string keys have been deleted and the deleteRow was created - if normalFlow { - existenceFragment.clearRow(writeTx, deletedRowID) - } else { - // this is if we are recovering from failure and cleaning up - rows, err := existenceFragment.rows(ctx, writeTx, 1) - if err != nil { - return false, err - } - for _, rowId := range rows { - existenceFragment.clearRow(writeTx, rowId) - } - } + if existenceFragment == nil { //a string keys have been deleted and the deleteRow was created + return changed, nil + } + + if normalFlow { + existenceFragment.clearRow(writeTx, deletedRowID) + return changed, nil + } + + // this is if we are recovering from failure and cleaning up + rows, err := existenceFragment.rows(ctx, writeTx, 1) + if err != nil { + return false, err + } + for _, rowId := range rows { + existenceFragment.clearRow(writeTx, rowId) } return changed, nil } -func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, normalFlow bool) (bool, error) { +func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, normalFlow bool) (change bool, err error) { if len(src.segments) == 0 { //nothing to remove return false, nil } @@ -8496,8 +8485,6 @@ func DeleteRowsWithFlow(ctx context.Context, src *Row, idx *Index, shard uint64, } return b } - var change bool - var err error limit := idx.holder.cfg.RBFConfig.MaxDelete for i := 0; i < len(bits); i += limit { batch := roaring.NewBitmap(bits[i:min(i+limit, len(bits))]...)