Merge pull request #46 from benbjohnson/cache-persistence

Add cache persistence
This commit is contained in:
tgruben 2016-02-15 09:29:57 -06:00
commit a16db08b32
7 changed files with 354 additions and 62 deletions

12
Godeps/Godeps.json generated
View file

@ -10,13 +10,19 @@
"Comment": "v0.1.0-21-g056c9bc",
"Rev": "056c9bc7be7190eaa7715723883caffa5f8fa3e4"
},
{
"ImportPath": "github.com/boltdb/bolt",
"Comment": "v1.1.0-65-gee4a088",
"Rev": "ee4a0888a9abe7eefe5a0992ca4cb06864839873"
},
{
"ImportPath": "github.com/davecgh/go-spew/spew",
"Rev": "e762b3d1320b76030bd7f6cc2bfc3d9acce874c0"
},
{
"ImportPath": "github.com/gogo/protobuf/proto",
"Rev": "499788908625f4d83de42a204d1350fde8588e4f"
"Comment": "v0.1-125-g82d16f7",
"Rev": "82d16f734d6d871204a3feb1a73cb220cc92574c"
},
{
"ImportPath": "github.com/golang/groupcache/lru",
@ -25,6 +31,10 @@
{
"ImportPath": "github.com/yasushi-saito/rbtree",
"Rev": "571e2538414bf914c7e2909b61217b4e3e5508f4"
},
{
"ImportPath": "golang.org/x/sys/unix",
"Rev": "50c6bc5e4292a1d4e65c6e9be5f53be28bcbe28e"
}
]
}

View file

@ -16,6 +16,9 @@ type Cache interface {
Get(bitmapID uint64) *Bitmap
Len() int
// Returns a list of all bitmap IDs.
BitmapIDs() []uint64
// Updates the cache, if necessary.
Invalidate()
@ -39,7 +42,7 @@ func NewLRUCache(maxEntries int) *LRUCache {
return c
}
// Get returns a bitmap with a given id.
// Add adds a bitmap to the cache.
func (c *LRUCache) Add(bitmapID uint64, bm *Bitmap) {
c.cache.Add(bitmapID, bm)
c.bitmaps[bitmapID] = bm
@ -60,6 +63,16 @@ func (c *LRUCache) Len() int { return c.cache.Len() }
// Invalidate is a no-op.
func (c *LRUCache) Invalidate() {}
// BitmapIDs returns a list of all bitmap IDs in the cache.
func (c *LRUCache) BitmapIDs() []uint64 {
a := make([]uint64, 0, len(c.bitmaps))
for id := range c.bitmaps {
a = append(a, id)
}
sort.Sort(uint64Slice(a))
return a
}
// Top returns all bitmaps in the cache.
func (c *LRUCache) Top() []BitmapPair {
a := make([]BitmapPair, 0, len(c.bitmaps))
@ -98,7 +111,7 @@ func NewRankCache() *RankCache {
}
}
// Get returns a bitmap with a given id.
// Add adds a bitmap to the cache.
func (c *RankCache) Add(bitmapID uint64, bm *Bitmap) {
// Ignore if the bit count on the bitmap is below the threshold.
if bm.Count() < c.ThresholdValue {
@ -125,6 +138,16 @@ func (c *RankCache) Get(bitmapID uint64) *Bitmap { return c.entries[bitmapID] }
// Len returns the number of items in the cache.
func (c *RankCache) Len() int { return len(c.entries) }
// BitmapIDs returns a list of all bitmap IDs in the cache.
func (c *RankCache) BitmapIDs() []uint64 {
a := make([]uint64, 0, len(c.entries))
for id := range c.entries {
a = append(a, id)
}
sort.Sort(uint64Slice(a))
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.
@ -247,3 +270,10 @@ func decodePairs(a []*internal.Pair) []Pair {
}
return other
}
// uint64Slice represents a sortable slice of uint64 numbers.
type uint64Slice []uint64
func (p uint64Slice) Swap(i, j int) { p[i], p[j] = p[j], p[i] }
func (p uint64Slice) Len() int { return len(p) }
func (p uint64Slice) Less(i, j int) bool { return p[i] < p[j] }

View file

@ -9,6 +9,7 @@ import (
"net"
"net/http"
"os"
"os/signal"
"os/user"
"path/filepath"
"runtime/pprof"
@ -55,8 +56,18 @@ func main() {
os.Exit(1)
}
// Wait indefinitely.
<-(chan struct{})(nil)
// First SIGKILL causes server to shut down gracefully.
// Second signal causes a hard shutdown.
c := make(chan os.Signal, 2)
signal.Notify(c, os.Interrupt)
sig := <-c
fmt.Fprintf(m.Stderr, "Received %s; gracefully shutting down...\n", sig.String())
go func() { <-c; os.Exit(1) }()
if err := m.Close(); err != nil {
fmt.Fprintln(m.Stderr, err)
os.Exit(1)
}
}
// Main represents the main program execution.

View file

@ -3,6 +3,9 @@ package pilosa
import (
"errors"
"fmt"
"io"
"io/ioutil"
"log"
"os"
"sort"
"strings"
@ -11,16 +14,28 @@ import (
"time"
"unsafe"
"github.com/gogo/protobuf/proto"
"github.com/umbel/pilosa/internal"
"github.com/umbel/pilosa/roaring"
)
// SliceWidth is the number of profile IDs in a slice.
const SliceWidth = 65536
const (
// SliceWidth is the number of profile IDs in a slice.
SliceWidth = 65536
// SnapshotExt is the file extension used for an in-process snapshot.
const SnapshotExt = ".snapshotting"
// SnapshotExt is the file extension used for an in-process snapshot.
SnapshotExt = ".snapshotting"
const MinThreshold = 10
// CacheExt is the file extension for persisted cache ids.
CacheExt = ".cache"
MinThreshold = 10
)
const (
// DefaultCacheFlushInterval is the default value for Fragment.CacheFlushInterval.
DefaultCacheFlushInterval = 1 * time.Minute
)
// Fragment represents the intersection of a frame and slice in a database.
type Fragment struct {
@ -40,6 +55,16 @@ type Fragment struct {
// Bitmap cache.
cache Cache
// Close management
wg sync.WaitGroup
closing chan struct{}
// The interval at which the cached bitmap ids are persisted to disk.
CacheFlushInterval time.Duration
// Writer used for out-of-band log entries.
LogOutput io.Writer
// Bitmap attribute storage.
// Typically this is the parent frame unless overridden for testing.
BitmapAttrStore interface {
@ -49,29 +74,24 @@ type Fragment struct {
// NewFragment returns a new instance of Fragment.
func NewFragment(path, db, frame string, slice uint64) *Fragment {
f := &Fragment{
path: path,
db: db,
frame: frame,
slice: slice,
}
return &Fragment{
path: path,
db: db,
frame: frame,
slice: slice,
closing: make(chan struct{}, 0),
// Determine cache type from frame name.
if strings.HasSuffix(frame, ".n") {
c := NewRankCache()
c.ThresholdLength = 50000
c.ThresholdIndex = 45000
f.cache = c
} else {
f.cache = NewLRUCache(50000)
LogOutput: os.Stderr,
CacheFlushInterval: DefaultCacheFlushInterval,
}
return f
}
// Path returns the path the fragment was initialized with.
func (f *Fragment) Path() string { return f.path }
// CachePath returns the path to the fragment's cache data.
func (f *Fragment) CachePath() string { return f.path + CacheExt }
// DB returns the database the fragment was initialized with.
func (f *Fragment) DB() string { return f.db }
@ -81,13 +101,28 @@ func (f *Fragment) Frame() string { return f.frame }
// Slice returns the slice the fragment was initialized with.
func (f *Fragment) Slice() uint64 { return f.slice }
// Cache returns the fragment's cache.
// This is not safe for concurrent use.
func (f *Fragment) Cache() Cache { return f.cache }
// Open opens the underlying storage.
func (f *Fragment) Open() error {
f.mu.Lock()
defer f.mu.Unlock()
// Initialize storage in a function so we can close if anything goes wrong.
if err := f.openStorage(); err != nil {
if err := func() error {
// Initialize storage in a function so we can close if anything goes wrong.
if err := f.openStorage(); err != nil {
return err
}
// Fill cache with bitmaps persisted to disk.
if err := f.openCache(); err != nil {
return err
}
return nil
}(); err != nil {
f.close()
return err
}
@ -152,6 +187,47 @@ func (f *Fragment) openStorage() error {
}
// openCache initializes the cache from bitmap ids persisted to disk.
func (f *Fragment) openCache() error {
// Determine cache type from frame name.
if strings.HasSuffix(f.frame, ".n") {
c := NewRankCache()
c.ThresholdLength = 50000
c.ThresholdIndex = 45000
f.cache = c
} else {
f.cache = NewLRUCache(50000)
}
// Read cache data from disk.
path := f.CachePath()
buf, err := ioutil.ReadFile(path)
if os.IsNotExist(err) {
return nil
} else if err != nil {
return fmt.Errorf("open cache: %s", err)
}
// Unmarshal cache data.
var pb internal.Cache
if err := proto.Unmarshal(buf, &pb); err != nil {
log.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
return nil
}
// Read in all bitmaps by ID.
// This will cause them to be added to the cache.
for _, bitmapID := range pb.GetBitmapIDs() {
f.bitmap(bitmapID)
}
// Periodically flush cache.
f.wg.Add(1)
go func() { defer f.wg.Done(); f.monitorCacheFlush() }()
return nil
}
// Close flushes the underlying storage, closes the file and unlocks it.
func (f *Fragment) Close() error {
f.mu.Lock()
@ -160,9 +236,22 @@ func (f *Fragment) Close() error {
}
func (f *Fragment) close() error {
if err := f.closeStorage(); err != nil {
return err
// Notify goroutines of closing and wait for completion.
close(f.closing)
f.mu.Unlock()
f.wg.Wait()
f.mu.Lock()
// Flush cache if closing gracefully.
if err := f.flushCache(); err != nil {
f.logger().Printf("error flushing cache on close: err=%s, path=%s", err, f.path)
}
// Close underlying storage.
if err := f.closeStorage(); err != nil {
f.logger().Printf("error closing storage: err=%s, path=%s", err, f.path)
}
return nil
}
@ -194,6 +283,9 @@ func (f *Fragment) closeStorage() error {
return nil
}
// logger returns a logger instance for the fragment.nt.
func (f *Fragment) logger() *log.Logger { return log.New(f.LogOutput, "", log.LstdFlags) }
// Bitmap returns a bitmap by ID.
func (f *Fragment) Bitmap(bitmapID uint64) *Bitmap {
f.mu.Lock()
@ -479,6 +571,55 @@ func (f *Fragment) snapshot() error {
return nil
}
// monitorCacheFlush periodically flushes the cache to disk.
// This is run in a goroutine.
func (f *Fragment) monitorCacheFlush() {
ticker := time.NewTicker(f.CacheFlushInterval)
defer ticker.Stop()
for {
select {
case <-f.closing:
return
case <-ticker.C:
if err := f.FlushCache(); err != nil {
f.logger().Printf("error flushing cache: err=%s, path=%s", err, f.CachePath())
}
}
}
}
// FlushCache writes the cache data to disk.
func (f *Fragment) FlushCache() error {
f.mu.Lock()
defer f.mu.Unlock()
return f.flushCache()
}
func (f *Fragment) flushCache() error {
if f.cache == nil {
return nil
}
// Retrieve a list of bitmap ids from the cache.
bitmapIDs := f.cache.BitmapIDs()
// Marshal cache data to bytes.
buf, err := proto.Marshal(&internal.Cache{
BitmapIDs: bitmapIDs,
})
if err != nil {
return err
}
// Write to disk.
if err := ioutil.WriteFile(f.CachePath(), buf, 0666); err != nil {
return err
}
return nil
}
func madvise(b []byte, advice int) (err error) {
_, _, e1 := syscall.Syscall(syscall.SYS_MADVISE, uintptr(unsafe.Pointer(&b[0])), uintptr(len(b)), uintptr(advice))
if e1 != 0 {

View file

@ -211,6 +211,70 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) {
}
}
// Ensure a fragment's cache can be persisted between restarts.
func TestFragment_LRUCache_Persistence(t *testing.T) {
f := MustOpenFragment("d", "f", 0)
defer f.Close()
// Set bits on the fragment.
for i := uint64(0); i < 1000; i++ {
if err := f.SetBit(i, 0); err != nil {
t.Fatal(err)
}
}
// Verify correct cache type and size.
if cache, ok := f.Cache().(*pilosa.LRUCache); !ok {
t.Fatalf("unexpected cache: %T", f.Cache())
} else if cache.Len() != 1000 {
t.Fatalf("unexpected cache len: %d", cache.Len())
}
// Reopen the fragment.
if err := f.Reopen(); err != nil {
t.Fatal(err)
}
// Re-verify correct cache type and size.
if cache, ok := f.Cache().(*pilosa.LRUCache); !ok {
t.Fatalf("unexpected cache: %T", f.Cache())
} else if cache.Len() != 1000 {
t.Fatalf("unexpected cache len: %d", cache.Len())
}
}
// Ensure a fragment's cache can be persisted between restarts.
func TestFragment_RankCache_Persistence(t *testing.T) {
f := MustOpenFragment("d", "f.n", 0)
defer f.Close()
// Set bits on the fragment.
for i := uint64(0); i < 1000; i++ {
if err := f.SetBit(i, 0); err != nil {
t.Fatal(err)
}
}
// Verify correct cache type and size.
if cache, ok := f.Cache().(*pilosa.RankCache); !ok {
t.Fatalf("unexpected cache: %T", f.Cache())
} else if cache.Len() != 1000 {
t.Fatalf("unexpected cache len: %d", cache.Len())
}
// Reopen the fragment.
if err := f.Reopen(); err != nil {
t.Fatal(err)
}
// Re-verify correct cache type and size.
if cache, ok := f.Cache().(*pilosa.RankCache); !ok {
t.Fatalf("unexpected cache: %T", f.Cache())
} else if cache.Len() != 1000 {
t.Fatalf("unexpected cache len: %d", cache.Len())
}
}
// Fragment is a test wrapper for pilosa.Fragment.
type Fragment struct {
*pilosa.Fragment
@ -245,17 +309,19 @@ func MustOpenFragment(db, frame string, slice uint64) *Fragment {
// Close closes the fragment and removes all underlying data.
func (f *Fragment) Close() error {
defer os.Remove(f.Path())
defer os.Remove(f.CachePath())
return f.Fragment.Close()
}
// Reopen closes the fragment and reopens it as a new instance.
func (f *Fragment) Reopen() error {
path := f.Path()
if err := f.Close(); err != nil {
if err := f.Fragment.Close(); err != nil {
return err
}
f = &Fragment{Fragment: pilosa.NewFragment(path, f.DB(), f.Frame(), f.Slice())}
f.Fragment = pilosa.NewFragment(path, f.DB(), f.Frame(), f.Slice())
f.Fragment.BitmapAttrStore = f.BitmapAttrStore
if err := f.Open(); err != nil {
return err
}

View file

@ -19,19 +19,22 @@ It has these top-level messages:
QueryResponse
ImportRequest
ImportResponse
Cache
*/
package internal
import proto "github.com/gogo/protobuf/proto"
import fmt "fmt"
import math "math"
// Reference imports to suppress errors if they are not otherwise used.
var _ = proto.Marshal
var _ = fmt.Errorf
var _ = math.Inf
type Bitmap struct {
Chunks []*Chunk `protobuf:"bytes,1,rep" json:"Chunks,omitempty"`
Attrs []*Attr `protobuf:"bytes,2,rep" json:"Attrs,omitempty"`
Chunks []*Chunk `protobuf:"bytes,1,rep,name=Chunks" json:"Chunks,omitempty"`
Attrs []*Attr `protobuf:"bytes,2,rep,name=Attrs" json:"Attrs,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -54,8 +57,8 @@ func (m *Bitmap) GetAttrs() []*Attr {
}
type Chunk struct {
Key *uint64 `protobuf:"varint,1,req" json:"Key,omitempty"`
Value []uint64 `protobuf:"varint,2,rep" json:"Value,omitempty"`
Key *uint64 `protobuf:"varint,1,req,name=Key" json:"Key,omitempty"`
Value []uint64 `protobuf:"varint,2,rep,name=Value" json:"Value,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -78,8 +81,8 @@ func (m *Chunk) GetValue() []uint64 {
}
type Pair struct {
Key *uint64 `protobuf:"varint,1,req" json:"Key,omitempty"`
Count *uint64 `protobuf:"varint,2,req" json:"Count,omitempty"`
Key *uint64 `protobuf:"varint,1,req,name=Key" json:"Key,omitempty"`
Count *uint64 `protobuf:"varint,2,req,name=Count" json:"Count,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -102,8 +105,8 @@ func (m *Pair) GetCount() uint64 {
}
type Bit struct {
BitmapID *uint64 `protobuf:"varint,1,req" json:"BitmapID,omitempty"`
ProfileID *uint64 `protobuf:"varint,2,req" json:"ProfileID,omitempty"`
BitmapID *uint64 `protobuf:"varint,1,req,name=BitmapID" json:"BitmapID,omitempty"`
ProfileID *uint64 `protobuf:"varint,2,req,name=ProfileID" json:"ProfileID,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -126,8 +129,8 @@ func (m *Bit) GetProfileID() uint64 {
}
type Profile struct {
ID *uint64 `protobuf:"varint,1,req" json:"ID,omitempty"`
Attrs []*Attr `protobuf:"bytes,2,rep" json:"Attrs,omitempty"`
ID *uint64 `protobuf:"varint,1,req,name=ID" json:"ID,omitempty"`
Attrs []*Attr `protobuf:"bytes,2,rep,name=Attrs" json:"Attrs,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -150,10 +153,10 @@ func (m *Profile) GetAttrs() []*Attr {
}
type Attr struct {
Key *string `protobuf:"bytes,1,req" json:"Key,omitempty"`
StringValue *string `protobuf:"bytes,2,opt" json:"StringValue,omitempty"`
IntValue *int64 `protobuf:"varint,3,opt" json:"IntValue,omitempty"`
BoolValue *bool `protobuf:"varint,4,opt" json:"BoolValue,omitempty"`
Key *string `protobuf:"bytes,1,req,name=Key" json:"Key,omitempty"`
StringValue *string `protobuf:"bytes,2,opt,name=StringValue" json:"StringValue,omitempty"`
IntValue *int64 `protobuf:"varint,3,opt,name=IntValue" json:"IntValue,omitempty"`
BoolValue *bool `protobuf:"varint,4,opt,name=BoolValue" json:"BoolValue,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -190,10 +193,10 @@ func (m *Attr) GetBoolValue() bool {
}
type QueryRequest struct {
DB *string `protobuf:"bytes,1,req" json:"DB,omitempty"`
Query *string `protobuf:"bytes,2,req" json:"Query,omitempty"`
Slices []uint64 `protobuf:"varint,3,rep" json:"Slices,omitempty"`
Profiles *bool `protobuf:"varint,4,opt" json:"Profiles,omitempty"`
DB *string `protobuf:"bytes,1,req,name=DB" json:"DB,omitempty"`
Query *string `protobuf:"bytes,2,req,name=Query" json:"Query,omitempty"`
Slices []uint64 `protobuf:"varint,3,rep,name=Slices" json:"Slices,omitempty"`
Profiles *bool `protobuf:"varint,4,opt,name=Profiles" json:"Profiles,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -230,11 +233,11 @@ func (m *QueryRequest) GetProfiles() bool {
}
type QueryResponse struct {
Err *string `protobuf:"bytes,1,opt" json:"Err,omitempty"`
Bitmap *Bitmap `protobuf:"bytes,2,opt" json:"Bitmap,omitempty"`
N *uint64 `protobuf:"varint,3,opt" json:"N,omitempty"`
Pairs []*Pair `protobuf:"bytes,4,rep" json:"Pairs,omitempty"`
Profiles []*Profile `protobuf:"bytes,5,rep" json:"Profiles,omitempty"`
Err *string `protobuf:"bytes,1,opt,name=Err" json:"Err,omitempty"`
Bitmap *Bitmap `protobuf:"bytes,2,opt,name=Bitmap" json:"Bitmap,omitempty"`
N *uint64 `protobuf:"varint,3,opt,name=N" json:"N,omitempty"`
Pairs []*Pair `protobuf:"bytes,4,rep,name=Pairs" json:"Pairs,omitempty"`
Profiles []*Profile `protobuf:"bytes,5,rep,name=Profiles" json:"Profiles,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -278,11 +281,11 @@ func (m *QueryResponse) GetProfiles() []*Profile {
}
type ImportRequest struct {
DB *string `protobuf:"bytes,1,req" json:"DB,omitempty"`
Frame *string `protobuf:"bytes,2,req" json:"Frame,omitempty"`
Slice *uint64 `protobuf:"varint,3,req" json:"Slice,omitempty"`
BitmapIDs []uint64 `protobuf:"varint,4,rep" json:"BitmapIDs,omitempty"`
ProfileIDs []uint64 `protobuf:"varint,5,rep" json:"ProfileIDs,omitempty"`
DB *string `protobuf:"bytes,1,req,name=DB" json:"DB,omitempty"`
Frame *string `protobuf:"bytes,2,req,name=Frame" json:"Frame,omitempty"`
Slice *uint64 `protobuf:"varint,3,req,name=Slice" json:"Slice,omitempty"`
BitmapIDs []uint64 `protobuf:"varint,4,rep,name=BitmapIDs" json:"BitmapIDs,omitempty"`
ProfileIDs []uint64 `protobuf:"varint,5,rep,name=ProfileIDs" json:"ProfileIDs,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -326,7 +329,7 @@ func (m *ImportRequest) GetProfileIDs() []uint64 {
}
type ImportResponse struct {
Err *string `protobuf:"bytes,1,opt" json:"Err,omitempty"`
Err *string `protobuf:"bytes,1,opt,name=Err" json:"Err,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
@ -341,5 +344,32 @@ func (m *ImportResponse) GetErr() string {
return ""
}
func init() {
type Cache struct {
BitmapIDs []uint64 `protobuf:"varint,1,rep,name=BitmapIDs" json:"BitmapIDs,omitempty"`
XXX_unrecognized []byte `json:"-"`
}
func (m *Cache) Reset() { *m = Cache{} }
func (m *Cache) String() string { return proto.CompactTextString(m) }
func (*Cache) ProtoMessage() {}
func (m *Cache) GetBitmapIDs() []uint64 {
if m != nil {
return m.BitmapIDs
}
return nil
}
func init() {
proto.RegisterType((*Bitmap)(nil), "internal.Bitmap")
proto.RegisterType((*Chunk)(nil), "internal.Chunk")
proto.RegisterType((*Pair)(nil), "internal.Pair")
proto.RegisterType((*Bit)(nil), "internal.Bit")
proto.RegisterType((*Profile)(nil), "internal.Profile")
proto.RegisterType((*Attr)(nil), "internal.Attr")
proto.RegisterType((*QueryRequest)(nil), "internal.QueryRequest")
proto.RegisterType((*QueryResponse)(nil), "internal.QueryResponse")
proto.RegisterType((*ImportRequest)(nil), "internal.ImportRequest")
proto.RegisterType((*ImportResponse)(nil), "internal.ImportResponse")
proto.RegisterType((*Cache)(nil), "internal.Cache")
}

View file

@ -58,3 +58,7 @@ message ImportRequest {
message ImportResponse {
optional string Err = 1;
}
message Cache {
repeated uint64 BitmapIDs = 1;
}