diff --git a/executor.go b/executor.go index c7b9d884f..0897d6392 100644 --- a/executor.go +++ b/executor.go @@ -8291,31 +8291,6 @@ 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, toClear []uint64) (changed bool, err error) { - - rowSet := make(map[uint64]struct{}) - toClear = toClear[:0] - callback := func(pos uint64) error { - toClear = append(toClear, pos) - rowID := pos / ShardWidth - rowSet[rowID] = struct{}{} - return nil - } - findExisting := roaring.NewBitmapBitmapFilter(columns, callback) - err = writeTx.ApplyFilter(frag.index(), frag.field(), frag.view(), frag.shard, 0, findExisting) - - if err != nil { - return false, err - } - if len(toClear) > 0 { - err = frag.importPositions(writeTx, []uint64{}, toClear, rowSet) - if err != nil { - return false, err - } - return true, nil - } - return false, nil -} func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, idx *Index, shard uint64, normalFlow bool) (bool, error) { var existenceFragment *fragment @@ -8368,22 +8343,19 @@ func DeleteRowsWithFlowWithKeys(ctx context.Context, columns *roaring.Bitmap, id }() - toClear := make([]uint64, 0) for _, field := range idx.Fields() { for _, view := range field.views() { - - frag, ok := view.fragments[shard] - if !ok { + frag := view.Fragment(shard) + if frag == nil { continue } - c, err := clearFragment(writeTx, columns, frag, toClear) + c, err := frag.clearRecordsByBitmap(writeTx, columns) if err != nil { return false, err } if c { changed = true } - } } if existenceFragment != nil { //a string keys have been deleted and the deleteRow was created @@ -8419,15 +8391,14 @@ func DeleteRowsWithOutKeysFlow(ctx context.Context, columns *roaring.Bitmap, idx return } }() - toClear := make([]uint64, 0) for _, field := range idx.Fields() { for _, view := range field.views() { - frag, ok := view.fragments[shard] - if !ok { + frag := view.Fragment(shard) + if frag == nil { continue } - c, err := clearFragment(writeTx, columns, frag, toClear) + c, err := frag.clearRecordsByBitmap(writeTx, columns) if err != nil { return false, err } diff --git a/field.go b/field.go index df87eea91..0d401a340 100644 --- a/field.go +++ b/field.go @@ -1310,7 +1310,8 @@ func (f *Field) ClearBits(tx Tx, shard uint64, recordIDs ...uint64) error { if frag == nil { return nil } - return frag.ClearRecords(tx, recordIDs) + _, err := frag.ClearRecords(tx, recordIDs) + return err } func groupCompare(a, b string, offset int) (lt, eq bool) { diff --git a/fragment.go b/fragment.go index c02e24699..6901762af 100644 --- a/fragment.go +++ b/fragment.go @@ -2045,7 +2045,14 @@ func (f *fragment) importPositions(tx Tx, set, clear []uint64, rowSet map[uint64 } f.stats.Count(MetricClearedN, int64(changedN), 1) } + return f.updateCaching(tx, rowSet) +} +// updateCaching clears checksums for rows, and clears any existing TopN +// cache for them, and marks the cache for needing updates. I'm not sure +// that's correct. This was originally the tail end of importPositions, but +// we want to be able to access the same logic from elsewhere. +func (f *fragment) updateCaching(tx Tx, rowSet map[uint64]struct{}) error { // Update cache counts for all affected rows. for rowID := range rowSet { // Invalidate block checksum. @@ -2171,28 +2178,46 @@ func (f *fragment) bulkImportMutex(tx Tx, rowIDs, columnIDs []uint64, options *I // ClearRecords deletes all bits for the given records. It's basically // the remove-only part of setting a mutex. -func (f *fragment) ClearRecords(tx Tx, recordIDs []uint64) error { - f.mu.Lock() - defer f.mu.Unlock() - +func (f *fragment) ClearRecords(tx Tx, recordIDs []uint64) (bool, error) { // create a mask of columns we care about columns := roaring.NewSliceBitmap(recordIDs...) + return f.clearRecordsByBitmap(tx, columns) +} - // we now need to find existing rows for these bits. +func (f *fragment) clearRecordsByBitmap(tx Tx, columns *roaring.Bitmap) (changed bool, err error) { + f.mu.Lock() + defer f.mu.Unlock() + return f.unprotectedClearRecordsByBitmap(tx, columns) +} + +// clearRecordsByBitmap clears bits in a fragment that correspond to those +// positions within the bitmap. +func (f *fragment) unprotectedClearRecordsByBitmap(tx Tx, columns *roaring.Bitmap) (changed bool, err error) { rowSet := make(map[uint64]struct{}) - var toClear []uint64 - callback := func(pos uint64) error { - toClear = append(toClear, pos) - rowID := pos / ShardWidth - rowSet[rowID] = struct{}{} + rewriteExisting := roaring.NewBitmapBitmapTrimmer(columns, func(key roaring.FilterKey, data *roaring.Container, filter *roaring.Container, writeback roaring.ContainerWriteback) error { + if filter.N() == 0 { + return nil + } + existing := data.N() + // nothing to delete. this can't happen normally, but the rewriter calls + // us with an empty data container when it's done. + if existing == 0 { + return nil + } + data = data.DifferenceInPlace(filter) + if data.N() != existing { + rowSet[key.Row()] = struct{}{} + changed = true + return writeback(key, data) + } return nil - } - findExisting := roaring.NewBitmapBitmapFilter(columns, callback) - err := tx.ApplyFilter(f.index(), f.field(), f.view(), f.shard, 0, findExisting) + }) + + err = tx.ApplyRewriter(f.index(), f.field(), f.view(), f.shard, 0, rewriteExisting) if err != nil { - return errors.Wrap(err, "finding existing positions") + return false, err } - return errors.Wrap(f.importPositions(tx, nil, toClear, rowSet), "clearing records") + return changed, f.updateCaching(tx, rowSet) } // importValue bulk imports a set of range-encoded values.