From 3fdf63f84fcf29074a701d9adf4bcec00618cc55 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 7 Jul 2016 10:03:38 -0600 Subject: [PATCH] optimize sparse bitmap block checksums This commit refactors the block checksumming by removing the iteration over each block and instead only checking blocks which have data. This requires merging the cache inspection with the roaring iterator to reduce CPU time and memory allocations. --- .gitignore | 1 + fragment.go | 140 ++++++++++++++++++++++++++++++++--------------- fragment_test.go | 69 ++++++++++++++++------- handler_test.go | 4 +- 4 files changed, 148 insertions(+), 66 deletions(-) diff --git a/.gitignore b/.gitignore index 4d4ef0e36..25442a3da 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,2 @@ default.etcd/ +*.test diff --git a/fragment.go b/fragment.go index 4c9be2782..7be7993e1 100644 --- a/fragment.go +++ b/fragment.go @@ -8,6 +8,7 @@ import ( "encoding/binary" "errors" "fmt" + "hash" "io" "io/ioutil" "log" @@ -594,8 +595,8 @@ func (f *Fragment) Range(bitmapID uint64, start, end time.Time) *Bitmap { // If two fragments have the same checksum then they have the same data. func (f *Fragment) Checksum() []byte { h := sha1.New() - for i, blockN := 0, f.BlockN(); i < blockN; i++ { - h.Write(f.BlockChecksum(i)) + for _, block := range f.Blocks() { + h.Write(block.Checksum) } return h.Sum(nil) } @@ -607,42 +608,6 @@ func (f *Fragment) BlockN() int { return int(f.storage.Max() / (HashBlockSize * SliceWidth)) } -// BlockChecksum returns the checksum for a single block in the fragment. -// Returns nil if there is no data for the block. -func (f *Fragment) BlockChecksum(i int) []byte { - f.mu.Lock() - defer f.mu.Unlock() - - // Use the cached checksum, if available. - if chksum, ok := f.checksums[i]; ok { - return chksum - } - - // Otherwise calculate the checksum from the data on disk. - h := sha1.New() - var written bool - f.storage.ForEachRange(uint64(i)*HashBlockSize*SliceWidth, (uint64(i)+1)*HashBlockSize*SliceWidth, func(i uint64) { - // Write value to the hash. - var buf [8]byte - binary.BigEndian.PutUint64(buf[:], i) - h.Write(buf[:]) - - // Mark the block has having data. - written = true - }) - - // If no data was written then return a nil checksum. - if !written { - return nil - } - - // Cache checksum for later use. - chksum := h.Sum(nil)[:] - f.checksums[i] = chksum - - return chksum -} - // InvalidateChecksums clears all cached block checksums. func (f *Fragment) InvalidateChecksums() { f.mu.Lock() @@ -652,19 +617,83 @@ func (f *Fragment) InvalidateChecksums() { // Blocks returns info for all blocks containing data. func (f *Fragment) Blocks() []FragmentBlock { + f.mu.Lock() + defer f.mu.Unlock() + var a []FragmentBlock - for i, blockN := 0, f.BlockN(); i <= blockN; i++ { - chksum := f.BlockChecksum(i) - if chksum == nil { + + // Initialize the iterator. + itr := f.storage.Iterator() + itr.Seek(0) + + // Initialize block hasher. + h := newBlockHasher() + + // Iterate over each value in the fragment. + v, eof := itr.Next() + if eof { + return nil + } + blockID := int(v / (HashBlockSize * SliceWidth)) + for { + // Check for multiple block checksums in a row. + if n := f.readContiguousChecksums(&a, blockID); n > 0 { + itr.Seek(uint64(blockID+n) * HashBlockSize * SliceWidth) + v, eof = itr.Next() + if eof { + break + } + blockID = int(v / (HashBlockSize * SliceWidth)) continue } + // Reset hasher. + h.blockID = blockID + h.Reset() + + // Read all values for the block. + for ; ; v, eof = itr.Next() { + // Once we hit the next block, save the value for the next iteration. + blockID = int(v / (HashBlockSize * SliceWidth)) + if blockID != h.blockID || eof { + break + } + + h.WriteValue(v) + } + + // Cache checksum. + chksum := h.Sum() + f.checksums[h.blockID] = chksum + + // Append block. a = append(a, FragmentBlock{ - ID: i, + ID: h.blockID, + Checksum: chksum, + }) + + // Exit if we're at the end. + if eof { + break + } + } + + return a +} + +// readContiguousChecksums appends multiple checksums in a row and returns the count added. +func (f *Fragment) readContiguousChecksums(a *[]FragmentBlock, blockID int) (n int) { + for i := 0; ; i++ { + chksum := f.checksums[blockID+i] + if chksum == nil { + return i + } + + *a = append(*a, FragmentBlock{ + ID: blockID + i, Checksum: chksum, }) } - return a } // BlockData returns bits in a block as bitmap & profile ID pairs. @@ -1164,6 +1193,31 @@ type FragmentBlock struct { Checksum []byte `json:"checksum"` } +type blockHasher struct { + blockID int + buf [8]byte + hash hash.Hash +} + +func newBlockHasher() blockHasher { + return blockHasher{ + blockID: -1, + hash: sha1.New(), + } +} +func (h *blockHasher) Reset() { + h.hash.Reset() +} + +func (h *blockHasher) Sum() []byte { + return h.hash.Sum(nil)[:] +} + +func (h *blockHasher) WriteValue(v uint64) { + binary.BigEndian.PutUint64(h.buf[:], v) + h.hash.Write(h.buf[:]) +} + // FragmentSyncer syncs a local fragment to one on a remote host. type FragmentSyncer struct { Fragment *Fragment diff --git a/fragment_test.go b/fragment_test.go index 54b43f7f5..6eeeb0fe8 100644 --- a/fragment_test.go +++ b/fragment_test.go @@ -2,6 +2,7 @@ package pilosa_test import ( "bytes" + "flag" "io/ioutil" "os" "reflect" @@ -11,6 +12,11 @@ import ( "github.com/umbel/pilosa" ) +// Test flags +var ( + FragmentPath = flag.String("fragment", "", "fragment path") +) + // SliceWidth is a helper reference to use when testing. const SliceWidth = pilosa.SliceWidth @@ -261,60 +267,58 @@ func TestFragment_Checksum(t *testing.T) { } // Ensure fragment can return a checksum for a given block. -func TestFragment_BlockChecksum(t *testing.T) { +func TestFragment_Blocks(t *testing.T) { f := MustOpenFragment("d", "f", 0) defer f.Close() // Retrieve initial checksum. - var chksum []byte - prev := f.Checksum() + var prev []pilosa.FragmentBlock // Set first bit. if _, err := f.SetBit(0, 0, nil, 0); err != nil { t.Fatal(err) } - chksum = f.BlockChecksum(0) - if bytes.Equal(chksum, prev) { - t.Fatalf("expected checksum to change: %x", chksum) + blocks := f.Blocks() + if blocks[0].Checksum == nil { + t.Fatalf("expected checksum: %x", blocks[0].Checksum) } - prev = chksum + prev = blocks // Set bit on different bitmap. if _, err := f.SetBit(20, 0, nil, 0); err != nil { t.Fatal(err) } - chksum = f.BlockChecksum(0) - if bytes.Equal(chksum, prev) { - t.Fatalf("expected checksum to change: %x", chksum) + blocks = f.Blocks() + if bytes.Equal(blocks[0].Checksum, prev[0].Checksum) { + t.Fatalf("expected checksum to change: %x", blocks[0].Checksum) } - prev = chksum + prev = blocks // Set bit on different profile. if _, err := f.SetBit(20, 100, nil, 0); err != nil { t.Fatal(err) } - chksum = f.BlockChecksum(0) - if bytes.Equal(chksum, prev) { - t.Fatalf("expected checksum to change: %x", chksum) + blocks = f.Blocks() + if bytes.Equal(blocks[0].Checksum, prev[0].Checksum) { + t.Fatalf("expected checksum to change: %x", blocks[0].Checksum) } } // Ensure fragment returns an empty checksum if no data exists for a block. -func TestFragment_BlockChecksum_Empty(t *testing.T) { +func TestFragment_Blocks_Empty(t *testing.T) { f := MustOpenFragment("d", "f", 0) defer f.Close() // Set bits on a different block. - if _, err := f.SetBit(1, 200, nil, 0); err != nil { + if _, err := f.SetBit(100, 1, nil, 0); err != nil { t.Fatal(err) } // Ensure checksum for block 1 is blank. - if chksum := f.BlockChecksum(0); chksum == nil { - t.Fatalf("expected chksum(0)") - } - if chksum := f.BlockChecksum(1); chksum != nil { - t.Fatalf("expected empty checksum: %x", chksum) + if blocks := f.Blocks(); len(blocks) != 1 { + t.Fatalf("unexpected block count: %d", len(blocks)) + } else if blocks[0].ID != 1 { + t.Fatalf("unexpected block id: %d", blocks[0].ID) } } @@ -436,6 +440,7 @@ func TestFragment_WriteTo_ReadFrom(t *testing.T) { } } +/* func BenchmarkFragment_BlockChecksum_Fill1(b *testing.B) { benchmarkFragmentBlockChecksum(b, 0.01) } func BenchmarkFragment_BlockChecksum_Fill10(b *testing.B) { benchmarkFragmentBlockChecksum(b, 0.10) } func BenchmarkFragment_BlockChecksum_Fill50(b *testing.B) { benchmarkFragmentBlockChecksum(b, 0.50) } @@ -462,6 +467,28 @@ func benchmarkFragmentBlockChecksum(b *testing.B, fillPercent float64) { } } } +*/ + +func BenchmarkFragment_Blocks(b *testing.B) { + if *FragmentPath == "" { + b.Skip("no fragment specified") + } + + // Open the fragment specified by the path. + f := pilosa.NewFragment(*FragmentPath, "d", "f", 0) + if err := f.Open(); err != nil { + b.Fatal(err) + } + defer f.Close() + + // Reset timer and execute benchmark. + b.ResetTimer() + for i := 0; i < b.N; i++ { + if a := f.Blocks(); len(a) == 0 { + b.Fatal("no blocks in fragment") + } + } +} // Fragment is a test wrapper for pilosa.Fragment. type Fragment struct { diff --git a/handler_test.go b/handler_test.go index f7723f1d3..89f7148e2 100644 --- a/handler_test.go +++ b/handler_test.go @@ -172,7 +172,7 @@ func TestHandler_Query_Bitmap_JSON(t *testing.T) { h.ServeHTTP(w, MustNewHTTPRequest("POST", "/query?db=d", strings.NewReader("Bitmap(100)"))) if w.Code != http.StatusOK { t.Fatalf("unexpected status code: %d", w.Code) - } else if body := w.Body.String(); body != `{"results":[{"attrs":{"a":"b","c":1,"d":true},"bits":[1,3,66,2097153]}]}`+"\n" { + } else if body := w.Body.String(); body != `{"results":[{"attrs":{"a":"b","c":1,"d":true},"bits":[1,3,66,1048577]}]}`+"\n" { t.Fatalf("unexpected body: %s", body) } } @@ -204,7 +204,7 @@ func TestHandler_Query_Bitmap_Profiles_JSON(t *testing.T) { h.ServeHTTP(w, MustNewHTTPRequest("POST", "/query?db=d&profiles=true", strings.NewReader("Bitmap(100)"))) if w.Code != http.StatusOK { t.Fatalf("unexpected status code: %d", w.Code) - } else if body := w.Body.String(); body != `{"results":[{"attrs":{"a":"b","c":1,"d":true},"bits":[1,3,66,2097153]}],"profiles":[{"id":3,"attrs":{"x":"y"}},{"id":66,"attrs":{"y":123,"z":false}}]}`+"\n" { + } else if body := w.Body.String(); body != `{"results":[{"attrs":{"a":"b","c":1,"d":true},"bits":[1,3,66,1048577]}],"profiles":[{"id":3,"attrs":{"x":"y"}},{"id":66,"attrs":{"y":123,"z":false}}]}`+"\n" { t.Fatalf("unexpected body: %s", body) } }