From 1e816b401cc65b99b3d8bfa004a20256ad0e7c97 Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 27 Jan 2017 13:16:27 -0600 Subject: [PATCH 01/14] Adds Todd's performance improvements: - Ingore asserts in fragment container. - Only log queries that take longer than 90 seconds. TODO: - address the TODOs that make the asserts configurable. --- fragment.go | 16 +++++++++++----- handler.go | 5 ++++- roaring/roaring.go | 24 ++++++++++++++---------- 3 files changed, 29 insertions(+), 16 deletions(-) diff --git a/fragment.go b/fragment.go index 1181f1739..aaef9323a 100644 --- a/fragment.go +++ b/fragment.go @@ -350,6 +350,11 @@ func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) return false, err } + // Don't update the cache if nothing changed. + if !changed { + return changed, nil + } + // Invalidate block checksum. delete(f.checksums, int(bitmapID/HashBlockSize)) @@ -389,6 +394,11 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { return false, err } + // Don't update the cache if nothing changed. + if !changed { + return changed, nil + } + // Invalidate block checksum. delete(f.checksums, int(bitmapID/HashBlockSize)) @@ -838,7 +848,6 @@ func (f *Fragment) Import(bitmapIDs, profileIDs []uint64) error { // Process every bit. // If an error occurs then reopen the storage. lastID := uint64(0) - bmCounter := 0 if err := func() error { set := make(map[uint64]struct{}) for i := range bitmapIDs { @@ -851,7 +860,7 @@ func (f *Fragment) Import(bitmapIDs, profileIDs []uint64) error { } // Write to storage. - changed, err := f.storage.Add(pos) + _, err = f.storage.Add(pos) if err != nil { return err } @@ -863,9 +872,6 @@ func (f *Fragment) Import(bitmapIDs, profileIDs []uint64) error { lastID = bitmapID set[bitmapID] = struct{}{} } - if changed { - bmCounter += 1 - } // Invalidate block checksum. delete(f.checksums, int(bitmapID/HashBlockSize)) diff --git a/handler.go b/handler.go index 175161384..9a704e6c8 100644 --- a/handler.go +++ b/handler.go @@ -202,7 +202,10 @@ func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { http.NotFound(w, r) } - h.logger().Printf("%s %s %.03fs", r.Method, r.URL.String(), time.Since(t).Seconds()) + dif := time.Since(t).Seconds() + if dif > 90 { + h.logger().Printf("%s %s %.03fs", r.Method, r.URL.String(), dif) + } } // handleGetSchema handles GET /schema requests. diff --git a/roaring/roaring.go b/roaring/roaring.go index 44b9c9341..efbf48374 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -460,8 +460,9 @@ func (b *Bitmap) WriteTo(w io.Writer) (n int64, err error) { c := b.containers[i] // Verify container count before writing. - count := c.count() - assert(c.count() == c.n, "cannot write container count, mismatch: count=%d, n=%d", count, c.n) + // TODO: instead of commenting this out, we need to make it a configuration option + //count := c.count() + //assert(c.count() == c.n, "cannot write container count, mismatch: count=%d, n=%d", count, c.n) binary.LittleEndian.PutUint64(buf[headerSize+i*12:], uint64(key)) binary.LittleEndian.PutUint32(buf[headerSize+i*12+8:], uint32(c.n-1)) @@ -532,9 +533,10 @@ func (b *Bitmap) UnmarshalBinary(data []byte) error { c := b.containers[i] if c.n <= ArrayMaxSize { c.array = (*[0xFFFFFFF]uint32)(unsafe.Pointer(&data[offset]))[:c.n] - for _, v := range c.array { - assert(lowbits(uint64(v)) == v, "array value out of range: %d", v) - } + // TODO: instead of commenting this out, we need to make it a configuration option + //for _, v := range c.array { + // assert(lowbits(uint64(v)) == v, "array value out of range: %d", v) + //} opsOffset = int(offset) + len(c.array)*4 } else { c.bitmap = (*[0xFFFFFFF]uint64)(unsafe.Pointer(&data[offset]))[:bitmapN] @@ -542,8 +544,9 @@ func (b *Bitmap) UnmarshalBinary(data []byte) error { } // Verify container count on load. - count := c.count() - assert(c.count() == c.n, "container count mismatch: count=%d, n=%d", count, c.n) + // TODO: instead of commenting this out, we need to make it a configuration option + //count := c.count() + //assert(c.count() == c.n, "container count mismatch: count=%d, n=%d", count, c.n) } // Read ops log until the end of the file. @@ -1074,9 +1077,10 @@ func (c *container) arrayWriteTo(w io.Writer) (n int64, err error) { } // Verify all elements are valid. - for _, v := range c.array { - assert(lowbits(uint64(v)) == v, "cannot write array value out of range: %d", v) - } + // TODO: instead of commenting this out, we need to make it a configuration option + //for _, v := range c.array { + // assert(lowbits(uint64(v)) == v, "cannot write array value out of range: %d", v) + //} nn, err := w.Write((*[0xFFFFFFF]byte)(unsafe.Pointer(&c.array[0]))[:4*c.n]) return int64(nn), err From e3f4eab5ac05007fead5d4e7316fe75efd3d83e4 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Fri, 27 Jan 2017 18:42:47 -0500 Subject: [PATCH 02/14] removed maxslice from setbit path --- executor.go | 32 ++++++++++++++++++++++++++------ 1 file changed, 26 insertions(+), 6 deletions(-) diff --git a/executor.go b/executor.go index 3a8152ff5..f97c35c79 100644 --- a/executor.go +++ b/executor.go @@ -52,13 +52,15 @@ func (e *Executor) Execute(ctx context.Context, db string, q *pql.Query, slices // If slices aren't specified, then include all of them. if len(slices) == 0 { - // Round up the number of slices. - maxSlice := e.Index.DB(db).MaxSlice() + if needsSlices(q.Calls){ + // Round up the number of slices. +maxSlice := e.Index.DB(db).MaxSlice() - // Generate a slices of all slices. - slices = make([]uint64, maxSlice+1) - for i := range slices { - slices[i] = uint64(i) + // Generate a slices of all slices. + slices = make([]uint64, maxSlice+1) + for i := range slices { + slices[i] = uint64(i) + } } } @@ -876,3 +878,21 @@ func hasOnlySetBitmapAttrs(calls pql.Calls) bool { } return true } + +func needsSlices(calls pql.Calls) bool { + if len(calls) == 0 { + return false + } + + for _, call := range calls { + if _, ok := call.(pql.BitmapCall); !ok { + return true + }else if _, ok := call.(*pql.Count); !ok { + return true + }else if _, ok := call.(*pql.TopN); !ok { + return true + } + + } + return false +} From 7105c89398904a7a429fab1c9190eeceabb2e413 Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 27 Jan 2017 17:48:44 -0600 Subject: [PATCH 03/14] run gofmt on the previous commit --- executor.go | 26 +++++++++++++------------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/executor.go b/executor.go index f97c35c79..6f213bafa 100644 --- a/executor.go +++ b/executor.go @@ -52,15 +52,15 @@ func (e *Executor) Execute(ctx context.Context, db string, q *pql.Query, slices // If slices aren't specified, then include all of them. if len(slices) == 0 { - if needsSlices(q.Calls){ + if needsSlices(q.Calls) { // Round up the number of slices. -maxSlice := e.Index.DB(db).MaxSlice() + maxSlice := e.Index.DB(db).MaxSlice() - // Generate a slices of all slices. - slices = make([]uint64, maxSlice+1) - for i := range slices { - slices[i] = uint64(i) - } + // Generate a slices of all slices. + slices = make([]uint64, maxSlice+1) + for i := range slices { + slices[i] = uint64(i) + } } } @@ -880,19 +880,19 @@ func hasOnlySetBitmapAttrs(calls pql.Calls) bool { } func needsSlices(calls pql.Calls) bool { - if len(calls) == 0 { - return false - } + if len(calls) == 0 { + return false + } for _, call := range calls { if _, ok := call.(pql.BitmapCall); !ok { return true - }else if _, ok := call.(*pql.Count); !ok { + } else if _, ok := call.(*pql.Count); !ok { return true - }else if _, ok := call.(*pql.TopN); !ok { + } else if _, ok := call.(*pql.TopN); !ok { return true } } - return false + return false } From 61da4a1001355875382817cb94df3cdd0bf882db Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 1 Feb 2017 13:40:12 -0500 Subject: [PATCH 04/14] limit the number of connections to single host --- cmd/pilosa/main.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/cmd/pilosa/main.go b/cmd/pilosa/main.go index 24d2a9fec..8513a5ce3 100644 --- a/cmd/pilosa/main.go +++ b/cmd/pilosa/main.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "math/rand" + "net/http" "os" "os/signal" "path/filepath" @@ -34,6 +35,7 @@ const ( ) func main() { + http.DefaultTransport.(*http.Transport).MaxIdleConnsPerHost = 64 m := NewMain() fmt.Fprintf(m.Stderr, "Pilosa %s\n", Build) From a35dd5e5861db6d9059dffa4963b7cd673692ebd Mon Sep 17 00:00:00 2001 From: Travis Date: Wed, 1 Feb 2017 12:57:06 -0600 Subject: [PATCH 05/14] add comments to the MaxIdleConnsPerHost tweak --- cmd/pilosa/main.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/cmd/pilosa/main.go b/cmd/pilosa/main.go index 8513a5ce3..4d19ef808 100644 --- a/cmd/pilosa/main.go +++ b/cmd/pilosa/main.go @@ -35,7 +35,10 @@ const ( ) func main() { + // Limit the number of connections that a server can make to a single node + // in order to prevent a cluster storm. http.DefaultTransport.(*http.Transport).MaxIdleConnsPerHost = 64 + m := NewMain() fmt.Fprintf(m.Stderr, "Pilosa %s\n", Build) From 631b3bc44c326cc39ea5aeb026e01dc68be39257 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 1 Feb 2017 16:29:48 -0500 Subject: [PATCH 06/14] added handling for float attribute --- attr.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/attr.go b/attr.go index 1425676cf..c85f77cad 100644 --- a/attr.go +++ b/attr.go @@ -265,6 +265,8 @@ func txUpdateAttrs(tx *bolt.Tx, id uint64, m map[string]interface{}) (map[string attr[k] = uint64(v) case uint: attr[k] = uint64(v) + case float64: + attr[k] = uint64(v) case int64: attr[k] = uint64(v) case string, uint64, bool: From e62bd5fad56f98a5b154b68b7c4839ec189b524b Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 16 Feb 2017 15:43:27 -0500 Subject: [PATCH 07/14] WIP snapshot optimization --- roaring/roaring.go | 48 ++++++++++++++++++++++++++++++++++------------ 1 file changed, 36 insertions(+), 12 deletions(-) diff --git a/roaring/roaring.go b/roaring/roaring.go index efbf48374..d714577da 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -444,17 +444,30 @@ func (b *Bitmap) removeEmptyContainers() { i++ } } +func (b *Bitmap) countEmptyContainers() int { + result := 0 + for i := 0; i < len(b.containers); { + c := b.containers[i] + + if c.n == 0 { + result++ + } + i++ + } + return result +} // WriteTo writes b to w. func (b *Bitmap) WriteTo(w io.Writer) (n int64, err error) { // Remove empty containers before persisting. - b.removeEmptyContainers() + //b.removeEmptyContainers() + containerCount := len(b.keys) - b.countEmptyContainers() // Build header before writing individual container blocks. - buf := make([]byte, headerSize+(len(b.keys)*(4+8+4))) + buf := make([]byte, headerSize+(containerCount*(4+8+4))) binary.LittleEndian.PutUint32(buf[0:], cookie) - binary.LittleEndian.PutUint32(buf[4:], uint32(len(b.keys))) - + binary.LittleEndian.PutUint32(buf[4:], uint32(containerCount)) + empty:=0 // Encode keys and cardinality. for i, key := range b.keys { c := b.containers[i] @@ -463,15 +476,24 @@ func (b *Bitmap) WriteTo(w io.Writer) (n int64, err error) { // TODO: instead of commenting this out, we need to make it a configuration option //count := c.count() //assert(c.count() == c.n, "cannot write container count, mismatch: count=%d, n=%d", count, c.n) - - binary.LittleEndian.PutUint64(buf[headerSize+i*12:], uint64(key)) - binary.LittleEndian.PutUint32(buf[headerSize+i*12+8:], uint32(c.n-1)) + if c.n > 0 { + binary.LittleEndian.PutUint64(buf[headerSize+(i-empty)*12:], uint64(key)) + binary.LittleEndian.PutUint32(buf[headerSize+(i-empty)*12+8:], uint32(c.n-1)) + }else{ + empty++ + } } // Write the offset for each container block. offset := uint32(len(buf)) + empty=0 for i, c := range b.containers { - binary.LittleEndian.PutUint32(buf[headerSize+(len(b.keys)*12)+(i*4):], uint32(offset)) + + if c.n > 0 { + binary.LittleEndian.PutUint32(buf[headerSize+(containerCount*12)+((i-empty)*4):], uint32(offset)) + }else{ + empty++ + } offset += uint32(c.size()) } @@ -484,10 +506,12 @@ func (b *Bitmap) WriteTo(w io.Writer) (n int64, err error) { // Write each container block. for _, c := range b.containers { - nn, err := c.WriteTo(w) - n += nn - if err != nil { - return n, err + if c.n > 0 { + nn, err := c.WriteTo(w) + n += nn + if err != nil { + return n, err + } } } From 9a506fa05477c4e227b739726ab6f731e63dfab7 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 16 Feb 2017 15:44:32 -0500 Subject: [PATCH 08/14] added max slice to schema --- db.go | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/db.go b/db.go index 2a4f1d6e0..60a018de4 100644 --- a/db.go +++ b/db.go @@ -426,8 +426,9 @@ func (p dbSlice) Less(i, j int) bool { return p[i].Name() < p[j].Name() } // DBInfo represents schema information for a database. type DBInfo struct { - Name string `json:"name"` - Frames []*FrameInfo `json:"frames"` + Name string `json:"name"` + Frames []*FrameInfo `json:"frames"` + MaxSlice uint64 } type dbInfoSlice []*DBInfo From 667d56eedf0b51e7c0e7541a17055ce50e3b7f7a Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 16 Feb 2017 15:54:47 -0500 Subject: [PATCH 09/14] added simple cache to fragment --- fragment.go | 36 ++++++++++++++++++++++++++++++++++-- 1 file changed, 34 insertions(+), 2 deletions(-) diff --git a/fragment.go b/fragment.go index aaef9323a..88b39663c 100644 --- a/fragment.go +++ b/fragment.go @@ -28,7 +28,8 @@ import ( const ( // SliceWidth is the number of profile IDs in a slice. - SliceWidth = 1048576 + //SliceWidth = 1048576 + SliceWidth = 262144 // SnapshotExt is the file extension used for an in-process snapshot. SnapshotExt = ".snapshotting" @@ -49,8 +50,25 @@ const ( const ( // DefaultFragmentMaxOpN is the default value for Fragment.MaxOpN. - DefaultFragmentMaxOpN = 1000 + //TODO CHANGING FOR TEST TO 10x + DefaultFragmentMaxOpN = 2000 ) +type BitmapCacher interface { + Fetch(id uint64)(*Bitmap,bool) + Add(id uint64, b *Bitmap) +} + + +type Simple struct{ + cache map[uint64]*Bitmap +} +func (s *Simple)Fetch(id uint64)(*Bitmap,bool){ + m,ok:=s.cache[id] + return m,ok +} +func (s *Simple)Add(id uint64,p*Bitmap){ + s.cache[id]=p +} // Fragment represents the intersection of a frame and slice in a database. type Fragment struct { @@ -87,6 +105,8 @@ type Fragment struct { BitmapAttrStore *AttrStore stats StatsClient + turbo BitmapCacher + } // NewFragment returns a new instance of Fragment. @@ -203,6 +223,7 @@ func (f *Fragment) openStorage() error { // Attach the file to the bitmap to act as a write-ahead log. f.storage.OpWriter = f.file + f.turbo = &Simple{make(map[uint64]*Bitmap)} return nil @@ -308,7 +329,12 @@ func (f *Fragment) Bitmap(bitmapID uint64) *Bitmap { return f.bitmap(bitmapID) } + func (f *Fragment) bitmap(bitmapID uint64) *Bitmap { + r,ok:=f.turbo.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) @@ -325,6 +351,7 @@ func (f *Fragment) bitmap(bitmapID uint64) *Bitmap { // Update cache. f.cache.Add(bitmapID, bm.Count()) + f.turbo.Add(bitmapID, bm) return bm } @@ -918,10 +945,15 @@ func (f *Fragment) Snapshot() error { defer f.mu.Unlock() return f.snapshot() } +func track(start time.Time, name string, logger *log.Logger) { + elapsed := time.Since(start) + logger.Printf("%s took %s", name, elapsed) +} func (f *Fragment) snapshot() error { logger := f.logger() logger.Printf("fragment: snapshotting %s/%s/%d", f.db, f.frame, f.slice) + defer track(time.Now(), fmt.Sprintf("fragment: snapshot complete %s/%s/%d", f.db, f.frame, f.slice), logger) // Create a temporary file to snapshot to. snapshotPath := f.path + SnapshotExt From 028539f10da8d27e0aa81d27e79214322c357544 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Fri, 17 Feb 2017 15:18:38 -0500 Subject: [PATCH 10/14] top bug in cache refresh --- cache.go | 16 +++++++++++----- fragment.go | 2 +- 2 files changed, 12 insertions(+), 6 deletions(-) diff --git a/cache.go b/cache.go index a1cba3dec..fb5f4defd 100644 --- a/cache.go +++ b/cache.go @@ -22,6 +22,7 @@ type Cache interface { // Updates the cache, if necessary. Invalidate() + Refresh() // Returns an ordered list of the top ranked bitmaps. Top() []BitmapPair @@ -86,6 +87,8 @@ func (c *LRUCache) Top() []BitmapPair { } func (c *LRUCache) onEvicted(key lru.Key, _ interface{}) { delete(c.counts, key.(uint64)) } +func (c *LRUCache) Refresh() { +} // Ensure LRUCache implements Cache. var _ Cache = &LRUCache{} @@ -117,17 +120,17 @@ func (c *RankCache) Add(bitmapID uint64, n uint64) { return } - // Add to cache. c.entries[bitmapID] = n + // If size is larger than the threshold then trim it. if len(c.entries) > c.ThresholdLength { c.update() - for id, n := range c.entries { - if n <= c.ThresholdValue { - delete(c.entries, id) + for id, n := range c.entries { + if n <= c.ThresholdValue { + delete(c.entries, id) + } } - } } } @@ -154,6 +157,9 @@ func (c *RankCache) Invalidate() { c.update() } } +func (c *RankCache) Refresh() { + c.update() +} // update reorders the entries by rank. func (c *RankCache) update() { diff --git a/fragment.go b/fragment.go index 88b39663c..51cb1dc61 100644 --- a/fragment.go +++ b/fragment.go @@ -909,7 +909,7 @@ func (f *Fragment) Import(bitmapIDs, profileIDs []uint64) error { f.cache.Add(bitmapID, f.bitmap(bitmapID).Count()) } - f.cache.Invalidate() + f.cache.Refresh() return nil }(); err != nil { _ = f.closeStorage() From d66e368843807d2307e90ac3e3c1ad17dba4f1fa Mon Sep 17 00:00:00 2001 From: Travis Date: Mon, 20 Feb 2017 18:02:24 -0600 Subject: [PATCH 11/14] handle `unionArrayBitmap` width of bitmaps correctly --- roaring/roaring.go | 1 + 1 file changed, 1 insertion(+) diff --git a/roaring/roaring.go b/roaring/roaring.go index efbf48374..c3e50a940 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -1368,6 +1368,7 @@ func unionArrayBitmap(a, b *container) *container { break } else if i >= len(a.array) { output.add(vb) + continue } else if eof { output.add(a.array[i]) i++ From 43ede4d9c8efe07e05acbc557d1fa61d5abebb5c Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Tue, 21 Feb 2017 09:28:01 -0600 Subject: [PATCH 12/14] WIP increase caching and topn adjustments --- cache.go | 42 ++++++++++++++-------------- fragment.go | 69 ++++++++++++++++++++++++++++------------------ roaring/roaring.go | 13 +++++---- 3 files changed, 71 insertions(+), 53 deletions(-) diff --git a/cache.go b/cache.go index fb5f4defd..c6fc48ec6 100644 --- a/cache.go +++ b/cache.go @@ -14,6 +14,7 @@ import ( // Cache represents a cache for bitmap counts. type Cache interface { Add(bitmapID uint64, n uint64) + BulkAdd(bitmapID uint64, n uint64) Get(bitmapID uint64) uint64 Len() int @@ -22,7 +23,6 @@ type Cache interface { // Updates the cache, if necessary. Invalidate() - Refresh() // Returns an ordered list of the top ranked bitmaps. Top() []BitmapPair @@ -44,6 +44,10 @@ func NewLRUCache(maxEntries int) *LRUCache { return c } +func (c *LRUCache) BulkAdd(bitmapID, n uint64) { + c.Add(bitmapID, n) +} + // Add adds a bitmap to the cache. func (c *LRUCache) Add(bitmapID, n uint64) { c.cache.Add(bitmapID, n) @@ -87,7 +91,7 @@ func (c *LRUCache) Top() []BitmapPair { } func (c *LRUCache) onEvicted(key lru.Key, _ interface{}) { delete(c.counts, key.(uint64)) } -func (c *LRUCache) Refresh() { +func (c *LRUCache) Refresh() { } // Ensure LRUCache implements Cache. @@ -122,18 +126,26 @@ 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 { - c.update() - for id, n := range c.entries { - if n <= c.ThresholdValue { - delete(c.entries, id) - } + for id, n := range c.entries { + if n <= c.ThresholdValue { + delete(c.entries, id) } + } } } +// Add adds a bitmap to the cache unsorted you should Invalidate after completion +func (c *RankCache) BulkAdd(bitmapID uint64, n uint64) { + if n < c.ThresholdValue { + return + } + + c.entries[bitmapID] = n +} + // Get returns a bitmap with a given id. func (c *RankCache) Get(bitmapID uint64) uint64 { return c.entries[bitmapID] } @@ -150,19 +162,9 @@ func (c *RankCache) BitmapIDs() []uint64 { return a } -// Invalidate reorders the entries, if necessary. -func (c *RankCache) Invalidate() { - // Update if there aren't many items or it hasn't been updated recently. - if len(c.rankings) < 50 || (c.updateN > 0 && time.Since(c.updateTime) > 5*time.Minute) { - c.update() - } -} -func (c *RankCache) Refresh() { - c.update() -} - // update reorders the entries by rank. -func (c *RankCache) update() { +func (c *RankCache) Invalidate() { + //fmt.Println("RankCache Update") // Convert cache to a sorted list. rankings := make([]BitmapPair, 0, len(c.entries)) for id, n := range c.entries { diff --git a/fragment.go b/fragment.go index 51cb1dc61..f16bf6ddb 100644 --- a/fragment.go +++ b/fragment.go @@ -53,21 +53,22 @@ const ( //TODO CHANGING FOR TEST TO 10x DefaultFragmentMaxOpN = 2000 ) + type BitmapCacher interface { - Fetch(id uint64)(*Bitmap,bool) + Fetch(id uint64) (*Bitmap, bool) Add(id uint64, b *Bitmap) } +type Simple struct { + cache map[uint64]*Bitmap +} -type Simple struct{ - cache map[uint64]*Bitmap +func (s *Simple) Fetch(id uint64) (*Bitmap, bool) { + m, ok := s.cache[id] + return m, ok } -func (s *Simple)Fetch(id uint64)(*Bitmap,bool){ - m,ok:=s.cache[id] - return m,ok -} -func (s *Simple)Add(id uint64,p*Bitmap){ - s.cache[id]=p +func (s *Simple) Add(id uint64, p *Bitmap) { + s.cache[id] = p } // Fragment represents the intersection of a frame and slice in a database. @@ -105,8 +106,7 @@ type Fragment struct { BitmapAttrStore *AttrStore stats StatsClient - turbo BitmapCacher - + turbo BitmapCacher } // NewFragment returns a new instance of Fragment. @@ -173,6 +173,7 @@ func (f *Fragment) Open() error { // openStorage opens the storage bitmap. func (f *Fragment) openStorage() error { + //f.logger().Printf("Open Storage %s/%s/%d", f.db, f.frame, f.slice) // Create a roaring bitmap to serve as storage for the slice. f.storage = roaring.NewBitmap() @@ -260,9 +261,11 @@ func (f *Fragment) openCache() error { // Read in all bitmaps by ID. // This will cause them to be added to the cache. for _, bitmapID := range pb.BitmapIDs { - n := f.storage.CountRange(bitmapID*SliceWidth, (bitmapID+1)*SliceWidth) - f.cache.Add(bitmapID, n) + //n := f.storage.CountRange(bitmapID*SliceWidth, (bitmapID+1)*SliceWidth) + n := f.bitmap(bitmapID, false).Count() + f.cache.BulkAdd(bitmapID, n) } + f.cache.Invalidate() return nil } @@ -326,15 +329,14 @@ 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) + return f.bitmap(bitmapID, false) } - -func (f *Fragment) bitmap(bitmapID uint64) *Bitmap { - r,ok:=f.turbo.Fetch(bitmapID) - if ok && r != nil{ - return r - } +func (f *Fragment) bitmap(bitmapID uint64, updateCache bool) *Bitmap { + r, ok := f.turbo.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) @@ -349,9 +351,11 @@ func (f *Fragment) bitmap(bitmapID uint64) *Bitmap { } bm.InvalidateCount() - // Update cache. - f.cache.Add(bitmapID, bm.Count()) - f.turbo.Add(bitmapID, bm) + if updateCache { + // Update cache. + f.cache.Add(bitmapID, bm.Count()) + f.turbo.Add(bitmapID, bm) + } return bm } @@ -391,7 +395,7 @@ func (f *Fragment) setBit(bitmapID, profileID uint64) (changed bool, bool error) } // Update the cache. - if f.bitmap(bitmapID).SetBit(profileID) { + if f.bitmap(bitmapID, true).SetBit(profileID) { changed = true } @@ -435,7 +439,7 @@ func (f *Fragment) clearBit(bitmapID, profileID uint64) (bool, error) { } // Update the cache. - if f.bitmap(bitmapID).ClearBit(profileID) { + if f.bitmap(bitmapID, true).ClearBit(profileID) { return true, nil } @@ -570,7 +574,16 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { return results, nil } +func debugDumpPairs(pairs []BitmapPair) { + fmt.Println("=====Start") + for i, pair := range pairs { + fmt.Println(i, pair.ID, pair.Count) + } + fmt.Println("=====Stop") +} + func (f *Fragment) topBitmapPairs(bitmapIDs []uint64) []BitmapPair { + //fmt.Println("DEBUG topBitmapPairs") // If no specific bitmaps are requested, retrieve top bitmaps. if len(bitmapIDs) == 0 { f.mu.Lock() @@ -597,6 +610,8 @@ func (f *Fragment) topBitmapPairs(bitmapIDs []uint64) []BitmapPair { Count: f.Bitmap(bitmapID).Count(), } } + sort.Sort(BitmapPairs(pairs)) + //debugDumpPairs(pairs) return pairs } @@ -906,10 +921,10 @@ func (f *Fragment) Import(bitmapIDs, profileIDs []uint64) error { // Update cache counts for all bitmaps. for bitmapID := range set { - f.cache.Add(bitmapID, f.bitmap(bitmapID).Count()) + f.cache.BulkAdd(bitmapID, f.bitmap(bitmapID, false).Count()) } - f.cache.Refresh() + f.cache.Invalidate() return nil }(); err != nil { _ = f.closeStorage() diff --git a/roaring/roaring.go b/roaring/roaring.go index d714577da..58b4d1e19 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -467,7 +467,7 @@ func (b *Bitmap) WriteTo(w io.Writer) (n int64, err error) { buf := make([]byte, headerSize+(containerCount*(4+8+4))) binary.LittleEndian.PutUint32(buf[0:], cookie) binary.LittleEndian.PutUint32(buf[4:], uint32(containerCount)) - empty:=0 + empty := 0 // Encode keys and cardinality. for i, key := range b.keys { c := b.containers[i] @@ -479,20 +479,20 @@ func (b *Bitmap) WriteTo(w io.Writer) (n int64, err error) { if c.n > 0 { binary.LittleEndian.PutUint64(buf[headerSize+(i-empty)*12:], uint64(key)) binary.LittleEndian.PutUint32(buf[headerSize+(i-empty)*12+8:], uint32(c.n-1)) - }else{ + } else { empty++ } } // Write the offset for each container block. offset := uint32(len(buf)) - empty=0 + empty = 0 for i, c := range b.containers { if c.n > 0 { - binary.LittleEndian.PutUint32(buf[headerSize+(containerCount*12)+((i-empty)*4):], uint32(offset)) - }else{ - empty++ + binary.LittleEndian.PutUint32(buf[headerSize+(containerCount*12)+((i-empty)*4):], uint32(offset)) + } else { + empty++ } offset += uint32(c.size()) } @@ -1392,6 +1392,7 @@ func unionArrayBitmap(a, b *container) *container { break } else if i >= len(a.array) { output.add(vb) + continue } else if eof { output.add(a.array[i]) i++ From 774dc14412c4a7205f5448d9cd6a32d6c09dd746 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 22 Feb 2017 18:11:06 -0600 Subject: [PATCH 13/14] handle empty Union/Intersect --- executor.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/executor.go b/executor.go index 6f213bafa..bfe1cccf1 100644 --- a/executor.go +++ b/executor.go @@ -281,7 +281,7 @@ func (e *Executor) executeBitmapSlice(ctx context.Context, db string, c *pql.Bit // executeIntersectSlice executes a intersect() call for a local slice. func (e *Executor) executeIntersectSlice(ctx context.Context, db string, c *pql.Intersect, slice uint64) (*Bitmap, error) { - var other *Bitmap + other := &Bitmap{} for i, input := range c.Inputs { bm, err := e.executeBitmapCallSlice(ctx, db, input, slice) if err != nil { @@ -331,7 +331,7 @@ func (e *Executor) executeRangeSlice(ctx context.Context, db string, c *pql.Rang // executeUnionSlice executes a union() call for a local slice. func (e *Executor) executeUnionSlice(ctx context.Context, db string, c *pql.Union, slice uint64) (*Bitmap, error) { - var other *Bitmap + other := &Bitmap{} for i, input := range c.Inputs { bm, err := e.executeBitmapCallSlice(ctx, db, input, slice) if err != nil { From ff52ffa0473a02e68433d11c7f7ac038886d7a6f Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 23 Feb 2017 15:08:18 -0600 Subject: [PATCH 14/14] WIP TopN optimization --- cache.go | 18 ++++++++++++++++++ fragment.go | 47 +++++++++++++++++++++-------------------------- 2 files changed, 39 insertions(+), 26 deletions(-) diff --git a/cache.go b/cache.go index c6fc48ec6..2f4761872 100644 --- a/cache.go +++ b/cache.go @@ -241,6 +241,24 @@ func (p Pairs) Swap(i, j int) { p[i], p[j] = p[j], p[i] } func (p Pairs) Len() int { return len(p) } func (p Pairs) Less(i, j int) bool { return p[i].Count > p[j].Count } +type PairHeap struct { + Pairs +} + +func (h *Pairs) Push(x interface{}) { + // Push and Pop use pointer receivers because they modify the slice's length, + // not just its contents. + *h = append(*h, x.(Pair)) +} + +func (h *Pairs) Pop() interface{} { + old := *h + n := len(old) + x := old[n-1] + *h = old[0 : n-1] + return x +} + // Add merges other into p and returns a new slice. func (p Pairs) Add(other []Pair) []Pair { // Create lookup of key/counts. diff --git a/fragment.go b/fragment.go index f16bf6ddb..e7d6086d1 100644 --- a/fragment.go +++ b/fragment.go @@ -4,6 +4,7 @@ import ( "archive/tar" "bufio" "bytes" + "container/heap" "context" "crypto/sha1" "encoding/binary" @@ -28,7 +29,7 @@ import ( const ( // SliceWidth is the number of profile IDs in a slice. - //SliceWidth = 1048576 + // SliceWidth = 1048576 SliceWidth = 262144 // SnapshotExt is the file extension used for an in-process snapshot. @@ -173,7 +174,6 @@ func (f *Fragment) Open() error { // openStorage opens the storage bitmap. func (f *Fragment) openStorage() error { - //f.logger().Printf("Open Storage %s/%s/%d", f.db, f.frame, f.slice) // Create a roaring bitmap to serve as storage for the slice. f.storage = roaring.NewBitmap() @@ -494,7 +494,8 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { } // Iterate over rankings and add to results until we have enough. - results := make([]Pair, 0, opt.N) + //results := make(PairHeap, 0, opt.N) + results := &PairHeap{} for _, pair := range pairs { bitmapID, n := pair.ID, pair.Count @@ -518,7 +519,7 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { } // The initial n pairs should simply be added to the results. - if opt.N == 0 || len(results) < opt.N { + if opt.N == 0 || results.Len() < opt.N { // Calculate count and append. count := n if opt.Src != nil { @@ -527,23 +528,26 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { if count == 0 { continue } - results = append(results, Pair{Key: bitmapID, Count: count}) + //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 // intersections then simply exit. If we are intersecting then sort // and then only keep pairs that are higher than the lowest count. - if opt.N > 0 && len(results) == opt.N { + if opt.N > 0 && results.Len() == opt.N { if opt.Src == nil { break } - sort.Sort(Pairs(results)) + // 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[len(results)-1].Count + + threshold := results.Pairs[0].Count if threshold < MinThreshold { break } @@ -556,34 +560,25 @@ func (f *Fragment) Top(opt TopOptions) ([]Pair, error) { // Calculate the intersecting bit count and skip if it's below our // last bitmap in our current result set. + count := opt.Src.IntersectionCount(f.Bitmap(bitmapID)) if count < threshold { continue } - // Swap out the last pair for this new count. - results[len(results)-1] = Pair{Key: bitmapID, Count: count} - - // If it's count is also higher than the second to last item then resort. - if len(results) >= 2 && count > results[len(results)-2].Count { - sort.Sort(Pairs(results)) - } + heap.Push(results, Pair{Key: bitmapID, Count: count}) } - - sort.Sort(Pairs(results)) - return results, nil -} - -func debugDumpPairs(pairs []BitmapPair) { - fmt.Println("=====Start") - for i, pair := range pairs { - fmt.Println(i, pair.ID, pair.Count) + r := make(Pairs, results.Len(), results.Len()) + x := results.Len() + i := 1 + for results.Len() > 0 { + r[x-i] = heap.Pop(results).(Pair) + i++ } - fmt.Println("=====Stop") + return r, nil } func (f *Fragment) topBitmapPairs(bitmapIDs []uint64) []BitmapPair { - //fmt.Println("DEBUG topBitmapPairs") // If no specific bitmaps are requested, retrieve top bitmaps. if len(bitmapIDs) == 0 { f.mu.Lock()