mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-09 04:17:51 +00:00
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.
This commit is contained in:
parent
74ae1fd598
commit
852a8be45a
3 changed files with 48 additions and 51 deletions
41
executor.go
41
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
|
||||
}
|
||||
|
|
|
|||
3
field.go
3
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) {
|
||||
|
|
|
|||
55
fragment.go
55
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.
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue