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