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 +}