added basic persitant for brands;general fragments

This commit is contained in:
Todd Gruben 2014-01-28 15:32:12 -06:00
parent c48ad2890e
commit df8dc3d927
7 changed files with 192 additions and 43 deletions

View file

@ -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
}

View file

@ -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)

View file

@ -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)

View file

@ -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)
}
}

View file

@ -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()
}
}
}

View file

@ -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)
}
}

View file

@ -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)
})
*/
})
}