From 7c415b42178bba9b421e791b26e34827ae4b4306 Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 9 Dec 2020 13:54:14 -0600 Subject: [PATCH] Implement rbf-specific ApplyFilter This gives RBF an ApplyFilter that can run without instantiating containers when the filter it's using doesn't need them instantiated. We can also seek ahead in cases where we know the next key we care about is not just the next key numerically. --- rbf.go | 4 +-- rbf/cursorx.go | 43 ++++++++++++++++++++++++ rbf/tx.go | 69 ++++++++++++++++++++++++++++++++++++++ roaring/container_stash.go | 19 +++++++++++ 4 files changed, 133 insertions(+), 2 deletions(-) diff --git a/rbf.go b/rbf.go index c0b50d838..c0d95502f 100644 --- a/rbf.go +++ b/rbf.go @@ -438,8 +438,8 @@ func (tx *RBFTx) UseRowCache() bool { return rbf.EnableRowCache() } -func (c *RBFTx) ApplyFilter(index, field, view string, shard uint64, ckey uint64, filter roaring.BitmapFilter) (err error) { - return GenericApplyFilter(c, index, field, view, shard, ckey, filter) +func (tx *RBFTx) ApplyFilter(index, field, view string, shard uint64, ckey uint64, filter roaring.BitmapFilter) (err error) { + return tx.tx.ApplyFilter(rbfName(index, field, view, shard), ckey, filter) } // rbfName returns a NULL-separated key used for identifying bitmap maps in RBF. diff --git a/rbf/cursorx.go b/rbf/cursorx.go index e487029c3..171bbd376 100644 --- a/rbf/cursorx.go +++ b/rbf/cursorx.go @@ -20,6 +20,7 @@ import ( "math" "os" "sync/atomic" + "unsafe" "github.com/pilosa/pilosa/v2/roaring" "github.com/pkg/errors" @@ -184,6 +185,48 @@ func (c *Cursor) CurrentPageType() ContainerType { return cell.Type } +func intoContainer(l leafCell, tx *Tx, replacing *roaring.Container, target []byte) (c *roaring.Container) { + if len(l.Data) == 0 { + return nil + } + orig := l.Data + var cpMaybe []byte + var mapped bool + if EnableRowCache() || tx.db.cfg.DoAllocZero { + // make a copy, otherwise the rowCache will see corrupted data + // or mmapped data that may disappear. + cpMaybe = target[:len(orig)] + copy(cpMaybe, orig) + mapped = false + } else { + // not a copy + cpMaybe = orig + mapped = true + } + switch l.Type { + case ContainerTypeArray: + c = roaring.RemakeContainerArray(replacing, toArray16(cpMaybe)) + case ContainerTypeBitmapPtr: + _, bm, _ := tx.leafCellBitmap(toPgno(cpMaybe)) + cloneMaybe := bm + if EnableRowCache() { + cloneMaybe = (*[1024]uint64)(unsafe.Pointer(&target[0]))[:1024] + copy(cloneMaybe, bm) + } + c = roaring.RemakeContainerBitmap(replacing, cloneMaybe) + case ContainerTypeBitmap: + c = roaring.RemakeContainerBitmap(replacing, toArray64(cpMaybe)) + case ContainerTypeRLE: + c = roaring.RemakeContainerRun(replacing, toInterval16(cpMaybe)) + } + // Note: If the "roaringparanoia" build tag isn't set, this + // should be optimized away entirely. Otherwise it's moderately + // expensive. + c.CheckN() + c.SetMapped(mapped) + return c +} + func toContainer(l leafCell, tx *Tx) (c *roaring.Container) { if len(l.Data) == 0 { diff --git a/rbf/tx.go b/rbf/tx.go index b4f1c7306..7c2ebbd49 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -1053,6 +1053,26 @@ func (tx *Tx) ContainerIterator(name string, key uint64) (citer roaring.Containe return &containerIterator{cursor: c}, exact, nil } +func (tx *Tx) ApplyFilter(name string, key uint64, filter roaring.BitmapFilter) (err error) { + tx.mu.RLock() + defer tx.mu.RUnlock() + + c, err := tx.cursor(name) + if err == ErrBitmapNotFound { + return nil // nothing available. + } else if err != nil { + return err + } + + _, err = c.Seek(key) + if err != nil { + return err + } + f := containerFilter{cursor: c, filter: filter, tx: tx} + defer f.Close() + return f.Apply() +} + func (tx *Tx) ForEach(name string, fn func(i uint64) error) error { return tx.ForEachRange(name, 0, math.MaxUint64, fn) } @@ -1398,6 +1418,55 @@ func (tx *Tx) OffsetRange(name string, offset, start, endx uint64) (*roaring.Bit return other, nil } +// containerFilter is like ContainerIterator, but implements ApplyFilter +type containerFilter struct { + cursor *Cursor + filter roaring.BitmapFilter + tx *Tx + header roaring.Container + body [8192]byte +} + +func (s *containerFilter) Close() { + s.cursor.Close() +} + +func (s *containerFilter) Apply() (err error) { + var minKey roaring.FilterKey + for err := s.cursor.Next(); err == nil; err = s.cursor.Next() { + elem := &s.cursor.stack.elems[s.cursor.stack.top] + leafPage, _, _ := s.cursor.tx.readPage(elem.pgno) + cell := readLeafCell(leafPage, elem.index) + key := roaring.FilterKey(cell.Key) + if key < minKey { + continue + } + s.tx.mu.RUnlock() + res := s.filter.ConsiderKey(key) + s.tx.mu.RLock() + if res.Err != nil { + return res.Err + } + if res.YesKey <= key && res.NoKey <= key { + data := intoContainer(cell, s.cursor.tx, &s.header, s.body[:]) + s.tx.mu.RUnlock() + res = s.filter.ConsiderData(key, data) + s.tx.mu.RLock() + if res.Err != nil { + return res.Err + } + } + minKey = res.NoKey + if minKey > key+1 { + _, err := s.cursor.Seek(uint64(minKey)) + if err != nil { + return err + } + } + } + return nil +} + // containerIterator wraps Cursor to implement roaring.ContainerIterator. type containerIterator struct { cursor *Cursor diff --git a/roaring/container_stash.go b/roaring/container_stash.go index 8a959fdde..2cde05743 100644 --- a/roaring/container_stash.go +++ b/roaring/container_stash.go @@ -125,6 +125,25 @@ func NewContainer() *Container { return NewContainerArray(nil) } +func RemakeContainerBitmap(c *Container, bitmap []uint64) *Container { + *c = Container{typeID: ContainerBitmap} + c.setBitmap(bitmap) + c.bitmapRepair() + return c +} + +func RemakeContainerArray(c *Container, array []uint16) *Container { + *c = Container{typeID: ContainerArray} + c.setArray(array) + return c +} + +func RemakeContainerRun(c *Container, intervals []Interval16) *Container { + *c = Container{typeID: ContainerRun} + c.setRuns(intervals) + return c +} + // NewContainerBitmap makes a bitmap container using the provided bitmap, or // an empty one if provided bitmap is nil. If the provided bitmap is too short, // it will be padded. This function's API is wrong; it should have been