From 59d89dda99c659700245ee6ec4b6e538d6886c06 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 7 Dec 2020 16:11:07 -0600 Subject: [PATCH] Allow arbitrary and potentially more efficient filtering of bitmaps This is a partial solution to a nasty performance problem, which is that a ContainerIterator has to *generate* all the containers. With roaring, this was cheap because they already exist in memory; with transactional backends, it's an allocation per container, *even for the containers we don't use*. This design admits filters which can distinguish between answers they can give just based on keys and times when they actually need containers instantiated, and can also give hints as to future answers -- saying "yes" or "no" to entire rows at a time, or indicating when they're done. This is only part of the solution; we also need a Tx API hook for doing scans like this which doesn't rely on ContainerIterator. --- executor.go | 12 +- executor_internal_test.go | 73 ---- fragment.go | 371 ++++++++--------- fragment_internal_test.go | 294 ++++++++++++- roaring/filter.go | 714 ++++++++++++++++++++++++++++++++ roaring/filter_internal_test.go | 226 ++++++++++ roaring/printutil.go | 7 - roaring/roaring.go | 42 ++ 8 files changed, 1461 insertions(+), 278 deletions(-) create mode 100644 roaring/filter.go create mode 100644 roaring/filter_internal_test.go diff --git a/executor.go b/executor.go index c4866d6d9..ae52457c7 100644 --- a/executor.go +++ b/executor.go @@ -3197,7 +3197,7 @@ func (e *executor) executeRowsShard(ctx context.Context, qcx *Qcx, index string, start = previous + 1 } - filters := []rowFilter{} + filters := []roaring.BitmapFilter{} if columnID, ok, err := c.UintArg("column"); err != nil { return nil, err } else if ok { @@ -3205,14 +3205,14 @@ func (e *executor) executeRowsShard(ctx context.Context, qcx *Qcx, index string, if colShard != shard { return rowIDs, nil } - filters = append(filters, filterColumn(columnID)) + filters = append(filters, roaring.NewBitmapColumnFilter(columnID)) } limit := int(^uint(0) >> 1) if lim, hasLimit, err := c.UintArg("limit"); err != nil { return nil, errors.Wrap(err, "getting limit") } else if hasLimit { - filters = append(filters, filterWithLimit(lim)) + filters = append(filters, roaring.NewBitmapRowLimitFilter(lim)) limit = int(lim) } @@ -3221,7 +3221,7 @@ func (e *executor) executeRowsShard(ctx context.Context, qcx *Qcx, index string, return nil, errors.Wrap(err, "getting like pattern") } else if hasLike { likeErr = make(chan error, 1) - filters = append(filters, filterLike(like, f.TranslateStore(), likeErr)) + filters = append(filters, NewBitmapLikeFilter(like, f.TranslateStore())) } tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard}) @@ -6894,9 +6894,9 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children if frag == nil { // this means this whole shard doesn't have all it needs to continue return nil, nil } - filters := []rowFilter{} + filters := []roaring.BitmapFilter{} if len(rowIDs[i]) > 0 { - filters = append(filters, filterWithRows(rowIDs[i])) + filters = append(filters, roaring.NewBitmapRowsFilter(rowIDs[i])) } tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard}) diff --git a/executor_internal_test.go b/executor_internal_test.go index 6ee76ca3e..06e271fde 100644 --- a/executor_internal_test.go +++ b/executor_internal_test.go @@ -92,79 +92,6 @@ func TestExecutor_TranslateRowsOnBool(t *testing.T) { } } -func TestFilterWithLimit(t *testing.T) { - f := filterWithLimit(5) - - for i := uint64(0); i < 5; i++ { - include, done := f(i, i*(1<= len(toClear) { + return toClear[:n] + } + cv = toClear[cn] + } + } + return toClear[:n] +} + // bulkImportMutex performs a bulk import on a fragment while ensuring // mutex restrictions. Because the mutex requirements must be checked // against storage, this method must acquire a write lock on the fragment @@ -2432,53 +2514,64 @@ func (f *fragment) bulkImportMutex(tx Tx, rowIDs, columnIDs []uint64) error { f.mu.Lock() defer f.mu.Unlock() - rowSet := make(map[uint64]struct{}) - // we have to maintain which columns are getting bits set as a map so that - // we don't end up setting multiple bits in the same column if a column is - // repeated within the import. - colSet := make(map[uint64]uint64) + // We don't use this right away, but if this fails nothing else is + // useful... + iter, _, err := tx.ContainerIterator(f.index, f.field, f.view, f.shard, 0) + if err != nil { + return errors.Wrap(err, "searching storage") + } - // Since each imported bit will at most set one bit and clear one bit, we - // can reuse the rowIDs and columnIDs slices as the set and clear slice - // arguments to importPositions. The set positions we'll get from the - // colSet, but we maintain clearIdx as we loop through row and col ids so - // that we know how many bits we need to clear and how far through columnIDs - // we are. - clearIdx := 0 + p := parallelSlices{cols: columnIDs, rows: rowIDs} + p.FullPrune() + columnIDs = p.cols + rowIDs = p.rows + + // create a mask of rows we care about + columns := roaring.NewSliceBitmap() + _, err = columns.AddN(columnIDs...) + if err != nil { + return errors.Wrap(err, "creating bitmap of affected columns") + } + // we now need to find existing rows for these bits. + rowSet := make(map[uint64]struct{}, len(rowIDs)) + unsorted := false + prev := uint64(0) for i := range rowIDs { rowID, columnID := rowIDs[i], columnIDs[i] - if existingRowID, found, err := f.mutexVector.Get(tx, columnID); err != nil { - return errors.Wrap(err, "getting mutex vector data") - } else if found && existingRowID != rowID { - // Determine the position of the bit in the storage. - clearPos, err := f.pos(existingRowID, columnID) - if err != nil { - return err - } - columnIDs[clearIdx] = clearPos - clearIdx++ - - rowSet[existingRowID] = struct{}{} - } else if found && existingRowID == rowID { - continue - } + rowSet[rowID] = struct{}{} pos, err := f.pos(rowID, columnID) if err != nil { return err } - colSet[columnID] = pos - rowSet[rowID] = struct{}{} - } - - // re-use rowIDs by populating positions to set from colSet. - i := 0 - for _, pos := range colSet { + // positions are sorted by columns, but not by absolute + // position. we might want them sorted, though. + if pos < prev { + unsorted = true + } rowIDs[i] = pos - i++ } - toSet := rowIDs[:i] - toClear := columnIDs[:clearIdx] - + toSet := rowIDs + 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{}{} + return nil + } + existing := roaring.NewBitmapBitmapFilter(columns, callback) + err = roaring.ApplyFilterToIterator(existing, iter) + if err != nil { + return errors.Wrap(err, "finding existing positions") + } + // if we're clearing things, anything being set that is being cleared + // should not be cleared + if len(toClear) > 0 { + toClear = unclearSets(toSet, toClear) + } return errors.Wrap(f.importPositions(tx, toSet, toClear, rowSet), "importing positions") } @@ -3032,77 +3125,39 @@ func (f *fragment) minRowID(tx Tx) (uint64, bool, error) { return min / ShardWidth, ok, err } -// rowFilter is a function signature for controlling iteration over containers -// in a fragment. It will be invoked on each container found and returns two -// booleans. The first is whether the row this container is in should be -// included or skipped, and the second is whether to stop processing or -// continue. -type rowFilter func(rowID, key uint64, c *roaring.Container) (include, done bool) - -// filterWithLimit returns a filter which will only allow a limited number of -// rows to be returned. It should be applied last so that it is only called (and -// therefore only updates its internal state) if the row is being included by -// every other filter. -func filterWithLimit(limit uint64) rowFilter { - return func(rowID, key uint64, c *roaring.Container) (include, done bool) { - if limit > 0 { - limit-- - return true, false - } - return false, true - } +// BitmapLikeFilter is a roaring.BitmapFilter which handles Like expressions. +type BitmapLikeFilter struct { + roaring.BitmapRowFilterBase + plan []filterStep + translator TranslateStore } -func filterColumn(col uint64) rowFilter { - return func(rowID, key uint64, c *roaring.Container) (include, done bool) { - colID := col % ShardWidth - colKey := ((rowID * ShardWidth) + colID) >> 16 - colVal := uint16(colID & 0xFFFF) // columnID within the container - return colKey == key && c.Contains(colVal), false +func (b *BitmapLikeFilter) ConsiderKey(key roaring.FilterKey) roaring.FilterResult { + res, done := b.DetermineByKey(key) + if done { + return res } + row := key.Row() + keyStr, err := b.translator.TranslateID(row) + if err != nil { + return b.SetResult(key, key.Fail(errors.Wrap(err, "translating key for row"))) + } + if matchLike(keyStr, b.plan...) { + return b.SetResult(key, key.MatchRow()) + } + return b.SetResult(key, key.RejectRow()) } -func filterLike(like string, t TranslateStore, e chan error) rowFilter { - plan := planLike(like) - - return func(rowID, key uint64, c *roaring.Container) (include, done bool) { - keyStr, err := t.TranslateID(rowID) - if err != nil { - select { - case e <- err: - default: - } - return false, true - } - return matchLike(keyStr, plan...), false - } +func (b *BitmapLikeFilter) ConsiderData(key roaring.FilterKey, data *roaring.Container) roaring.FilterResult { + b.FilterResult.Err = errors.New("like filter should not need to look at data") + return b.FilterResult } -// TODO: this works, but it would be more performant if the fragment could seek -// to the next row in the rows list rather than asking the filter for each -// container serially. The container iterator would need to expose a seek -// method, and the rowFilter would need some way of communicating to -// fragment.rows what the next rowID to seek to is. -func filterWithRows(rows []uint64) rowFilter { - loc := 0 - return func(rowID, key uint64, c *roaring.Container) (include, done bool) { - if loc >= len(rows) { - return false, true - } - i := sort.Search(len(rows[loc:]), func(i int) bool { - return rows[loc+i] >= rowID - }) - loc += i - if loc >= len(rows) { - return false, true - } - if rows[loc] == rowID { - if loc == len(rows)-1 { - done = true - } - return true, done - } - return false, false +func NewBitmapLikeFilter(like string, translator TranslateStore) *BitmapLikeFilter { + return &BitmapLikeFilter{ + BitmapRowFilterBase: *roaring.NewBitmapRowFilterBase(nil), + plan: planLike(like), + translator: translator, } } @@ -3115,15 +3170,15 @@ func filterWithRows(rows []uint64) rowFilter { // returning done == true will cause processing to stop after all filters for // this container have been processed. The rows accumulated up to this point // (including this row if all filters passed) will be returned. -func (f *fragment) rows(ctx context.Context, tx Tx, start uint64, filters ...rowFilter) ([]uint64, error) { +func (f *fragment) rows(ctx context.Context, tx Tx, start uint64, filters ...roaring.BitmapFilter) ([]uint64, error) { f.mu.RLock() defer f.mu.RUnlock() return f.unprotectedRows(ctx, tx, start, filters...) } // unprotectedRows calls rows without grabbing the mutex. -func (f *fragment) unprotectedRows(ctx context.Context, tx Tx, start uint64, filters ...rowFilter) ([]uint64, error) { - rows := make([]uint64, 0) +func (f *fragment) unprotectedRows(ctx context.Context, tx Tx, start uint64, filters ...roaring.BitmapFilter) ([]uint64, error) { + var rows []uint64 startKey := rowToKey(start) i, _, err := tx.ContainerIterator(f.index(), f.field(), f.view(), f.shard, startKey) if err != nil { @@ -3131,49 +3186,13 @@ func (f *fragment) unprotectedRows(ctx context.Context, tx Tx, start uint64, fil } else if i == nil { return rows, nil } - defer i.Close() // must close iterators allocated on a Tx - var lastRow uint64 = math.MaxUint64 - - // Loop over the existing containers. - var k uint16 - for i.Next() { - if k == 0 { - if err := ctx.Err(); err != nil { - // caller doesn't need a result anymore. - return nil, err - } - } - k++ - - key, c := i.Value() - - // virtual row for the current container - vRow := key >> shardVsContainerExponent - - // skip dups - if vRow == lastRow { - continue - } - - // apply filters - addRow, done := true, false - for _, filter := range filters { - var d bool - addRow, d = filter(vRow, key, c) - done = done || d - if !addRow { - break - } - } - if addRow { - lastRow = vRow - rows = append(rows, vRow) - } - if done { - return rows, nil - } + callback := func(row uint64) error { + rows = append(rows, row) + return nil } - return rows, nil + filter := roaring.NewBitmapRowFilter(callback, filters...) + err = roaring.ApplyFilterToIterator(filter, i) + return rows, err } // blockToRoaringData converts a fragment block into a roaring.Bitmap @@ -3246,7 +3265,7 @@ type rowIterator interface { Next() (*Row, uint64, *int64, bool, error) } -func (f *fragment) rowIterator(tx Tx, wrap bool, filters ...rowFilter) (rowIterator, error) { +func (f *fragment) rowIterator(tx Tx, wrap bool, filters ...roaring.BitmapFilter) (rowIterator, error) { if strings.HasPrefix(f.view(), viewBSIGroupPrefix) { return f.intRowIterator(tx, wrap, filters...) } @@ -3264,7 +3283,7 @@ type intRowIterator struct { wrap bool } -func (f *fragment) intRowIterator(tx Tx, wrap bool, filters ...rowFilter) (rowIterator, error) { +func (f *fragment) intRowIterator(tx Tx, wrap bool, filters ...roaring.BitmapFilter) (rowIterator, error) { it := intRowIterator{ f: f, colIDs: make(map[int64][]uint64), @@ -3283,7 +3302,7 @@ func (f *fragment) intRowIterator(tx Tx, wrap bool, filters ...rowFilter) (rowIt f.mu.RLock() defer f.mu.RUnlock() } - if err := f.foreachRow(tx, filters, func(rid uint64) error { + callback := func(rid uint64) error { // skip exist(0) and sign(1) rows if rid == bsiExistsBit || rid == bsiSignBit { return nil @@ -3298,7 +3317,8 @@ func (f *fragment) intRowIterator(tx Tx, wrap bool, filters ...rowFilter) (rowIt acc[cid] |= val } return nil - }); err != nil { + } + if err := f.foreachRow(tx, filters, callback); err != nil { return nil, err } @@ -3341,47 +3361,20 @@ func (f *fragment) intRowIterator(tx Tx, wrap bool, filters ...rowFilter) (rowIt return &it, nil } +<<<<<<< HEAD func (f *fragment) foreachRow(tx Tx, filters []rowFilter, fn func(rid uint64) error) error { var lastRow uint64 = math.MaxUint64 i, _, err := tx.ContainerIterator(f.index(), f.field(), f.view(), f.shard, rowToKey(0)) +======= +func (f *fragment) foreachRow(tx Tx, filters []roaring.BitmapFilter, fn func(rid uint64) error) error { + i, _, err := tx.ContainerIterator(f.index, f.field, f.view, f.shard, rowToKey(0)) +>>>>>>> Allow arbitrary and potentially more efficient filtering of bitmaps if err != nil { return err } - defer i.Close() // must close tx allocated iterators when done. - - // Loop over the existing containers. - for i.Next() { - key, c := i.Value() - // virtual row for the current container - vRow := key >> shardVsContainerExponent - // skip dups - if vRow == lastRow { - continue - } - - // apply filters - addRow, done := true, false - for _, filter := range filters { - var d bool - addRow, d = filter(vRow, key, c) - done = done || d - if !addRow { - break - } - } - if addRow { - lastRow = vRow - if fn != nil { - if err := fn(vRow); err != nil { - return err - } - } - } - if done { - break - } - } - return nil + filter := roaring.NewBitmapRowFilter(fn, filters...) + err = roaring.ApplyFilterToIterator(filter, i) + return err } func (it *intRowIterator) Seek(rowID uint64) { @@ -3416,7 +3409,7 @@ type setRowIterator struct { wrap bool } -func (f *fragment) setRowIterator(tx Tx, wrap bool, filters ...rowFilter) (rowIterator, error) { +func (f *fragment) setRowIterator(tx Tx, wrap bool, filters ...roaring.BitmapFilter) (rowIterator, error) { rows, err := f.rows(context.Background(), tx, 0, filters...) if err != nil { return nil, err @@ -3867,7 +3860,7 @@ func newRowsVector(f *fragment) *rowsVector { // otherwise it returns false. Ensure that you already // have the mutex before calling this. func (v *rowsVector) Get(tx Tx, colID uint64) (uint64, bool, error) { - rows, err := v.f.unprotectedRows(context.Background(), tx, 0, filterColumn(colID)) + rows, err := v.f.unprotectedRows(context.Background(), tx, 0, roaring.NewBitmapColumnFilter(colID)) if err != nil { return 0, false, err } else if len(rows) > 1 { @@ -3903,7 +3896,7 @@ func newBoolVector(f *fragment) *boolVector { // otherwise it returns false. Ensure that you already // have the fragment mutex before calling this. func (v *boolVector) Get(tx Tx, colID uint64) (uint64, bool, error) { - rows, err := v.f.unprotectedRows(context.Background(), tx, 0, filterColumn(colID)) + rows, err := v.f.unprotectedRows(context.Background(), tx, 0, roaring.NewBitmapColumnFilter(colID)) if err != nil { return 0, false, err } else if len(rows) > 1 { diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 970596705..e46c6f61d 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -29,7 +29,9 @@ import ( "runtime" "runtime/debug" "sort" + "strconv" "strings" + "sync" "testing" "testing/quick" @@ -3599,7 +3601,7 @@ func TestFragment_RowsIteration(t *testing.T) { t.Fatalf("Do not match %v %v", expectedAll, ids) } - ids, err = f.rows(context.Background(), tx, 0, filterColumn(1)) + ids, err = f.rows(context.Background(), tx, 0, roaring.NewBitmapColumnFilter(1)) if err != nil { t.Fatal(err) } else if !reflect.DeepEqual(expectedOdd, ids) { @@ -3633,7 +3635,7 @@ func TestFragment_RowsIteration(t *testing.T) { t.Fatalf("Do not match %v %v", expected, ids) } - ids, err = f.rows(context.Background(), tx, 0, filterColumn(66000)) + ids, err = f.rows(context.Background(), tx, 0, roaring.NewBitmapColumnFilter(66000)) if err != nil { t.Fatal(err) } else if !reflect.DeepEqual(expected, ids) { @@ -3661,7 +3663,7 @@ func TestFragment_RowsIteration(t *testing.T) { } else if !reflect.DeepEqual(expectedRows, ids) { t.Fatalf("Do not match %v %v", expectedRows, ids) } - ids, err = f.rows(context.Background(), tx, 0, filterColumn(c)) + ids, err = f.rows(context.Background(), tx, 0, roaring.NewBitmapColumnFilter(c)) if err != nil { t.Fatal(err) } else if !reflect.DeepEqual(expectedRows, ids) { @@ -5432,3 +5434,289 @@ func notBlueGreenTest(t *testing.T) { t.Skip("skip under blue green") } } + +var mutexSamplesPrepared sync.Once + +func requireMutexSampleData() { + mutexSamplesPrepared.Do(prepareMutexSampleData) +} + +// a few mutex tests want common largeish pools of mutex data +type mutexSampleData struct { + name string + colIDs, rowIDs []uint64 +} + +func (m *mutexSampleData) scratchSpace(cols, rows []uint64) ([]uint64, []uint64) { + if cap(cols) < len(m.colIDs) { + cols = make([]uint64, len(m.colIDs)) + } else { + cols = cols[:len(m.colIDs)] + } + copy(cols, m.colIDs) + if cap(rows) < len(m.rowIDs) { + rows = make([]uint64, len(m.rowIDs)) + } else { + rows = rows[:len(m.rowIDs)] + } + copy(rows, m.rowIDs) + return cols, rows +} + +// mutex data wants to exist for differing numbers of rows, and different +// densities, and reasonably-large N. But we don't want a huge pool of nested +// maps, so... +type mutexSampleRange int16 + +// density should be able to range from every bit filled to almost no +// bits filled. So, how many bits per container? How about pow(2, N), where N +// can be negative. 16 is the highest possible value, and -16 is the lowest, +// and 0 means about one bit per container on average. +func (m mutexSampleRange) density() int8 { + return int8(m >> 8) +} + +// bottom 8 bits are log2 of number of rows. +func (m mutexSampleRange) rows() uint32 { + return uint32(1) << (m & 0x1F) +} + +// just a convenience thing to hide the implementation +func newMutexSampleRange(density int8, rows uint8) mutexSampleRange { + if density > 16 { + density = 16 + } + if density < -16 { + density = -16 + } + return mutexSampleRange((int16(density) << 8) | int16(rows)) +} + +var sampleMutexData = map[mutexSampleRange]*mutexSampleData{} + +type mutexDensity struct { + name string + density int8 +} + +type mutexSize struct { + name string + rows uint8 +} + +var mutexDensities = []mutexDensity{ + {"64K", 16}, + // {"32K", 15}, // 50-50 + // {"16K", 14}, // 1/4 + // {"4K", 12}, // a fair number of things + // {"1", 0}, // about one per container + // {"empty", -14}, // almost none +} + +var mutexSizes = []mutexSize{ + {"4r", 2}, + {"16r", 4}, + {"256r", 8}, + // {"65Kr", 16}, +} + +var mutexCaches = []string{ + // "ranked", + "none", +} + +const mutexSampleDataSize = ShardWidth + +func prepareMutexSampleData() { + for _, d := range mutexDensities { + for _, s := range mutexSizes { + rng := newMutexSampleRange(d.density, s.rows) + col := uint64(0) + // at density 16, we want everything to be adjacent. + // at density 0, we want about 65k between items. + // The average spacing we want is 1<<(16 - density), + // so random numbers between 0 and twice that would + // be close, but we never want 0, so, subtract 1 from + // "twice that", then add 1 to the result. + // + // So for density 16, we compute spacing of 1, then + // draw random numbers in [0,1), and add 1 to them. + spacing := ((int64(1) << (16 - d.density)) * 2) - 1 + rows := (int64(1) << s.rows) + expected := int64(mutexSampleDataSize) + if (ShardWidth / spacing) < mutexSampleDataSize { + expected = (ShardWidth / spacing) * 2 + if expected < 2 { + expected = 2 + } + } + data := &mutexSampleData{ + name: d.name + "/" + s.name, + colIDs: make([]uint64, mutexSampleDataSize), + rowIDs: make([]uint64, mutexSampleDataSize), + } + for i := int64(0); i < expected; i++ { + col += uint64(rand.Int63n(spacing)) + 1 + // can only import one fragment at a time, + // though! + if col >= ShardWidth { + expected = i + break + } + row := uint64(rand.Int63n(rows)) + data.colIDs[i] = col + data.rowIDs[i] = row + } + // if we came up short, because we overran the shard + // size, we reduced expected above. + data.colIDs = data.colIDs[:expected] + data.rowIDs = data.rowIDs[:expected] + sampleMutexData[rng] = data + } + } +} + +var importBatchSizes = []int{65536} + +func TestImportMutexSampleData(t *testing.T) { + requireMutexSampleData() + var scratchCols []uint64 + var scratchRows []uint64 + for rng, data := range sampleMutexData { + // skip the larger ones, they'll be slow + if rng.rows() > 256 { + continue + } + scratchCols, scratchRows = data.scratchSpace(scratchCols, scratchRows) + t.Run(data.name, func(t *testing.T) { + for _, batchSize := range importBatchSizes { + t.Run(fmt.Sprintf("%d", batchSize), func(t *testing.T) { + f, _, tx := mustOpenMutexFragment(t, "i", "f", viewStandard, 0, "") + defer f.Clean(t) + // Set import. + var err error + for i := 0; i < len(scratchCols); i += batchSize { + max := i + batchSize + if len(scratchCols) < max { + max = len(scratchCols) + } + err = f.bulkImport(tx, scratchRows[i:max], scratchCols[i:max], &ImportOptions{}) + if err != nil { + t.Fatalf("bulk importing ids [%d:%d]: %v", i, max, err) + } + } + count := uint64(0) + for k := uint32(0); k < rng.rows(); k++ { + count += f.mustRow(tx, uint64(k)).Count() + } + if int(count) != len(data.colIDs) { + t.Fatalf("for %d rows, %d density: expected %d results, got %d", + rng.rows(), rng.density(), len(data.colIDs), count) + } + }) + } + }) + } +} + +func BenchmarkImportMutexSampleData(b *testing.B) { + requireMutexSampleData() + var scratchCols []uint64 + var scratchRows []uint64 + var data *mutexSampleData + var cache string + var batchSize int + // this exists just to manage indentation below + benchmarkFragmentImports := func(b *testing.B) { + f, _, tx := mustOpenMutexFragment(b, "i", "f", viewStandard, 0, cache) + defer f.Clean(b) + scratchCols, scratchRows = data.scratchSpace(scratchCols, scratchRows) + var err error + for i := 0; i < len(scratchCols) && i < (batchSize*b.N); i += batchSize { + max := i + batchSize + if len(scratchCols) < max { + max = len(scratchCols) + } + err = f.bulkImport(tx, scratchRows[i:max], scratchCols[i:max], &ImportOptions{}) + if err != nil { + b.Fatalf("bulk importing ids [%d:%d]: %v", i, max, err) + } + } + } + for _, data = range sampleMutexData { + b.Run(data.name, func(b *testing.B) { + for _, batchSize = range importBatchSizes { + b.Run(strconv.Itoa(batchSize), func(b *testing.B) { + for _, cache = range mutexCaches { + b.Run(cache, benchmarkFragmentImports) + } + }) + } + }) + } +} + +func testOneParallelSlice(t *testing.T, p *parallelSlices) { + // the easy answer + seen := make(map[uint64]uint64, len(p.cols)) + t.Logf("cols %d, rows %d", p.cols, p.rows) + for i, c := range p.cols { + seen[c] = p.rows[i] + } + unsorted := p.Prune() + t.Logf("unsorted %t, cols %d, rows %d", unsorted, p.cols, p.rows) + if unsorted { + sort.Stable(p) + } + unsorted = p.Prune() + if unsorted { + t.Fatalf("slice still unsorted after sort") + } + t.Logf("pruned/sorted cols %d, rows %d", p.cols, p.rows) + if len(p.cols) != len(seen) { + t.Fatalf("expected %d entries, found %d", len(seen), len(p.cols)) + } + for i, c := range p.cols { + if seen[c] != p.rows[i] { + t.Fatalf("expected %d:%d, found :%d", c, seen[c], p.rows[i]) + } + } +} + +func TestParallelSlices(t *testing.T) { + cols := make([]uint64, 256) + rows := make([]uint64, 256) + // ensure at least some overlap by coercing columns into a range + // smaller than number of entries + for i := range cols { + cols[i] = rand.Uint64() & ((uint64(len(cols)) / 2) - 1) + rows[i] = rand.Uint64() & 0xff + } + t.Run("random", func(t *testing.T) { + testOneParallelSlice(t, ¶llelSlices{cols: cols, rows: rows}) + }) + cols = cols[:cap(cols)] + rows = rows[:cap(rows)] + // in-order but no overlap + col := uint64(0) + for i := range cols { + cols[i] = col + col = col + (rand.Uint64() & 3) + 1 + rows[i] = rand.Uint64() & 0xff + } + t.Run("ordered", func(t *testing.T) { + testOneParallelSlice(t, ¶llelSlices{cols: cols, rows: rows}) + }) + cols = cols[:cap(cols)] + rows = rows[:cap(rows)] + // in-order with + col = uint64(0) + for i := range cols { + cols[i] = col + col = col + (rand.Uint64() & 3) + rows[i] = rand.Uint64() & 0xff + } + t.Run("orderlapping", func(t *testing.T) { + testOneParallelSlice(t, ¶llelSlices{cols: cols, rows: rows}) + }) +} diff --git a/roaring/filter.go b/roaring/filter.go new file mode 100644 index 000000000..2f826ab36 --- /dev/null +++ b/roaring/filter.go @@ -0,0 +1,714 @@ +// Copyright 2020 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package roaring + +import ( + "errors" + "fmt" + + "github.com/pilosa/pilosa/v2/shardwidth" +) + +// We want BitmapScanner to be accessible from both the pilosa package, and +// the rbf package. Pilosa imports rbf, so rbf can't import pilosa, but they +// both import roaring, and this package is closely tied to roaring structures +// like Containers and the key/container mapping, so it mostly makes sense for +// this to be here. +// +// Unfortunately, this really needs to be capable of being row-aware, which +// means it needs access to the shard width stuff, which roaring otherwise +// studiously avoids knowing about. +const ( + rowExponent = (shardwidth.Exponent - 16) + rowWidth = 1 << rowExponent // containers per row + keyMask = (rowWidth - 1) // a mask for offset within the row + rowMask = ^FilterKey(keyMask) // a mask for the row bits, without converting them to a row ID +) + +type FilterKey uint64 + +// FilterResult represents the results of a BitmapFilter considering a +// key, or data. The values are represented as exclusive upper bounds on +// a series of matches followed by a series of rejections. So for instance, +// if called on key 23, the result {YesKey: 23, NoKey: 24} indicates that +// key 23 is a "no". This may seem confusing but it makes the math a lot +// easier to write. It can also report an error, which indicates that the +// entire operation should be stopped with that error. +type FilterResult struct { + YesKey FilterKey // The lowest container key this container is known NOT to match. + NoKey FilterKey // The highest container key after that that this filter is known to not match. + Err error // An error which should terminate processing. +} + +// Row() computes the row number of a key. +func (f FilterKey) Row() uint64 { + return uint64(f >> rowExponent) +} + +// Add adds an offset to a key. +func (f FilterKey) Add(x uint64) FilterKey { + return f + FilterKey(x) +} + +// Sub determines the distance from o to f. +func (f FilterKey) Sub(o FilterKey) uint64 { + return uint64(f - o) +} + +// MatchReject just sets Yes and No appropriately. +func (f FilterKey) MatchReject(y, n FilterKey) FilterResult { + return FilterResult{YesKey: y, NoKey: n} +} + +func (f FilterKey) MatchOne() FilterResult { + return FilterResult{YesKey: f + 1} +} + +// NeedData() is only really meaningful for ConsiderKey, and indicates +// that a decision can't be made from the key alone. +func (f FilterKey) NeedData() FilterResult { + return FilterResult{} +} + +// Fail() reports a fatal error that should terminate processing. +func (f FilterKey) Fail(err error) FilterResult { + return FilterResult{Err: err} +} + +// Failf() is just like Errorf, etc +func (f FilterKey) Failf(msg string, args ...interface{}) FilterResult { + return FilterResult{Err: fmt.Errorf(msg, args...)} +} + +// MatchRow indicates that the current row matches the filter. +func (f FilterKey) MatchRow() FilterResult { + return FilterResult{YesKey: (f & rowMask) + rowWidth} +} + +// MatchOneRejectRow indicates that this item matched but no further +// items in this row can match. +func (f FilterKey) MatchOneRejectRow() FilterResult { + return FilterResult{YesKey: f + 1, NoKey: (f & rowMask) + rowWidth} +} + +// Reject rejects this item only. +func (f FilterKey) RejectOne() FilterResult { + return FilterResult{NoKey: f + 1} +} + +// Reject rejects N items. +func (f FilterKey) Reject(n uint64) FilterResult { + return FilterResult{NoKey: f.Add(n)} +} + +// RejectRow indicates that this entire row is rejected. +func (f FilterKey) RejectRow() FilterResult { + return FilterResult{NoKey: (f & rowMask) + rowWidth} +} + +// RejectUntil rejects everything up to the given key. +func (f FilterKey) RejectUntil(until FilterKey) FilterResult { + return FilterResult{NoKey: until} +} + +// RejectUntilRow rejects everything until the given row ID. +func (f FilterKey) RejectUntilRow(rowID uint64) FilterResult { + return FilterResult{NoKey: FilterKey(rowID) << rowExponent} +} + +// MatchRowUntilRow matches this row, then rejects everything else until +// the given row ID. +func (f FilterKey) MatchRowUntilRow(rowID uint64) FilterResult { + // if rows are 16 wide, "yes" will be 16 minus our current position + // within a row, and "no" will be the distance from the end of our + // current row to the start of rowID, which is also the distance from + // the beginning of our current row to the start of rowID-1. + return FilterResult{ + YesKey: (f & rowMask) + rowWidth, + NoKey: FilterKey(rowID) << rowExponent, + } +} + +// RejectUntilOffset rejects this container, and any others until the given +// in-row offset. +func (f FilterKey) RejectUntilOffset(offset uint64) FilterResult { + next := (f & rowMask).Add(offset) + if next <= f { + next += rowWidth + } + return FilterResult{NoKey: next} +} + +// MatchUntilOffset matches the current container, then skips any other +// containers until the given offset. +func (f FilterKey) MatchOneUntilOffset(offset uint64) FilterResult { + r := f.RejectUntilOffset(offset) + r.YesKey = f + 1 + return r +} + +// Done indicates that nothing can ever match. +func (f FilterKey) Done() FilterResult { + return FilterResult{ + NoKey: ^FilterKey(0), + } +} + +// MatchRowAndDone matches this row and nothing after that. +func (f FilterKey) MatchRowAndDone() FilterResult { + return FilterResult{ + YesKey: (f & rowMask) + rowWidth, + NoKey: ^FilterKey(0), + } +} + +// Match the current container, then skip any others until the same offset +// is reached again. +func (f FilterKey) MatchOneUntilSameOffset() FilterResult { + return f.MatchOneUntilOffset(uint64(f) & keyMask) +} + +// A BitmapFilter, given a series of key/data pairs, is considered to "match" +// some of those containers. Matching may be dependent on key values alone, +// or on the contents of the container. +// +// The ConsiderData function must not retain the container, or the data +// from the container; if it needs access to that information later, it needs +// to make a copy. +// +// Many filters are, by virtue of how they operate, able to predict their +// results on future keys. To accommodate this, and allow operations to +// avoid processing keys they don't need to process, the result of a filter +// operation can indicate not just whether a given key matches, but whether +// some upcoming keys will, or won't, match. If ConsiderKey yields a non-zero +// number of matches or non-matches for a given key, ConsiderData will not be +// called for that key. +// +// If multiple filters are combined, they are only called if their input is +// needed to determine a value. +type BitmapFilter interface { + ConsiderKey(key FilterKey) FilterResult + ConsiderData(key FilterKey, data *Container) FilterResult +} + +// BitmapColumnFilter is a BitmapFilter which checks for containers matching +// a given column within a row; thus, only the one container per row which +// matches the column needs to be evaluated, and it's evaluated as matching +// if it contains the relevant bit. +type BitmapColumnFilter struct { + key, offset uint16 +} + +var _ BitmapFilter = &BitmapColumnFilter{} + +func (f *BitmapColumnFilter) ConsiderKey(key FilterKey) FilterResult { + if uint16(key&keyMask) != f.key { + return key.RejectUntilOffset(uint64(f.key)) + } + return key.NeedData() +} + +func (f *BitmapColumnFilter) ConsiderData(key FilterKey, data *Container) FilterResult { + if data.Contains(f.offset) { + return key.MatchOneUntilSameOffset() + } + return key.RejectUntilOffset(uint64(f.key)) +} + +// BitmapRowsFilter is a BitmapFilter which checks for containers that are +// in any of a provided list of rows. The row list should be sorted. +type BitmapRowsFilter struct { + rows []uint64 + i int +} + +func (f *BitmapRowsFilter) ConsiderKey(key FilterKey) FilterResult { + if f.i == -1 { + return key.Done() + } + row := uint64(key) >> rowExponent + for f.rows[f.i] < row { + f.i++ + if f.i >= len(f.rows) { + f.i = -1 + return key.Done() + } + } + if f.rows[f.i] > row { + return key.RejectUntilRow(f.rows[f.i]) + } + // rows[f.i] must be equal, so we should match this row, until the + // next row, if there is a next row. + if f.i+1 < len(f.rows) { + return key.MatchRowUntilRow(f.rows[f.i+1]) + } + return key.MatchRowAndDone() +} + +func (f *BitmapRowsFilter) ConsiderData(key FilterKey, data *Container) FilterResult { + return key.Fail(errors.New("bitmap rows filter should never consider data")) +} + +func NewBitmapRowsFilter(rows []uint64) BitmapFilter { + if len(rows) == 0 { + return &BitmapRowsFilter{rows: rows, i: -1} + } + return &BitmapRowsFilter{rows: rows, i: 0} +} + +// BitmapRowFilterBase is a generic form of a row-aware wrapper; it +// handles making decisions about keys once you tell it a yesKey and noKey +// that it should be using, and makes callbacks per row. +type BitmapRowFilterBase struct { + FilterResult + callback func(row uint64) error + lastRow uint64 +} + +var _ BitmapFilter = &BitmapRowFilterBase{} + +// DetermineByKey decides whether it can produce a meaningful FilterResult +// for a given key. This encapsulates all the logic for row callbacks and +// figuring out when to wrap a row. +func (b *BitmapRowFilterBase) DetermineByKey(key FilterKey) (FilterResult, bool) { + if b.FilterResult.Err != nil { + return b.FilterResult, true + } + row := key.Row() + if b.YesKey <= key && b.NoKey > key { + return key.RejectUntil(b.NoKey), true + } + if b.lastRow == row { + return key.RejectRow(), true + } + // If we got here: Either b.noKey is less than key, or b.yesKey is + // greater than key. If yesKey is greater, we match this row, and + // possibly update to mark that we've said no through to the end + // of this row. + if b.YesKey > key { + b.lastRow = row + if b.callback != nil { + err := b.callback(row) + if err != nil { + return key.Fail(err), true + } + } + res := key.MatchOneRejectRow() + // This is probably unnecessary, but the idea is, since + // we've decided that we're rejecting everything up to the + // end of this row, we want to be sure that a later call + // doesn't produce a different answer. + if b.NoKey < res.NoKey { + b.NoKey = res.NoKey + } + // if our run of yes answers ends before the rejected row + // ends, and our run of no answers extends beyond this row, + // we can reject until then. note that we can't round that + // up to a full row; if our inner filter were a column + // filter, for instance, that only wanted to see the 7th + // key in each row, we would want to reject up to that 7th + // key, but then look at it. + if b.YesKey <= res.NoKey && b.NoKey > res.NoKey { + res.NoKey = b.NoKey + } + return res, true + } + // Both keys are <= key, err is nil, so this is basically a + // NeedData. + return b.FilterResult, false +} + +// SetResult is a convenience function so that things embedding this +// can just call this instead of using a long series of dotted names. +// It returns the new result of DetermineByKey after this change. +func (b *BitmapRowFilterBase) SetResult(key FilterKey, result FilterResult) FilterResult { + b.FilterResult = result + result, _ = b.DetermineByKey(key) + return result +} + +// Without a sub-filter, we always-succeed; if we get a key that isn't +// already answered by our YesKey/NoKey/lastRow, we will match this key, +// reject the rest of the row, and update our keys accordingly. We'll +// also hit the callback, and return an error from it if appropriate. +func (b *BitmapRowFilterBase) ConsiderKey(key FilterKey) FilterResult { + var done bool + b.FilterResult, done = b.DetermineByKey(key) + if done { + return b.FilterResult + } + b.FilterResult = key.MatchOneRejectRow() + row := key.Row() + b.lastRow = row + if b.callback != nil { + b.Err = b.callback(row) + } + return b.FilterResult +} + +// This should probably never be reached? +func (b *BitmapRowFilterBase) ConsiderData(key FilterKey, data *Container) FilterResult { + b.Err = errors.New("base iterator should never consider data") + return b.FilterResult +} + +func NewBitmapRowFilterBase(callback func(row uint64) error) *BitmapRowFilterBase { + return &BitmapRowFilterBase{lastRow: ^uint64(0), callback: callback} +} + +type BitmapRowLimitFilter struct { + BitmapRowFilterBase + limit uint64 +} + +var _ BitmapFilter = &BitmapRowLimitFilter{} + +// Without a sub-filter, we always-succeed; if we get a key that isn't +// already answered by our YesKey/NoKey/lastRow, we will match the whole +// row. +func (b *BitmapRowLimitFilter) ConsiderKey(key FilterKey) FilterResult { + var done bool + b.FilterResult, done = b.DetermineByKey(key) + if done { + return b.FilterResult + } + if b.limit > 0 { + b.FilterResult = key.MatchRow() + b.limit-- + } else { + b.FilterResult = key.Done() + } + return b.FilterResult +} + +// This should probably never be reached? +func (b *BitmapRowLimitFilter) ConsiderData(key FilterKey, data *Container) FilterResult { + b.Err = errors.New("limit iterator should never consider data") + return b.FilterResult +} + +func NewBitmapRowLimitFilter(limit uint64) *BitmapRowLimitFilter { + return &BitmapRowLimitFilter{BitmapRowFilterBase: *NewBitmapRowFilterBase(nil), limit: limit} +} + +// BitmapRowFilterSingleFilter is a row iterator with a single +// filter, which is simpler than one with multiple filters where +// it coincidentally turns out that N==1. +type BitmapRowFilterSingleFilter struct { + BitmapRowFilterBase + filter BitmapFilter +} + +func (b *BitmapRowFilterSingleFilter) ConsiderKey(key FilterKey) FilterResult { + res, done := b.DetermineByKey(key) + if done { + return res + } + return b.SetResult(key, b.filter.ConsiderKey(key)) +} + +func (b *BitmapRowFilterSingleFilter) ConsiderData(key FilterKey, data *Container) FilterResult { + // We already handled any consideration of the key above, in principle. + b.FilterResult = b.filter.ConsiderData(key, data) + if b.FilterResult.Err != nil { + return b.FilterResult + } + res, done := b.DetermineByKey(key) + if done { + return res + } + // We could just return the res, which would say nothing, but I + // think it should be a visible error if that happens. + b.FilterResult.Err = errors.New("inner filter didn't make a decision") + return b.FilterResult +} + +func NewBitmapRowFilterSingleFilter(callback func(row uint64) error, filter BitmapFilter) *BitmapRowFilterSingleFilter { + return &BitmapRowFilterSingleFilter{ + BitmapRowFilterBase: BitmapRowFilterBase{lastRow: ^uint64(0), callback: callback}, + filter: filter, + } +} + +// BitmapRowFilterMultiFilter is a BitmapFilter which wraps other bitmap filters, +// calling a callback function once per row whenever it finds a container +// for which all the filters returned true. +type BitmapRowFilterMultiFilter struct { + BitmapRowFilterBase + filters []BitmapFilter + yesKeys, noKeys []FilterKey + toDo []int +} + +func (b *BitmapRowFilterMultiFilter) ConsiderKey(key FilterKey) FilterResult { + res, done := b.DetermineByKey(key) + if done { + return res + } + // highestNo: The highest No value that we have that isn't preceeded + // by a relevant Yes. + highestNo := key + lowestYes := ^FilterKey(0) + // The length of the "no" run after the lowest "yes" + lowestYesNo := FilterKey(0) + b.toDo = b.toDo[:0] + // We scan for any no values that don't have an earlier yes that's + // still greater than this key. If there are any, we can skip to + // the highest such value immediately. We also build a todo list + // of items for which we have neither a yes nor a no answer greater + // than this key. + for i, yk := range b.yesKeys { + if yk > key { + if yk < lowestYes { + lowestYes = yk + lowestYesNo = b.noKeys[i] + } + continue + } + nk := b.noKeys[i] + if nk > highestNo { + highestNo = nk + continue + } + b.toDo = append(b.toDo, i) + } + // We have an unambiguous no, so we can set our internal state to + // be aware that we have a No until then. We can unconditionally + // return the result; it can't be not-done, because we just set + // it to a known done state. + if highestNo > key { + return b.SetResult(key, key.RejectUntil(highestNo)) + } + // Everything either has a yes value which is at least as high + // as lowestYes, or is in f.toDo now. Now we call ConsiderKey + // for everything in f.toDo, and accumulate a new list of the + // values still don't know, using the same backing store. + newToDo := b.toDo[:0] + for _, filter := range b.toDo { + result := b.filters[filter].ConsiderKey(key) + if result.Err != nil { + return key.Fail(result.Err) + } + yk, nk := result.YesKey, result.NoKey + b.yesKeys[filter], b.noKeys[filter] = yk, nk + if yk > key { + if lowestYes == 0 || yk < lowestYes { + lowestYes = yk + lowestYesNo = nk + } + continue + } + if nk > highestNo { + highestNo = nk + continue + } + newToDo = append(newToDo, filter) + } + // Same logic as before; if we have a highestNo, we don't need more + // information. + if highestNo > key { + return b.SetResult(key, key.RejectUntil(highestNo)) + } + b.toDo = newToDo + if len(b.toDo) > 0 { + return key.NeedData() + } + // this shouldn't be possible + if lowestYes <= key { + return key.Failf("got lowest yes %d for key %d, this shouldn't happen", lowestYes, key) + } + // Flag that we have a definite Yes as far as the lowest yes, and a + // definite No after that to the corresponding No. + return b.SetResult(key, key.MatchReject(lowestYes, lowestYesNo)) +} + +// ConsiderData only gets called in cases where f.toDo had a list of filters +// for which we needed to get data to make a decision. That means that +// everything but the indexes in f.toDo must be a "yes" right now. +func (b *BitmapRowFilterMultiFilter) ConsiderData(key FilterKey, data *Container) FilterResult { + res, done := b.DetermineByKey(key) + if done { + return res + } + highestNo := key + for _, filter := range b.toDo { + result := b.filters[filter].ConsiderData(key, data) + if result.Err != nil { + return key.Fail(result.Err) + } + yk, nk := result.YesKey, result.NoKey + b.yesKeys[filter], b.noKeys[filter] = yk, nk + if yk <= key && nk > highestNo { + highestNo = nk + } + } + if highestNo > key { + return b.SetResult(key, key.RejectUntil(highestNo)) + } + // if we got here, either something was buggy, or everything has a yes + // > key. + lowestYes := ^FilterKey(0) + lowestYesNo := key + for i, yk := range b.yesKeys { + if yk < lowestYes { + lowestYes = yk + lowestYesNo = b.noKeys[i] + } + } + // this shouldn't be possible + if lowestYes <= key { + return key.Failf("got lowest yes %d on data for key %d, this shouldn't happen", lowestYes, key) + } + return b.SetResult(key, key.MatchReject(lowestYes, lowestYesNo)) +} + +func NewBitmapColumnFilter(col uint64) BitmapFilter { + return &BitmapColumnFilter{key: uint16((col >> 16) & keyMask), offset: uint16(col & 0xFFFF)} +} + +// BitmapBitmap filter builds a list of positions in the bitmap which +// match those in a provided bitmap. It is shard-agnostic; no matter what +// offsets the input bitmap's containers have, it matches them against +// corresponding keys. +type BitmapBitmapFilter struct { + filter *Bitmap // We don't use this, but in ludicrous edge cases it might be holding a generation we need. + containers []*Container + nextOffsets []uint64 + callback func(uint64) error +} + +func (b *BitmapBitmapFilter) ConsiderKey(key FilterKey) FilterResult { + pos := key & keyMask + if b.containers[pos] == nil { + return key.RejectUntilOffset(b.nextOffsets[pos]) + } + return key.NeedData() +} + +func (b *BitmapBitmapFilter) ConsiderData(key FilterKey, data *Container) FilterResult { + pos := key & keyMask + base := uint64(key << 16) + filter := b.containers[pos] + if filter == nil || !IntersectionAny(data, filter) { + key.RejectUntilOffset(b.nextOffsets[pos]) + } + matching := intersect(data, filter) + offsets := matching.Slice() + for _, v := range offsets { + err := b.callback(base + uint64(v)) + if err != nil { + return key.Fail(err) + } + } + return key.MatchOneUntilOffset(b.nextOffsets[pos]) +} + +// NewBitmapBitmapFilter creates a filter which can report all the positions +// within a bitmap which are set, and which have positions corresponding to +// the specified columns. It calls the provided callback function on +// each value it finds, terminating early if that returns an error. +func NewBitmapBitmapFilter(filter *Bitmap, callback func(uint64) error) *BitmapBitmapFilter { + b := &BitmapBitmapFilter{ + filter: filter, + callback: callback, + containers: make([]*Container, rowWidth), + nextOffsets: make([]uint64, rowWidth), + } + iter, _ := filter.Containers.Iterator(0) + last := uint64(0) + count := 0 + for iter.Next() { + k, v := iter.Value() + k = k & keyMask + b.containers[k] = v + last = k + count++ + } + // if there's only one container, we need to populate everything with + // its position. + if count == 1 { + for i := range b.containers { + b.nextOffsets[i] = last + } + } else { + // Point each container at the offset of the next valid container. + // With sparse bitmaps this will potentially make skipping faster. + for i := range b.containers { + if b.containers[i] != nil { + for int(last) != i { + b.nextOffsets[last] = uint64(i) + last = (last + 1) % rowWidth + } + } + } + } + return b +} + +// BitmapRowFilterMultiFilter will call a +func NewBitmapRowFilterMultiFilter(callback func(row uint64) error, filters ...BitmapFilter) BitmapFilter { + return &BitmapRowFilterMultiFilter{ + filters: filters, + yesKeys: make([]FilterKey, len(filters)), + noKeys: make([]FilterKey, len(filters)), + BitmapRowFilterBase: BitmapRowFilterBase{ + callback: callback, + lastRow: ^uint64(0), + }, + } +} + +// BitmapRowLister returns a pointer to a slice which it will populate when invoked +// as a bitmap filter. +func NewBitmapRowFilter(callback func(uint64) error, filters ...BitmapFilter) BitmapFilter { + if len(filters) == 0 { + return NewBitmapRowFilterBase(callback) + } + if len(filters) == 1 { + return NewBitmapRowFilterSingleFilter(callback, filters[0]) + } + return NewBitmapRowFilterMultiFilter(callback, filters...) +} + +// ApplyFilterToIterator is a simplistic implementation that applies a bitmap +// filter to a ContainerIterator, returning an error if it encounters an error. +// +// This mostly exists for testing purposes; a Tx implementation where generating +// containers is expensive should almost certainly implement a better way to +// use filters which only generates data if it needs to. +func ApplyFilterToIterator(filter BitmapFilter, iter ContainerIterator) error { + defer iter.Close() + var until = uint64(0) + for (until < ^uint64(0)) && iter.Next() { + key, data := iter.Value() + if key < until { + continue + } + result := filter.ConsiderKey(FilterKey(key)) + if result.Err != nil { + return result.Err + } + until = uint64(result.NoKey) + if key < until { + continue + } + result = filter.ConsiderData(FilterKey(key), data) + if result.Err != nil { + return result.Err + } + until = uint64(result.NoKey) + } + return nil +} diff --git a/roaring/filter_internal_test.go b/roaring/filter_internal_test.go new file mode 100644 index 000000000..29f3ba7f7 --- /dev/null +++ b/roaring/filter_internal_test.go @@ -0,0 +1,226 @@ +// Copyright 2020 Pilosa Corp. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package roaring + +import ( + "fmt" + "sort" + "sync" + "testing" + + "github.com/pilosa/pilosa/v2/shardwidth" +) + +// For each container key i from 1 to (shard width in containers), we +// populate Bit i of that container in rows 0, K, etc. up to 100. +var filterSampleData *Bitmap +var filterSamplePositions [][]uint64 +var prepareSampleData sync.Once + +const sampleDataSize = 100 + +func requireSampleData(tb testing.TB) { + prepareSampleData.Do(func() { + b := NewBitmap(0) + filterSamplePositions = make([][]uint64, rowWidth) + for i := 1; i < rowWidth; i++ { + filterSamplePositions[i] = make([]uint64, 0, sampleDataSize/i) + base := (uint64(i) << 16) + uint64(i) + for row := 0; row < sampleDataSize; row += i { + pos := (uint64(row) << shardwidth.Exponent) + base + filterSamplePositions[i] = append(filterSamplePositions[i], pos) + _, _ = b.Add(pos) + } + } + filterSampleData = b + }) +} + +func applyFilter(tb testing.TB, data *Bitmap, filter *BitmapBitmapFilter) error { + iter, _ := data.Containers.Iterator(0) + return ApplyFilterToIterator(filter, iter) +} + +func getRows(tb testing.TB, data *Bitmap, filters ...BitmapFilter) []uint64 { + var rows []uint64 + callback := func(row uint64) error { + rows = append(rows, row) + return nil + } + iter, _ := data.Containers.Iterator(0) + filter := NewBitmapRowFilter(callback, filters...) + err := ApplyFilterToIterator(filter, iter) + if err != nil { + tb.Fatalf("unexpected error: %v", err) + } + return rows +} + +func compareSlices(tb testing.TB, name string, s1, s2 []uint64) { + if len(s1) != len(s2) { + tb.Fatalf("slice length mismatch %q: expected %d items %d, got %d items %d", + name, len(s1), s1, len(s2), s2) + } + for i, v := range s1 { + if s2[i] != v { + tb.Fatalf("row mismatch %q: expected item %d to be %d, got %d", + name, i, s1[i], s2[i]) + } + } +} + +func TestBaseFilter(t *testing.T) { + requireSampleData(t) + expected := make([]uint64, sampleDataSize) + // we expect every row, because stride1 hits them all + for i := range expected { + expected[i] = uint64(i) + } + rows := getRows(t, filterSampleData) + compareSlices(t, "base", expected, rows) +} + +func TestColumnFilter(t *testing.T) { + requireSampleData(t) + expected := make([]uint64, 0, sampleDataSize) + for i := 1; i < rowWidth; i++ { + t.Run(fmt.Sprintf("stride%d", i), func(t *testing.T) { + expected = expected[:0] + base := (uint64(i) << 16) + uint64(i) + for row := 0; row < sampleDataSize; row += i { + expected = append(expected, uint64(row)) + } + rows := getRows(t, filterSampleData, NewBitmapColumnFilter(base)) + compareSlices(t, "stride", expected, rows) + }) + } +} + +func TestRowsFilter(t *testing.T) { + requireSampleData(t) + rowSet := []uint64{0, 1, 2, 3} + expected := []uint64{0, 2} + // we expect every row to be present, because every row is present in stride1 + rows := getRows(t, filterSampleData, NewBitmapRowsFilter(rowSet)) + compareSlices(t, "initial", rowSet, rows) + // Check for stride-2 values that overlap with rowSet + rows = getRows(t, filterSampleData, NewBitmapRowsFilter(rowSet), NewBitmapColumnFilter((2<<16)+2)) + compareSlices(t, "rows-column", expected, rows) + rows = getRows(t, filterSampleData, NewBitmapColumnFilter((2<<16)+2), NewBitmapRowsFilter(rowSet)) + compareSlices(t, "column-rows", expected, rows) + rows = getRows(t, filterSampleData, NewBitmapColumnFilter((2<<16)+2), NewBitmapRowsFilter(rowSet), NewBitmapRowLimitFilter(1)) + compareSlices(t, "limit", expected[:1], rows) +} + +func TestBitmapFilter(t *testing.T) { + requireSampleData(t) + bm := NewBitmap() + expected := []uint64{} + positions := make([]uint64, 0, 80) + callback := func(pos uint64) error { + positions = append(positions, pos) + return nil + } + for i := 1; i < rowWidth; i++ { + positions = positions[:0] + expected = append(expected, filterSamplePositions[i]...) + sort.Slice(expected, func(i, j int) bool { return expected[i] < expected[j] }) + _, _ = bm.Add(uint64((i << 16) + i)) + filter := NewBitmapBitmapFilter(bm, callback) + err := applyFilter(t, filterSampleData, filter) + if err != nil { + t.Fatalf("unexpected error applying filter: %v", err) + } + compareSlices(t, fmt.Sprintf("stride-%d", i), expected, positions) + } +} + +func TestLimitFilter(t *testing.T) { + f := NewBitmapRowLimitFilter(5) + + for i := FilterKey(0); i < 5; i++ { + key := i << rowExponent + res := f.ConsiderKey(key) + if res.NoKey == ^FilterKey(0) { + t.Fatalf("limit filter ended early on iteration %d", i) + } + if res.YesKey <= key { + t.Fatalf("limit filter should always include until done") + } + } + res := f.ConsiderKey(FilterKey(5) << rowExponent) + if res.NoKey != ^FilterKey(0) { + t.Fatalf("limit filter should have thought it was done, reported last key %d", res.NoKey) + } +} + +func TestFilterWithRows(t *testing.T) { + tests := []struct { + rows []uint64 + callWith []uint64 + expect [][2]bool + }{ + { + rows: []uint64{}, + callWith: []uint64{0}, + expect: [][2]bool{{false, true}}, + }, + { + rows: []uint64{0}, + callWith: []uint64{0}, + expect: [][2]bool{{true, true}}, + }, + { + rows: []uint64{1}, + callWith: []uint64{0, 2}, + expect: [][2]bool{{false, false}, {false, true}}, + }, + { + rows: []uint64{0}, + callWith: []uint64{1, 2}, + expect: [][2]bool{{false, true}, {false, true}}, + }, + { + rows: []uint64{3, 9}, + callWith: []uint64{1, 2, 3, 10}, + expect: [][2]bool{{false, false}, {false, false}, {true, false}, {false, true}}, + }, + { + rows: []uint64{0, 1, 2}, + callWith: []uint64{0, 1, 2}, + expect: [][2]bool{{true, false}, {true, false}, {true, true}}, + }, + } + + for num, test := range tests { + t.Run(fmt.Sprintf("%d_%v_with_%v", num, test.rows, test.callWith), func(t *testing.T) { + if len(test.callWith) != len(test.expect) { + t.Fatalf("Badly specified test - must expect the same number of values as calls.") + } + f := NewBitmapRowsFilter(test.rows) + for i, id := range test.callWith { + key := FilterKey(id) << rowExponent + res := f.ConsiderKey(key) + inc := res.YesKey > key + done := res.NoKey == ^FilterKey(0) + if inc != test.expect[i][0] || done != test.expect[i][1] { + t.Logf("rows %d, calling row %d, result %#v\n", test.rows, id, res) + t.Fatalf("Calling with %d\nexp: %v,%v\ngot: %v,%v", id, test.expect[i][0], test.expect[i][1], inc, done) + } + } + }) + } + +} diff --git a/roaring/printutil.go b/roaring/printutil.go index 1f4e8feaa..01fb4e6e2 100644 --- a/roaring/printutil.go +++ b/roaring/printutil.go @@ -21,13 +21,6 @@ import ( "github.com/pilosa/pilosa/v2/shardwidth" ) -// TODO: once fastRows2 merges, take out these redundant const definitions -const ( - rowExponent = (shardwidth.Exponent - 16) // e.g. 20 - 16 == 4 - rowWidth = 1 << rowExponent // containers per row // e.g. 1 << 4 == 16 - keyMask = (rowWidth - 1) // a mask for offset within the row // e.g. 0x0000000f or 15 -) - func (b *Bitmap) String() (r string) { r = "c(" slc := b.Slice() diff --git a/roaring/roaring.go b/roaring/roaring.go index c0926caaa..f0dffd560 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -7164,6 +7164,10 @@ func Difference(a, b *Container) *Container { return difference(a, b) } +func Intersect(x, y *Container) *Container { + return intersect(x, y) +} + func IntersectionCount(x, y *Container) int32 { return intersectionCount(x, y) } @@ -7205,3 +7209,41 @@ func (c *Container) Difference(other *Container) *Container { func NewSliceContainers() *sliceContainers { return newSliceContainers() } + +// Slice returns an array of the values in the container as uint16. +// Do NOT modify the result; it could be the container's actual storage. +func (c *Container) Slice() (r []uint16) { + if c == nil { + return r + } + switch c.typ() { + case ContainerArray: + r = c.array() + case ContainerBitmap: + r = make([]uint16, c.N()) + n := int32(0) + for i, word := range c.bitmap() { + for word != 0 { + t := word & -word + if roaringParanoia { + if n >= c.N() { + panic("bitmap has more bits set than container.n") + } + } + r[n] = uint16((i*64 + int(popcount(t-1)))) + n++ + word ^= t + } + } + case ContainerRun: + r = make([]uint16, c.N()) + n := 0 + for _, run := range c.runs() { + for v := int(run.Start); v <= int(run.Last); v++ { + r[n] = uint16(v) + n++ + } + } + } + return r +}