From 9237e7b6995618af7f6aa818707e88ab4208038b Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 27 Feb 2017 12:42:36 -0600 Subject: [PATCH 1/4] 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)) From cbac9ec88bf80a864f95bb19115329b7ca469c39 Mon Sep 17 00:00:00 2001 From: Travis Date: Thu, 2 Mar 2017 15:17:21 -0600 Subject: [PATCH 2/4] move BitmapCache into cache.go. adjust some variable names for clarity --- bitmap.go | 2 +- cache.go | 35 +++++++++++++++++++++++++----- executor.go | 1 + fragment.go | 62 ++++++++++++++++------------------------------------- 4 files changed, 51 insertions(+), 49 deletions(-) diff --git a/bitmap.go b/bitmap.go index d30a9192a..651950e45 100644 --- a/bitmap.go +++ b/bitmap.go @@ -17,7 +17,7 @@ type Bitmap struct { // Attributes associated with the bitmap. Attrs map[string]interface{} - cacheoveride uint64 + adjustedCount uint64 } // NewBitmap returns a new instance of Bitmap. diff --git a/cache.go b/cache.go index 4345c93fd..184310361 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,26 @@ 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) { + return s.cache[id] +} + +func (s *SimpleCache) Add(id uint64, b *Bitmap) { + s.cache[id] = b +} diff --git a/executor.go b/executor.go index 4a2d954c9..26b2f4747 100644 --- a/executor.go +++ b/executor.go @@ -185,6 +185,7 @@ func (e *Executor) executeTopN(ctx context.Context, db string, c *pql.Call, slic other := c.Clone() // Double the size of n for other calls in order to... + // TODO: travis review other.Args["n"] = len(bitmapIDs) * 2 ids := Pairs(pairs).Keys() diff --git a/fragment.go b/fragment.go index e203dd711..8f8698192 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. @@ -351,10 +331,10 @@ func (f *Fragment) bitmap(bitmapID uint64, updateCache bool) *Bitmap { slice: f.slice, writable: false, }}, - cacheoveride: 0, + adjustedCount: 0, } bm.InvalidateCount() - bm.cacheoveride = bm.Count() + bm.adjustedCount = bm.Count() if updateCache { // Update cache. @@ -382,7 +362,6 @@ 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 } @@ -400,18 +379,17 @@ func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) return false, err } - // Update the cache. + // If adjustedCount is set, then apply that value to bitmap.n instead. bm := f.bitmap(bitmapID, true) - if bm.cacheoveride > 0 { - bm.SetCount(profileID, bm.cacheoveride) - bm.cacheoveride = 0 + if bm.adjustedCount > 0 { + bm.SetCount(profileID, bm.adjustedCount) + bm.adjustedCount = 0 } else { bm.IncrementCount(profileID) + f.cache.Add(bitmapID, bm.Count()) } bm.SetBit(profileID) - f.cache.Add(bitmapID, bm.Count()) - f.stats.Count("setN", 1) return changed, nil @@ -451,15 +429,16 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { return false, err } - // Update the cache. - bm := f.bitmap(bitmapID, false) - if bm.cacheoveride > 0 { - bm.SetCount(profileID, bm.cacheoveride) - bm.cacheoveride = 0 + // If adjustedCount is set, then apply that value to bitmap.n instead. + bm := f.bitmap(bitmapID, true) + if bm.adjustedCount > 0 { + bm.SetCount(profileID, bm.adjustedCount) + bm.adjustedCount = 0 } else { bm.DecrementCount(profileID) + f.cache.Add(bitmapID, bm.Count()) } - f.cache.Add(bitmapID, bm.Count()) + bm.ClearBit(profileID) f.stats.Count("clearN", 1) @@ -546,7 +525,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,7 +534,6 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { if opt.Src == nil { break } - // sort.Sort(Pairs(results)) } continue } @@ -622,7 +599,6 @@ func (f *Fragment) topBitmapPairs(bitmapIDs []uint64) []BitmapPair { } } sort.Sort(BitmapPairs(pairs)) - //debugDumpPairs(pairs) return pairs } From c013661880d0682521813f9374e16e8704636e7e Mon Sep 17 00:00:00 2001 From: Travis Date: Thu, 2 Mar 2017 16:00:34 -0600 Subject: [PATCH 3/4] bug fix in SimpleCache.Fetch (really just reverting my change) --- cache.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/cache.go b/cache.go index 184310361..7c2e96f2d 100644 --- a/cache.go +++ b/cache.go @@ -397,7 +397,8 @@ type SimpleCache struct { } func (s *SimpleCache) Fetch(id uint64) (*Bitmap, bool) { - return s.cache[id] + m, ok := s.cache[id] + return m, ok } func (s *SimpleCache) Add(id uint64, b *Bitmap) { From c3ed12c20ebe09c12efc431bfcd9666f1c588305 Mon Sep 17 00:00:00 2001 From: Travis Date: Mon, 6 Mar 2017 11:06:22 -0600 Subject: [PATCH 4/4] Refactor fragment.bitmap() so that it leverages `bitmapCache` and so that it's no longer reponsible for updating the count cache. This commit also helps SetBit/ClearBit performance by allowing them to work against data from `bitmapCache` instead of loading bitmaps from fragment.storage every time. --- bitmap.go | 2 -- fragment.go | 72 ++++++++++++++++++++++------------------------ roaring/roaring.go | 2 +- 3 files changed, 35 insertions(+), 41 deletions(-) diff --git a/bitmap.go b/bitmap.go index 651950e45..3a2e4c96a 100644 --- a/bitmap.go +++ b/bitmap.go @@ -16,8 +16,6 @@ type Bitmap struct { // Attributes associated with the bitmap. Attrs map[string]interface{} - - adjustedCount uint64 } // NewBitmap returns a new instance of Bitmap. diff --git a/fragment.go b/fragment.go index 8f8698192..51259585a 100644 --- a/fragment.go +++ b/fragment.go @@ -245,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() @@ -312,33 +312,35 @@ 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, }}, - adjustedCount: 0, } bm.InvalidateCount() - bm.adjustedCount = bm.Count() - if updateCache { - // Update cache. - f.cache.Add(bitmapID, bm.Count()) + if updateBitmapCache { f.bitmapCache.Add(bitmapID, bm) } @@ -353,9 +355,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 @@ -374,22 +376,18 @@ 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 } - // If adjustedCount is set, then apply that value to bitmap.n instead. - bm := f.bitmap(bitmapID, true) - if bm.adjustedCount > 0 { - bm.SetCount(profileID, bm.adjustedCount) - bm.adjustedCount = 0 - } else { - bm.IncrementCount(profileID) - f.cache.Add(bitmapID, bm.Count()) - } + // 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) return changed, nil @@ -403,7 +401,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 { @@ -411,8 +410,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 } @@ -429,17 +427,13 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { return false, err } - // If adjustedCount is set, then apply that value to bitmap.n instead. - bm := f.bitmap(bitmapID, true) - if bm.adjustedCount > 0 { - bm.SetCount(profileID, bm.adjustedCount) - bm.adjustedCount = 0 - } else { - bm.DecrementCount(profileID) - f.cache.Add(bitmapID, bm.Count()) - } + // Get the bitmap from bitmapCache or fragment.storage. + bm := f.bitmap(bitmapID, true, true) bm.ClearBit(profileID) + // Update the cache. + f.cache.Add(bitmapID, bm.Count()) + f.stats.Count("clearN", 1) return changed, nil @@ -908,7 +902,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() @@ -1256,7 +1253,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.