diff --git a/core/service.go b/core/service.go index a4f229810..b1b4a15fb 100644 --- a/core/service.go +++ b/core/service.go @@ -32,11 +32,11 @@ type Service struct { Index *index.FragmentContainer Hold *hold.Holder version string - name string + name string } func NewService() *Service { - log.Println(spew.Sdump("NewService") + log.Println(spew.Sdump("NewService")) service := new(Service) service.init_id() service.Etcd = etcd.NewClient(nil) @@ -48,7 +48,7 @@ func NewService() *Service { service.Index = index.NewFragmentContainer() service.Hold = hold.NewHolder() service.version = "0.0.8" - service.name ="Cruncher" + service.name = "Cruncher" service.PrepareLogging() return service } @@ -58,7 +58,7 @@ func (self *Service) PrepareLogging() { if base_path == "" { base_path = "/tmp" } - f, err := os.OpenFile(fmt.Sprintf("%s/%s.%s", base_path, self.name,self.Id), os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) + f, err := os.OpenFile(fmt.Sprintf("%s/%s.%s", base_path, self.name, self.Id), os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) if err != nil { log.Println("error opening file: %v", err) } @@ -115,6 +115,7 @@ func (service *Service) Run() { // TODO: reload configuration case <-sigterm: log.Println("SIGTERM! Cleaning up...") + service.Index.Shutdown() service.Stop() return } diff --git a/cruncher/cruncher.go b/cruncher/cruncher.go index d11dc6159..622931021 100644 --- a/cruncher/cruncher.go +++ b/cruncher/cruncher.go @@ -4,7 +4,6 @@ import ( "pilosa/core" "pilosa/dispatch" "pilosa/executor" - "pilosa/index" "pilosa/transport" "github.com/davecgh/go-spew/spew" @@ -13,18 +12,33 @@ import ( type Cruncher struct { *core.Service close_chan chan bool - api *index.FragmentContainer + // api *index.FragmentContainer } func (cruncher *Cruncher) Run() { spew.Dump("Cruncher.Run") + // go cruncher.Cleanup() cruncher.Service.Run() } +/* +func (cruncher *Cruncher) Cleanup() { + exit, done := cruncher.Service.GetExitChannels() + for { + select { + case <-exit: + spew.Dump(cruncher) + cruncher.api.Shutdown() + done <- 1 + } + } +} +*/ + func NewCruncher() *Cruncher { service := core.NewService() - fragment_container := index.NewFragmentContainer() - cruncher := Cruncher{service, make(chan bool), fragment_container} + //fragment_container := index.NewFragmentContainer() + cruncher := Cruncher{service, make(chan bool)} cruncher.Transport = transport.NewTcpTransport(service) cruncher.Dispatch = dispatch.NewDispatch(service) cruncher.Executor = executor.NewExecutor(service) diff --git a/db/topology.go b/db/topology.go index 1e7f7fdf6..267ef4124 100644 --- a/db/topology.go +++ b/db/topology.go @@ -312,6 +312,7 @@ func (d *Database) GetFragmentForBitmap(slice *Slice, bitmap *Bitmap) (*Fragment //d.mutex.Lock() //defer d.mutex.Unlock() frame, _ := d.getFrame(bitmap.FrameType) + log.Println(frame, slice) fsi, err := d.GetFrameSliceIntersect(frame, slice) if err != nil { log.Println(err) diff --git a/index/brand.go b/index/brand.go index 6a7b9f74c..70ff7f5db 100644 --- a/index/brand.go +++ b/index/brand.go @@ -1,6 +1,13 @@ package index -import "sort" +import ( + "encoding/json" + "fmt" + "log" + "pilosa/config" + "pilosa/util" + "sort" +) type Pair struct { Key, Count uint64 @@ -238,3 +245,52 @@ func (self *Brand) TopN(src_bitmap IBitmap, n int) []Pair { } return packagePairs(results) } +func (self *Brand) getFileName() string { + base := config.GetString("fragment_base") + if base == "" { + base = "." + } + + return fmt.Sprintf("%s/Brand.%s.%d.json", base, self.db, self.slice) +} + +func (self *Brand) Persist() error { + log.Println("Brand Persist:", self.getFileName()) + + w, err := util.Create(self.getFileName()) + if err != nil { + log.Println("Error opening outfile %s", self.getFileName()) + log.Println(err) + return err + } + defer w.Close() + + encoder := json.NewEncoder(w) + + asize := len(self.bitmap_cache) + results := make([]uint64, asize) + i := 0 + for k := range self.bitmap_cache { // map[uint64]*Rank + results[i] = k + i += 1 + } + return encoder.Encode(results) +} + +func (self *Brand) Load() { + log.Println("Brand Load") + r, err := util.Open(self.getFileName()) + if err != nil { + log.Println("NO Brand Init File:", self.getFileName()) + return + } + dec := json.NewDecoder(r) + var keys []uint64 + if err := dec.Decode(&keys); err != nil { + return + //log.Println("Bad mojo") + } + for _, k := range keys { + self.Get(k) + } +} diff --git a/index/fragment_container.go b/index/fragment_container.go index 9ff78c03f..f7fc0777c 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -3,9 +3,11 @@ package index import ( "encoding/gob" "errors" + "fmt" "log" "pilosa/config" . "pilosa/util" + "sync" "time" "github.com/golang/groupcache/lru" @@ -26,6 +28,18 @@ func init() { gob.Register(vh) } +func (self *FragmentContainer) Shutdown() { + log.Println("Container Shutdown Started") + var wg sync.WaitGroup + wg.Add(len(self.fragments)) + + for _, v := range self.fragments { + v.exit <- &wg + } + wg.Wait() + log.Println("Container Shutdown Complete") +} + func (self *FragmentContainer) LoadBitmap(frag_id SUUID, bitmap_id uint64, compressed_bitmap string) { if fragment, found := self.GetFragment(frag_id); found { request := NewLoader(bitmap_id, compressed_bitmap) @@ -154,6 +168,7 @@ func (self *FragmentContainer) AddFragment(db string, frame string, slice int, i log.Println("ADD FRAGMENT", frame) f := NewFragment(id, db, slice, frame) self.fragments[id] = f + f.Load() go f.ServeFragment() } @@ -164,6 +179,8 @@ type Pilosa interface { Clear() bool Store(bitmap_id uint64, bm IBitmap) Stats() interface{} + Persist() error + Load() } type Fragment struct { @@ -175,6 +192,7 @@ type Fragment struct { cache *lru.Cache mesg_count uint64 mesg_time time.Duration + exit chan *sync.WaitGroup } func getStorage(db string, slice int, frame string) Storage { @@ -206,14 +224,14 @@ func getStorage(db string, slice int, frame string) Storage { func NewFragment(frag_id SUUID, db string, slice int, frame string) *Fragment { var impl Pilosa - log.Println("XXXXXXXXXXXXXXXXXXXXXXXXXXXX", frame) + log.Println(fmt.Sprintf("XXXXXXXXXXXXXXXXXXXXXXXXXXX(%s)", frame)) switch frame { + case "brand": + log.Println("Brand") + impl = NewBrand(db, slice, getStorage(db, slice, frame), 50000, 45000, 100) default: log.Println("General") impl = NewGeneral(db, slice, getStorage(db, slice, frame)) - case "Brand": - log.Println("Brand") - impl = NewBrand(db, slice, getStorage(db, slice, frame), 50000, 45000, 100) } f := new(Fragment) @@ -222,6 +240,7 @@ func NewFragment(frag_id SUUID, db string, slice int, frame string) *Fragment { f.cache = lru.New(10000) f.impl = impl //NewGeneral(db, slice, NewMemoryStorage()) f.slice = slice + f.exit = make(chan *sync.WaitGroup) return f } @@ -283,24 +302,33 @@ func (self *Fragment) intersect(bitmaps []BitmapHandle) BitmapHandle { } return self.AllocHandle(result) } +func (self *Fragment) Persist() { + err := self.impl.Persist() + if err != nil { + log.Println("Error saving:", err) + } +} +func (self *Fragment) Load() { + self.impl.Load() +} func (self *Fragment) ServeFragment() { for { - req := <-self.requestChan - self.mesg_count++ - start := time.Now() - answer := req.Execute(self) - delta := time.Since(start) - self.mesg_count += 1 - self.mesg_time += delta - /* - var buffer bytes.Buffer - buffer.WriteString(`{ "results":`) - buffer.WriteString(answer) - buffer.WriteString(fmt.Sprintf(`,"query type": "%s"`, responder.QueryType())) - buffer.WriteString(fmt.Sprintf(`, "elapsed": "%s"}`, delta)) - */ - req.ResponseChannel() <- Result{answer, delta} + select { + case req := <-self.requestChan: + self.mesg_count++ + start := time.Now() + answer := req.Execute(self) + delta := time.Since(start) + self.mesg_count += 1 + self.mesg_time += delta + req.ResponseChannel() <- Result{answer, delta} + + case wg := <-self.exit: + log.Println("Fragment Shutdown") + self.Persist() + wg.Done() + } } } diff --git a/index/general.go b/index/general.go index 07dcd1d04..df932ac41 100644 --- a/index/general.go +++ b/index/general.go @@ -1,9 +1,18 @@ package index -import "github.com/golang/groupcache/lru" +import ( + "encoding/json" + "fmt" + "log" + "os" + "pilosa/config" + + "github.com/golang/groupcache/lru" +) type General struct { bitmap_cache *lru.Cache + keys map[uint64]interface{} db string slice int storage Storage @@ -15,12 +24,14 @@ func NewGeneral(db string, slice int, s Storage) *General { f.slice = slice f.db = db f.Clear() + f.keys = make(map[uint64]interface{}) //f.bitmap_cache = lru.New(10000) return f } func (self *General) Clear() bool { self.bitmap_cache = lru.New(10000) + self.bitmap_cache.OnEvicted = self.OnEvicted return true } @@ -31,6 +42,7 @@ func (self *General) Get(bitmap_id uint64) IBitmap { } bm = self.storage.Fetch(bitmap_id, self.db, self.slice) self.bitmap_cache.Add(bitmap_id, bm) + self.keys[bitmap_id] = 0 return bm.(*Bitmap) } func (self *General) SetBit(bitmap_id uint64, bit_pos uint64) bool { @@ -53,6 +65,10 @@ func (self *General) Store(bitmap_id uint64, bm IBitmap) { //nbm = Union(oldbm, bm) self.storage.Store(int64(bitmap_id), self.db, self.slice, bm.(*Bitmap)) self.bitmap_cache.Add(bitmap_id, bm) + self.keys[bitmap_id] = 0 +} +func (self *General) OnEvicted(key lru.Key, value interface{}) { + delete(self.keys, key.(uint64)) } func (self *General) Stats() interface{} { @@ -61,3 +77,38 @@ func (self *General) Stats() interface{} { "total size of cache in items": self.bitmap_cache.Len()} return stats } + +func (self *General) getFileName() string { + base := config.GetString("fragment_base") + if base == "" { + base = "." + } + return fmt.Sprintf("%s/General.%s.%d", base, self.db, self.slice) +} + +func (self *General) Persist() error { + log.Println("General Persist") + w, _ := os.Create(self.getFileName()) + + encoder := json.NewEncoder(w) + return encoder.Encode(self.keys) +} + +func (self *General) Load() { + log.Println("General Load") + r, err := os.Open(self.getFileName()) + if err != nil { + log.Println("NO General Init File:", self.getFileName()) + return + } + + dec := json.NewDecoder(r) + var keys map[uint64]interface{} + if err := dec.Decode(&keys); err != nil { + return + //log.Println("Bad mojo") + } + for k := range keys { + self.Get(k) + } +} diff --git a/index/storage_test.go b/index/storage_test.go index 7c27451ee..eb3e00631 100644 --- a/index/storage_test.go +++ b/index/storage_test.go @@ -25,21 +25,19 @@ func TestStorage(t *testing.T) { So(BitCount(bm), ShouldEqual, 3) }) - /* - Convey("cassandra", t, func() { - storage, _ := NewCassStorage("127.0.0.1", "hotbox") + Convey("cassandra", t, func() { + storage := NewCassStorage("127.0.0.1", "hotbox") - bm := storage.Fetch(bitmap_id, db, slice) - SetBit(bm, 0) - SetBit(bm, 1) - SetBit(bm, 2) - storage.Store(int64(bitmap_id), db, slice, bm.(*Bitmap)) - bm2 := storage.Fetch(bitmap_id, db, slice) - So(BitCount(bm), ShouldEqual, BitCount(bm2)) - So(BitCount(bm), ShouldEqual, bm.Count()) - So(BitCount(bm), ShouldEqual, 3) + bm := storage.Fetch(bitmap_id, db, slice) + SetBit(bm, 0) + SetBit(bm, 1) + SetBit(bm, 2) + storage.Store(int64(bitmap_id), db, slice, bm.(*Bitmap)) + bm2 := storage.Fetch(bitmap_id, db, slice) + So(BitCount(bm), ShouldEqual, BitCount(bm2)) + So(BitCount(bm), ShouldEqual, bm.Count()) + So(BitCount(bm), ShouldEqual, 3) - }) - */ + }) }