diff --git a/cache.go b/cache.go index 4345c93fd..7c2e96f2d 100644 --- a/cache.go +++ b/cache.go @@ -176,16 +176,17 @@ func (c *RankCache) Invalidate() { } func (c *RankCache) invalidate() { + // Don't invalidate more than once every X seconds. + // TODO: consider making this configurable. if time.Now().Sub(c.updateTime).Seconds() < 10 { return } - //fmt.Println("RankCache Update") // Convert cache to a sorted list. rankings := make([]BitmapPair, 0, len(c.entries)) - for id, n := range c.entries { + for id, cnt := range c.entries { rankings = append(rankings, BitmapPair{ ID: id, - Count: n, + Count: cnt, }) } sort.Sort(BitmapPairs(rankings)) @@ -200,10 +201,11 @@ func (c *RankCache) invalidate() { // Reset counters. c.updateTime, c.updateN = time.Now(), 0 + // If size is larger than the threshold then trim it. if len(c.entries) > c.ThresholdLength { - for id, n := range c.entries { - if n <= c.ThresholdValue { + for id, cnt := range c.entries { + if cnt <= c.ThresholdValue { delete(c.entries, id) } } @@ -378,3 +380,27 @@ func (p uint64Slice) merge(other []uint64) []uint64 { return ret } + +// BitmapCache provides an interface for caching full bitmaps. +type BitmapCache interface { + Fetch(id uint64) (*Bitmap, bool) + Add(id uint64, b *Bitmap) +} + +// SimpleCache implements BitmapCache +// it is meant to be a short-lived cache for cases where writes are continuing to access +// the same bit within a short time frame (i.e. good for write-heavy loads) +// A read-heavy use case would cause the cache to get bigger, potentially causing the +// node to run out of memory. +type SimpleCache struct { + cache map[uint64]*Bitmap +} + +func (s *SimpleCache) Fetch(id uint64) (*Bitmap, bool) { + m, ok := s.cache[id] + return m, ok +} + +func (s *SimpleCache) Add(id uint64, b *Bitmap) { + s.cache[id] = b +} diff --git a/executor.go b/executor.go index f1f1adb3c..26b2f4747 100644 --- a/executor.go +++ b/executor.go @@ -183,27 +183,27 @@ func (e *Executor) executeTopN(ctx context.Context, db string, c *pql.Call, slic } // Only the original caller should refetch the full counts. other := c.Clone() - //other.Args["n"] = 0 + + // Double the size of n for other calls in order to... + // TODO: travis review other.Args["n"] = len(bitmapIDs) * 2 ids := Pairs(pairs).Keys() sort.Sort(uint64Slice(ids)) other.Args["ids"] = ids - trimedlist, x := e.executeTopNSlices(ctx, db, other, slices, opt) - if x != nil { - return nil, x + trimmedList, err := e.executeTopNSlices(ctx, db, other, slices, opt) + if err != nil { + return nil, err } - if int(n) < len(trimedlist) { - trimedlist = trimedlist[0:n] + if int(n) < len(trimmedList) { + trimmedList = trimmedList[0:n] } - return trimedlist, nil + return trimmedList, nil } func (e *Executor) executeTopNSlices(ctx context.Context, db string, c *pql.Call, slices []uint64, opt *ExecOptions) ([]Pair, error) { - //n, _ := c.Args["n"].(uint64) - // Execute calls in bulk on each remote node and merge. mapFn := func(slice uint64) (interface{}, error) { return e.executeTopNSlice(ctx, db, c, slice) @@ -224,11 +224,6 @@ func (e *Executor) executeTopNSlices(ctx context.Context, db string, c *pql.Call // Sort final merged results. sort.Sort(Pairs(results)) - // Only keep the top n after sorting. - // if n > 0 && len(results) > int(n) { - // results = results[0:n] - //} - return results, nil } diff --git a/fragment.go b/fragment.go index ac4d63926..fd73ff343 100644 --- a/fragment.go +++ b/fragment.go @@ -53,28 +53,6 @@ const ( DefaultFragmentMaxOpN = 2000 ) -// BitmapCacher implements SimpleCache -// it is meant to be a short-lived cache for cases where writes are continuing to access -// the same bit withing a short time frame (i.e. good for write-heavy loads) -// A read-heavy use case would cause the cache to get bigger, potentially causing the -// node to run out of memory. -type BitmapCacher interface { - Fetch(id uint64) (*Bitmap, bool) - Add(id uint64, b *Bitmap) -} - -type SimpleCache struct { - cache map[uint64]*Bitmap -} - -func (s *SimpleCache) Fetch(id uint64) (*Bitmap, bool) { - m, ok := s.cache[id] - return m, ok -} -func (s *SimpleCache) Add(id uint64, p *Bitmap) { - s.cache[id] = p -} - // Fragment represents the intersection of a frame and slice in a database. type Fragment struct { mu sync.Mutex @@ -91,9 +69,12 @@ type Fragment struct { storageData []byte opN int // number of ops since snapshot - // Bitmap cache. + // Cache for bitmap counts. cache Cache + // Cache containing full bitmaps (not just counts). + bitmapCache BitmapCache + // Cached checksums for each block. checksums map[int][]byte @@ -109,8 +90,7 @@ type Fragment struct { // This is set by the parent frame unless overridden for testing. BitmapAttrStore *AttrStore - stats StatsClient - bitmapCache BitmapCacher + stats StatsClient } // NewFragment returns a new instance of Fragment. @@ -265,7 +245,7 @@ func (f *Fragment) openCache() error { // This will cause them to be added to the cache. for _, bitmapID := range pb.BitmapIDs { //n := f.storage.CountRange(bitmapID*SliceWidth, (bitmapID+1)*SliceWidth) - n := f.bitmap(bitmapID, false).Count() + n := f.bitmap(bitmapID, true, false).Count() f.cache.BulkAdd(bitmapID, n) } f.cache.Invalidate() @@ -332,22 +312,28 @@ func (f *Fragment) logger() *log.Logger { return log.New(f.LogOutput, "", log.Ls func (f *Fragment) Bitmap(bitmapID uint64) *Bitmap { f.mu.Lock() defer f.mu.Unlock() - return f.bitmap(bitmapID, false) + return f.bitmap(bitmapID, true, true) } -func (f *Fragment) bitmap(bitmapID uint64, updateCache bool) *Bitmap { - r, ok := f.bitmapCache.Fetch(bitmapID) - if ok && r != nil { - return r +func (f *Fragment) bitmap(bitmapID uint64, checkBitmapCache bool, updateBitmapCache bool) *Bitmap { + + if checkBitmapCache { + r, ok := f.bitmapCache.Fetch(bitmapID) + if ok && r != nil { + return r + } } + // Only use a subset of the containers. // NOTE: The start & end ranges must be divisible by data := f.storage.OffsetRange(f.slice*SliceWidth, bitmapID*SliceWidth, (bitmapID+1)*SliceWidth) // Reference bitmap subrange in storage. + // We Clone() data because otherwise bm will contains pointers to containers in storage. + // This causes unexpected results when we cache the bitmap and try to use it later. bm := &Bitmap{ segments: []BitmapSegment{{ - data: *data, + data: *data.Clone(), slice: f.slice, writable: false, }}, @@ -356,9 +342,7 @@ func (f *Fragment) bitmap(bitmapID uint64, updateCache bool) *Bitmap { bm.InvalidateCount() bm.cacheoveride = bm.Count() - if updateCache { - // Update cache. - f.cache.Add(bitmapID, bm.Count()) + if updateBitmapCache { f.bitmapCache.Add(bitmapID, bm) } @@ -373,9 +357,9 @@ func (f *Fragment) SetBit(bitmapID, profileID uint64) (changed bool, err error) return f.setBit(bitmapID, profileID) } -func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) { - // Determine the position of the bit in the storage. +func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, err error) { changed = false + // Determine the position of the bit in the storage. pos, err := f.pos(bitmapID, profileID) if err != nil { return false, err @@ -395,21 +379,16 @@ func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) // Invalidate block checksum. delete(f.checksums, int(bitmapID/HashBlockSize)) - // If the number of operations exceeds the limit then snapshot. + // Increment number of operations until snapshot is required. if err := f.incrementOpN(); err != nil { return false, err } - // Update the cache. - bm := f.bitmap(bitmapID, true) - if bm.cacheoveride > 0 { - bm.SetCount(profileID, bm.cacheoveride) - bm.cacheoveride = 0 - } else { - bm.IncrementCount(profileID) - } + // Get the bitmap from bitmapCache or fragment.storage. + bm := f.bitmap(bitmapID, true, true) bm.SetBit(profileID) + // Update the cache. f.cache.Add(bitmapID, bm.Count()) f.stats.Count("setN", 1) @@ -425,7 +404,8 @@ func (f *Fragment) ClearBit(bitmapID, profileID uint64) (bool, error) { return f.clearBit(bitmapID, profileID) } -func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { +func (f *Fragment) clearBit(bitmapID, profileID uint64) (changed bool, err error) { + changed = false // Determine the position of the bit in the storage. pos, err := f.pos(bitmapID, profileID) if err != nil { @@ -433,8 +413,7 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { } // Write to storage. - changed, err := f.storage.Remove(pos) - if err != nil { + if changed, err = f.storage.Remove(pos); err != nil { return false, err } @@ -451,14 +430,11 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { return false, err } + // Get the bitmap from bitmapCache or fragment.storage. + bm := f.bitmap(bitmapID, true, true) + bm.ClearBit(profileID) + // Update the cache. - bm := f.bitmap(bitmapID, false) - if bm.cacheoveride > 0 { - bm.SetCount(profileID, bm.cacheoveride) - bm.cacheoveride = 0 - } else { - bm.DecrementCount(profileID) - } f.cache.Add(bitmapID, bm.Count()) f.stats.Count("clearN", 1) @@ -546,7 +522,6 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { if count == 0 { continue } - //results = append(results, Pair{Key: bitmapID, Count: count}) heap.Push(results, Pair{Key: bitmapID, Count: count}) // If we reach the requested number of pairs and we are not computing @@ -556,19 +531,13 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { if opt.Src == nil { break } - // sort.Sort(Pairs(results)) } continue } // Retrieve the lowest count we have. // If it's too low then don't try finding anymore pairs. - //threshold := results[len(results)-1].Count - threshold := results.Pairs[0].Count - //if threshold < MinThreshold { - //break - //} // If the bitmap doesn't have enough bits set before the intersection // then we can assume that any remaing bitmaps also have a count too low. @@ -627,7 +596,6 @@ func (f *Fragment) topBitmapPairs(bitmapIDs []uint64) []BitmapPair { } } sort.Sort(BitmapPairs(pairs)) - //debugDumpPairs(pairs) return pairs } @@ -937,7 +905,10 @@ func (f *Fragment) Import(bitmapIDs, profileIDs []uint64) error { // Update cache counts for all bitmaps. for bitmapID := range set { - f.cache.BulkAdd(bitmapID, f.bitmap(bitmapID, false).Count()) + // Import should ALWAYS have bitmap() load a new bm from fragment.storage + // because the bitmap that's in bitmapCache hasn't been updated with + // this import's data. + f.cache.BulkAdd(bitmapID, f.bitmap(bitmapID, false, false).Count()) } f.cache.Invalidate() @@ -1285,7 +1256,6 @@ func (s *FragmentSyncer) SyncFragment() error { // Determine replica set. nodes := s.Cluster.FragmentNodes(s.Fragment.DB(), s.Fragment.Slice()) if len(nodes) == 1 { - //fmt.Println("no place to replicate", s.Fragment.DB(), s.Fragment.Frame(), s.Fragment.Slice()) return nil } diff --git a/roaring/roaring.go b/roaring/roaring.go index 58b4d1e19..e74739311 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -1084,7 +1084,7 @@ func (c *container) clone() *container { copy(other.bitmap, c.bitmap) } - return c + return other } // WriteTo writes c to w.