From 919eb3f0e847ff1ef0f5d4134a02a1442a859f33 Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 30 Mar 2022 15:05:59 -0500 Subject: [PATCH] rework bulkImportMutex to use ApplyRewriter Now that we have ApplyRewriter, it's a viable way to implement ImportMutex. It can be slower on low-density writes, because it's checking more things than it otherwise might -- the other filter form can skip ahead and only check the containers it's modfying, in principle, while this one doesn't know it can do that. (The decision as to how far to skip ahead has to be made in the BitmapBitmapTrimmer, while it's the callback provided to it that knows when it next has data to write.) On the other hand, it's probably faster in some cases, and would be more-faster if we could improve the cursor management a bit, and it's skipping at least some seeking because it doesn't need to use ImportPositions after reading the whole thing. --- fragment.go | 91 +++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 75 insertions(+), 16 deletions(-) diff --git a/fragment.go b/fragment.go index 6901762af..cff68282e 100644 --- a/fragment.go +++ b/fragment.go @@ -2155,25 +2155,84 @@ func (f *fragment) bulkImportMutex(tx Tx, rowIDs, columnIDs []uint64, options *I if unsorted { sort.Slice(toSet, func(i, j int) bool { return toSet[i] < toSet[j] }) } - // we'll reuse the row IDs as the values to clear, if any. - toClear := columnIDs[:0] - callback := func(pos uint64) error { - toClear = append(toClear, pos) - rowID := pos / ShardWidth - rowSet[rowID] = struct{}{} + + nextKey := toSet[0] >> 16 + scratchContainer := roaring.NewContainerArray([]uint16{}) + rewriteExisting := roaring.NewBitmapBitmapTrimmer(columns, func(key roaring.FilterKey, data *roaring.Container, filter *roaring.Container, writeback roaring.ContainerWriteback) error { + var inserting []uint64 + for roaring.FilterKey(nextKey) < key { + thisKey := roaring.FilterKey(nextKey) + inserting, toSet, nextKey = roaring.GetMatchingKeysFrom(toSet, nextKey) + scratchContainer = roaring.RemakeContainerFrom(scratchContainer, inserting) + + err := writeback(thisKey, scratchContainer) + if err != nil { + return err + } + } + // we only wanted to insert the data we had before the end, there's + // no actual data here to modify. + if data == nil { + return nil + } + if roaring.FilterKey(nextKey) > key { + // simple path: we only have to remove things from the filter, + // if there are any. + if filter.N() == 0 { + return nil + } + existing := data.N() + data = data.DifferenceInPlace(filter) + if data.N() != existing { + rowSet[key.Row()] = struct{}{} + return writeback(key, data) + } + return nil + } + // nextKey has to be the same as key. we have values to insert, and + // necessarily have a filter to remove which matches them. so we're + // going to remove everything in the filter, then add all the values + // we have to insert. but! in the case where a bit is already set, + // and we remove it and re-add it, we don't want to count that. + inserting, toSet, nextKey = roaring.GetMatchingKeysFrom(toSet, nextKey) + + existing := data.N() + + // so, we want to remove anything that's in the filter, *but*, if a + // thing is in the filter, and we then add it back, that doesn't + // count. but if a thing is in the filter, but *wasn't originally + // there*, that counts. But we can't check that *after* we compute + // the difference, so... + reAdds := 0 + for _, v := range inserting { + if filter.Contains(uint16(v)) && data.Contains(uint16(v)) { + reAdds++ + } + } + data = data.DifferenceInPlace(filter) + removes := int(existing - data.N()) + var changed bool + adds := 0 + for _, v := range inserting { + data, changed = data.Add(uint16(v)) + if changed { + adds++ + } + } + // if we added more things than were being readded, or removed more + // things than were being readded, we changed something. + if adds > reAdds || removes > reAdds { + rowSet[key.Row()] = struct{}{} + 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 err } - // if we're clearing things, anything being set that is being cleared - // should not be cleared - if len(toClear) > 0 { - toClear = sliceDifference(toClear, toSet) - } - return errors.Wrap(f.importPositions(tx, toSet, toClear, rowSet), "importing positions") + return f.updateCaching(tx, rowSet) } // ClearRecords deletes all bits for the given records. It's basically