From 9237e7b6995618af7f6aa818707e88ab4208038b Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 27 Feb 2017 12:42:36 -0600 Subject: [PATCH] rank cache update after count heap sort order backwards WIP TopN accuracy adjusted first phase topn to collect all slices id's incorrect handling of large topns cache performance enhancement fix for failed test TestMain_FrameRestore remove unused code and fix some variable names --- bitmap.go | 23 ++++++++++++++++++++++- cache.go | 49 ++++++++++++++++++++++++++++++++++++++----------- executor.go | 24 ++++++++++++++---------- fragment.go | 47 ++++++++++++++++++++++++++++++----------------- 4 files changed, 104 insertions(+), 39 deletions(-) diff --git a/bitmap.go b/bitmap.go index 55664b24b..d30a9192a 100644 --- a/bitmap.go +++ b/bitmap.go @@ -16,6 +16,8 @@ type Bitmap struct { // Attributes associated with the bitmap. Attrs map[string]interface{} + + cacheoveride uint64 } // NewBitmap returns a new instance of Bitmap. @@ -174,6 +176,26 @@ func (b *Bitmap) InvalidateCount() { } } +//increment the bitmap cached counter, note this is an optimization that assumes that the caller is aware the size increased +func (b *Bitmap) IncrementCount(i uint64) { + seg := b.segment(i / SliceWidth) + if seg != nil { + seg.n++ + } +} +func (b *Bitmap) DecrementCount(i uint64) { + seg := b.segment(i / SliceWidth) + if seg != nil { + if seg.n > 0 { + seg.n-- + } + } +} +func (b *Bitmap) SetCount(i uint64, count uint64) { + seg := b.segment(i / SliceWidth) + seg.n = count +} + // Count returns the number of set bits in the bitmap. func (b *Bitmap) Count() uint64 { var n uint64 @@ -312,7 +334,6 @@ func (s *BitmapSegment) Difference(other *BitmapSegment) *BitmapSegment { // SetBit sets the i-th bit of the bitmap. func (s *BitmapSegment) SetBit(i uint64) (changed bool) { s.ensureWritable() - changed, _ = s.data.Add(i) if changed { s.n++ diff --git a/cache.go b/cache.go index 86cab6b5a..4345c93fd 100644 --- a/cache.go +++ b/cache.go @@ -5,6 +5,7 @@ import ( "fmt" "io" "sort" + "sync" "time" "github.com/golang/groupcache/lru" @@ -97,6 +98,7 @@ var _ Cache = &LRUCache{} // RankCache represents a cache with sorted entries. type RankCache struct { + mu sync.Mutex entries map[uint64]uint64 rankings []BitmapPair // cached, ordered list @@ -117,6 +119,8 @@ func NewRankCache() *RankCache { // Add adds a bitmap to the cache. func (c *RankCache) Add(bitmapID uint64, n uint64) { + c.mu.Lock() + defer c.mu.Unlock() // Ignore if the bit count on the bitmap is below the threshold. if n < c.ThresholdValue { return @@ -124,19 +128,13 @@ func (c *RankCache) Add(bitmapID uint64, n uint64) { c.entries[bitmapID] = n - c.Invalidate() - // 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 { - delete(c.entries, id) - } - } - } + c.invalidate() } // BulkAdd adds a bitmap to the cache unsorted. You should Invalidate after completion. func (c *RankCache) BulkAdd(bitmapID uint64, n uint64) { + c.mu.Lock() + defer c.mu.Unlock() if n < c.ThresholdValue { return } @@ -145,13 +143,23 @@ func (c *RankCache) BulkAdd(bitmapID uint64, n uint64) { } // Get returns a bitmap with a given id. -func (c *RankCache) Get(bitmapID uint64) uint64 { return c.entries[bitmapID] } +func (c *RankCache) Get(bitmapID uint64) uint64 { + c.mu.Lock() + defer c.mu.Unlock() + return c.entries[bitmapID] +} // Len returns the number of items in the cache. -func (c *RankCache) Len() int { return len(c.entries) } +func (c *RankCache) Len() int { + c.mu.Lock() + defer c.mu.Unlock() + return len(c.entries) +} // BitmapIDs returns a list of all bitmap IDs in the cache. func (c *RankCache) BitmapIDs() []uint64 { + c.mu.Lock() + defer c.mu.Unlock() a := make([]uint64, 0, len(c.entries)) for id := range c.entries { a = append(a, id) @@ -162,6 +170,15 @@ func (c *RankCache) BitmapIDs() []uint64 { // update reorders the entries by rank. func (c *RankCache) Invalidate() { + c.mu.Lock() + defer c.mu.Unlock() + c.invalidate() + +} +func (c *RankCache) invalidate() { + 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)) @@ -183,6 +200,14 @@ 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 { + delete(c.entries, id) + } + } + } } // Top returns an ordered list of bitmaps. @@ -243,6 +268,8 @@ type PairHeap struct { Pairs } +func (p PairHeap) Less(i, j int) bool { return p.Pairs[i].Count < p.Pairs[j].Count } + func (h *Pairs) Push(x interface{}) { // Push and Pop use pointer receivers because they modify the slice's length, // not just its contents. diff --git a/executor.go b/executor.go index 5f6b766c9..4a2d954c9 100644 --- a/executor.go +++ b/executor.go @@ -168,6 +168,7 @@ func (e *Executor) executeBitmapCallSlice(ctx context.Context, db string, c *pql // requeries to retrieve the full counts for each of the top results. func (e *Executor) executeTopN(ctx context.Context, db string, c *pql.Call, slices []uint64, opt *ExecOptions) ([]Pair, error) { bitmapIDs, _ := c.Args["ids"].([]uint64) + n := c.Args["n"].(uint64) // Execute original query. pairs, err := e.executeTopNSlices(ctx, db, c, slices, opt) @@ -180,21 +181,28 @@ func (e *Executor) executeTopN(ctx context.Context, db string, c *pql.Call, slic if len(pairs) == 0 || len(bitmapIDs) > 0 || opt.Remote { return pairs, nil } - // 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... + other.Args["n"] = len(bitmapIDs) * 2 ids := Pairs(pairs).Keys() sort.Sort(uint64Slice(ids)) other.Args["ids"] = ids - return e.executeTopNSlices(ctx, db, other, slices, opt) + trimmedList, err := e.executeTopNSlices(ctx, db, other, slices, opt) + if err != nil { + return nil, err + } + + if int(n) < len(trimmedList) { + trimmedList = trimmedList[0:n] + } + 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) @@ -215,11 +223,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 } @@ -896,6 +899,7 @@ func (e *Executor) mapper(ctx context.Context, ch chan mapResponse, nodes []*Nod if n.Host == e.Host { resp.result, resp.err = e.mapperLocal(ctx, nodeSlices, mapFn, reduceFn) } else if !opt.Remote { + results, err := e.exec(ctx, n, db, &pql.Query{Calls: []*pql.Call{c}}, nodeSlices, opt) if len(results) > 0 { resp.result = results[0] diff --git a/fragment.go b/fragment.go index d6eead246..e203dd711 100644 --- a/fragment.go +++ b/fragment.go @@ -351,8 +351,10 @@ func (f *Fragment) bitmap(bitmapID uint64, updateCache bool) *Bitmap { slice: f.slice, writable: false, }}, + cacheoveride: 0, } bm.InvalidateCount() + bm.cacheoveride = bm.Count() if updateCache { // Update cache. @@ -380,6 +382,7 @@ func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) } // Write to storage. + if changed, err = f.storage.Add(pos); err != nil { return false, err } @@ -398,9 +401,16 @@ func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) } // Update the cache. - if f.bitmap(bitmapID, true).SetBit(profileID) { - changed = true + bm := f.bitmap(bitmapID, true) + if bm.cacheoveride > 0 { + bm.SetCount(profileID, bm.cacheoveride) + bm.cacheoveride = 0 + } else { + bm.IncrementCount(profileID) } + bm.SetBit(profileID) + + f.cache.Add(bitmapID, bm.Count()) f.stats.Count("setN", 1) @@ -442,9 +452,14 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { } // Update the cache. - if f.bitmap(bitmapID, true).ClearBit(profileID) { - return true, nil + 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) @@ -548,12 +563,7 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { // 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. @@ -591,21 +601,24 @@ func (f *Fragment) topBitmapPairs(bitmapIDs []uint64) []BitmapPair { } // Otherwise retrieve specific bitmaps. - pairs := make([]BitmapPair, len(bitmapIDs)) - for i, bitmapID := range bitmapIDs { + pairs := make([]BitmapPair, 0, len(bitmapIDs)) + for _, bitmapID := range bitmapIDs { // Look up cache first, if available. if n := f.cache.Get(bitmapID); n > 0 { - pairs[i] = BitmapPair{ + pairs = append(pairs, BitmapPair{ ID: bitmapID, Count: n, - } + }) continue } - // Otherwise load from storage. - pairs[i] = BitmapPair{ - ID: bitmapID, - Count: f.Bitmap(bitmapID).Count(), + bm := f.Bitmap(bitmapID) + if bm.Count() > 0 { + // Otherwise load from storage. + pairs = append(pairs, BitmapPair{ + ID: bitmapID, + Count: bm.Count(), + }) } } sort.Sort(BitmapPairs(pairs))