diff --git a/deps.json b/deps.json index b2f8e23e0..55f138ef9 100644 --- a/deps.json +++ b/deps.json @@ -9,6 +9,11 @@ "version": "35c98e8c885757fab1824daa77cc0e24bf30e76b", "type": "git" }, + "github.com/cactus/go-statsd-client/statsd": { + "repo": "github.com/cactus/go-statsd-client", + "version": "912f30c35e9cdf51f50bae24f071e227cca152fb", + "type": "git" + }, "gkvlite": { "repo": "github.com/steveyen/gkvlite", "version": "master", @@ -24,6 +29,11 @@ "version": "44eda643c1ae69e866e491b1c935c5b22e42350e", "type": "git" }, + "goleveldb": { + "repo": "github.com/syndtr/goleveldb/leveldb", + "version": "ff3719c6816e2cd194f05058452d660608e178ac", + "type": "git" + }, "google-crypto": { "repo": "code.google.com/p/go.crypto/ssh", "version": "7aa593ce8cea", @@ -83,11 +93,5 @@ "repo": "github.com/gorilla/websocket", "version": "92334662baa9cbebc2e6e68b8d56bc1233f85a4c", "type": "git" - }, - "github.com/cactus/go-statsd-client/statsd": { - "repo": "github.com/cactus/go-statsd-client", - "version": "912f30c35e9cdf51f50bae24f071e227cca152fb", - "type": "git" } - } diff --git a/index/brand.go b/index/brand.go index 7f644fe96..6403db638 100644 --- a/index/brand.go +++ b/index/brand.go @@ -85,6 +85,7 @@ func (self *Brand) cache_it(bm IBitmap, bitmap_id uint64, category uint64) { if bm.Count() >= self.threshold_value { self.bitmap_cache[bitmap_id] = &Rank{&Pair{bitmap_id, bm.Count()}, bm, category} if len(self.bitmap_cache) > self.threshold_length { + self.Rank() self.trim() } } @@ -309,6 +310,7 @@ func (self *Brand) Persist() error { return err } defer w.Close() + defer self.storage.Close() asize := len(self.bitmap_cache) var list RankList diff --git a/index/fragment_container.go b/index/fragment_container.go index 66707602f..887538dc7 100644 --- a/index/fragment_container.go +++ b/index/fragment_container.go @@ -233,12 +233,16 @@ type Fragment struct { exit chan *sync.WaitGroup } -func getStorage(db string, slice int, frame string) Storage { +func getStorage(db string, slice int, frame string, fid SUUID) Storage { storage_method := config.GetString("storage_backend") switch storage_method { default: return NewMemoryStorage() + case "leveldb": + base_path := config.GetString("level_db_path") + full_dir := fmt.Sprintf("%s/%s/%d/%s/%s", base_path, db, slice, frame, SUUID_to_Hex(fid)) + return NewLevelDBStorage(full_dir) case "cassandra": host := config.GetString("cassandra_host") if host == "" { @@ -258,10 +262,10 @@ func NewFragment(frag_id SUUID, db string, slice int, frame string) *Fragment { log.Println(fmt.Sprintf("XXXXXXXXXXXXXXXXXXXXXXXXXXX(%s)", frame)) if strings.HasSuffix(frame, ".n") { log.Println(frame + "TOP") - impl = NewBrand(db, frame, slice, getStorage(db, slice, frame), 50000, 45000, 100) + impl = NewBrand(db, frame, slice, getStorage(db, slice, frame, frag_id), 50000, 45000, 100) } else { log.Println(frame) - impl = NewGeneral(db, frame, slice, getStorage(db, slice, frame)) + impl = NewGeneral(db, frame, slice, getStorage(db, slice, frame, frag_id)) } f := new(Fragment) diff --git a/index/general.go b/index/general.go index c2c267ad5..95b899bf6 100644 --- a/index/general.go +++ b/index/general.go @@ -99,6 +99,7 @@ func (self *General) Persist() error { return err } defer w.Close() + defer self.storage.Close() encoder := json.NewEncoder(w) return encoder.Encode(self.keys) diff --git a/index/storage.go b/index/storage.go index 6d44a0501..ac774066e 100644 --- a/index/storage.go +++ b/index/storage.go @@ -6,4 +6,5 @@ type Storage interface { StoreBlock(id int64, db string, frame string, slice int, filter uint64, chunk int64, block_index int32, block int64) error BeginBatch() EndBatch() + Close() } diff --git a/index/storage_cass.go b/index/storage_cass.go index cca90a2b4..93f5f737f 100644 --- a/index/storage_cass.go +++ b/index/storage_cass.go @@ -53,6 +53,8 @@ func NewCassStorage(host, keyspace string) Storage { return obj } +func (c *CassandraStorage) Close() { +} func (c *CassandraStorage) Fetch(bitmap_id uint64, db string, frame string, slice int) (IBitmap, uint64) { var dumb = COUNTERMASK last_key := int64(dumb) diff --git a/index/storage_mem.go b/index/storage_mem.go index 95afc42b1..38043ba78 100644 --- a/index/storage_mem.go +++ b/index/storage_mem.go @@ -18,6 +18,8 @@ func NewMemoryStorage() Storage { func (c *MemoryStorage) BeginBatch() { } +func (c *MemoryStorage) Close() { +} func (c *MemoryStorage) EndBatch() { } func (c *MemoryStorage) Fetch(bitmap_id uint64, db string, frame string, slice int) (IBitmap, uint64) { diff --git a/index/storage_test.go b/index/storage_test.go index 93535fd84..9a1d00e11 100644 --- a/index/storage_test.go +++ b/index/storage_test.go @@ -2,12 +2,11 @@ package index import ( "fmt" - "net" "testing" - "time" // "io/ioutil" // "time" + "github.com/davecgh/go-spew/spew" . "github.com/smartystreets/goconvey/convey" ) @@ -18,42 +17,61 @@ func TestStorage(t *testing.T) { filter := 10 bitmap_id := uint64(1234) /* Convey("KV ", t, func() { - storage, _ := NewKVStorage("/tmp/", 0, db) - bm := storage.Fetch(bitmap_id, db, slice) - SetBit(bm, 0) - SetBit(bm, 1) - SetBit(bm, 2) - storage.Store(int64(bitmap_id), db, frame, slice, filter, 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) + storage, _ := NewKVStorage("/tmp/", 0, db) + bm := storage.Fetch(bitmap_id, db, slice) + SetBit(bm, 0) + SetBit(bm, 1) + SetBit(bm, 2) + storage.Store(int64(bitmap_id), db, frame, slice, filter, 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) - }) + }) + c, err := net.DialTimeout("tcp", "127.0.0.1:9042", 100*time.Millisecond) + if err != nil { + fmt.Println("NO cassandra. Skipping test.") + } else { + c.Close() + Convey("cassandra", t, func() { + fmt.Println("GO") + storage := NewCassStorage("127.0.0.1", "hotbox") + + fmt.Println("FETCH") + bm, _ := storage.Fetch(bitmap_id, db, frame, slice) + SetBit(bm, 0) + SetBit(bm, 1) + SetBit(bm, 2) + fmt.Println("STORE") + storage.Store(int64(bitmap_id), db, frame, slice, uint64(filter), bm.(*Bitmap)) + fmt.Println("FETCH") + bm2, _ := storage.Fetch(bitmap_id, db, frame, slice) + So(BitCount(bm), ShouldEqual, BitCount(bm2)) + So(BitCount(bm), ShouldEqual, bm.Count()) + So(BitCount(bm), ShouldEqual, 3) + + }) */ - c, err := net.DialTimeout("tcp", "127.0.0.1:9042", 100*time.Millisecond) - if err != nil { - fmt.Println("NO cassandra. Skipping test.") - } else { - c.Close() - Convey("cassandra", t, func() { - fmt.Println("GO") - storage := NewCassStorage("127.0.0.1", "hotbox") + Convey("leveldb", t, func() { + storage := NewLevelDBStorage("./basic/one") - fmt.Println("FETCH") - bm, _ := storage.Fetch(bitmap_id, db, frame, slice) - SetBit(bm, 0) - SetBit(bm, 1) - SetBit(bm, 2) - fmt.Println("STORE") - storage.Store(int64(bitmap_id), db, frame, slice, uint64(filter), bm.(*Bitmap)) - fmt.Println("FETCH") - bm2, _ := storage.Fetch(bitmap_id, db, frame, slice) - So(BitCount(bm), ShouldEqual, BitCount(bm2)) - So(BitCount(bm), ShouldEqual, bm.Count()) - So(BitCount(bm), ShouldEqual, 3) + fmt.Println("FETCH") + bm, _ := storage.Fetch(bitmap_id, db, frame, slice) + spew.Dump(bm) + SetBit(bm, 0) + SetBit(bm, 1) + SetBit(bm, 2) + fmt.Println("STORE") + storage.Store(int64(bitmap_id), db, frame, slice, uint64(filter), bm.(*Bitmap)) + //storage.FlushBatch() + fmt.Println("FETCH") + bm2, _ := storage.Fetch(bitmap_id, db, frame, slice) + So(BitCount(bm), ShouldEqual, BitCount(bm2)) + So(BitCount(bm), ShouldEqual, bm.Count()) + So(BitCount(bm), ShouldEqual, 3) + storage.Close() - }) - } + }) }