From 852a8be45a25c0e1d0578be90c4e049cc1c21de5 Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 30 Mar 2022 14:22:17 -0500 Subject: [PATCH] use ApplyRewriter for ClearRecords This uses the shiny new ApplyRewriter logic for ClearRecords, mostly to verify that ApplyRewriter works at all. This also implies separating the cache update code out from importPositions so it can be used also by this. We also use fragment.ClearRecords instead of the different clearFragment code in executor. The clearFragment implementation did not update TopN caches and the like. Standardize it on the clearRecords implementation which does. --- executor.go | 41 ++++++--------------------------------- field.go | 3 ++- fragment.go | 55 ++++++++++++++++++++++++++++++++++++++--------------- 3 files changed, 48 insertions(+), 51 deletions(-) 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.