From 77da1e0f1545f3c1696dc8ae9cfbc3c46f08eb00 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 28 Oct 2013 12:48:16 -0500 Subject: [PATCH] added storage interface and factored out cass dependancy --- db/shard.go | 9 ++-- db/shard_test.go | 12 ++--- index/bitmap.go | 106 ------------------------------------------ index/storage.go | 7 +++ index/storage_cass.go | 91 ++++++++++++++++++++++++++++++++++++ index/storage_mem.go | 47 +++++++++++++++++++ 6 files changed, 155 insertions(+), 117 deletions(-) create mode 100644 index/storage.go create mode 100644 index/storage_cass.go create mode 100644 index/storage_mem.go diff --git a/db/shard.go b/db/shard.go index 56c2b6cd0..b982cc0c9 100644 --- a/db/shard.go +++ b/db/shard.go @@ -4,7 +4,6 @@ import ( "fmt" "log" "lru" - "tux21b.org/v1/gocql" "pilosa/index" ) @@ -60,16 +59,16 @@ type Shard struct { done chan bool bitmap_cache *lru.Cache shard_key int32 - db *gocql.Session + storage index.Storage } -func NewShard(shard_key int32) *Shard { +func NewShard(shard_key int32, s index.Storage) *Shard { f := new(Shard) f.inbound = make(chan Request, 512) f.done = make(chan bool) f.bitmap_cache = lru.New(10000) f.shard_key = shard_key - f.db = index.GetDB() + f.storage = s return f } @@ -100,7 +99,7 @@ func (f *Shard) Get(bitmap_id uint64) index.IBitmap { return bm.(*index.Bitmap) } //need to stick it in the cache - bm = index.FetchCass(f.db, bitmap_id, f.shard_key) + bm = f.storage.Fetch(bitmap_id, f.shard_key) f.bitmap_cache.Add(bitmap_id, bm) //return umbel.FetchCass(f.db,bitmap_id,f.shard_key) return bm.(*index.Bitmap) diff --git a/db/shard_test.go b/db/shard_test.go index 292c17be9..4bbdcaa8a 100644 --- a/db/shard_test.go +++ b/db/shard_test.go @@ -2,19 +2,19 @@ package db import ( "testing" + "pilosa/index" + "log" ) - func TestFragment(t *testing.T) { - /* - f := NewFragment(1638400) +func TestShard(t *testing.T) { + f := NewShard(1638400,index.NewMemoryStorage()) go f.Run() log.Println("REQUEST") f. inbound <- Request{IdCount{7812}} f.inbound <- Request{UnionCount{7812 , 227149}} f.inbound <- Request{IntersectCount{7812, 227149}} - f.inbou nd <- Request{Quit{}} + f.inbound <- Request{Quit{}} <-f.done - */ - } +} diff --git a/index/bitmap.go b/index/bitmap.go index 72c84b497..71bfc267b 100644 --- a/index/bitmap.go +++ b/index/bitmap.go @@ -3,28 +3,10 @@ package index // #cgo CFLAGS:-mpopcnt import ( - // "C" - // "bufio" - // "encoding/binary" - // "database/sql" - // "flag" "log" - // "fmt" - // "encoding/base64" - // "compress/gzip" - // _ "github.com/tux21b/gocql" - // "reflect" - //"github.com/ugorji/go/codec" "bytes" "encoding/gob" "github.com/yasushi-saito/rbtree" - // "os" - // "runtime" - // "sort" - // "strconv" - // "strings" - // "time" - "tux21b.org/v1/gocql" ) const ( @@ -422,91 +404,3 @@ func BitCount(b IBitmap) uint64 { return total } -/* -func sendResults(ret_chan chan map[uint64]uint64, results CachedBitmapList, end int) { - ret_val := make(map[uint64]uint64) - for i, v := range results { - if i >= end { - break - } - ret_val[v.Key] = v.Count - } - ret_chan <- ret_val -} -*/ - -func FetchCass(db *gocql.Session, bitmap_id uint64, shard int32) IBitmap { - var dumb = COUNTERMASK - last_key := int64(dumb) - marker := int64(dumb) - var id = int64(bitmap_id) - - var ( - chunk *Chunk - chunk_key, block int64 - block_index uint32 - s8 uint8 - ) - log.Println("FETCHING ", bitmap_id, shard) - - bitmap := CreateRBBitmap() - iter := db.Query("SELECT Chunkkey,BlockIndex,block FROM bitmap WHERE bitmap_id=? AND shard_id=? ", id, shard).Iter() - count := int64(0) - for iter.Scan(&chunk_key, &block_index, &block) { - s8 = uint8(block_index) - if chunk_key != marker { - if chunk_key != last_key { - chunk = &Chunk{uint64(chunk_key), BlockArray{}} - bitmap.AddChunk(chunk) - } - chunk.Value.Block[s8] = uint64(block) - - } else { - count = block - } - last_key = chunk_key - - } - bitmap.SetCount(uint64(count)) - return bitmap -} -func GetDB() *gocql.Session { - cluster := gocql.NewCluster("127.0.0.1") - cluster.Keyspace = "hotbox" - cluster.Consistency = gocql.Quorum - //cluster.ProtoVersion = 1 - // cluster.CQLVersion = "3.0.0" - session := cluster.CreateSession() - if err := session.Query("USE hotbox").Exec(); err != nil { - } - return session -} - -func StoreCassandra(db *gocql.Session, id int64, shard_key int32, bitmap *Bitmap) error { - for i := bitmap.Min(); !i.Limit(); i = i.Next() { - var chunk = i.Item() - for idx, block := range chunk.Value.Block { - block_index := int32(idx) - iblock := int64(block) - if iblock != 0 { - StoreBlock(db, id, shard_key, int64(chunk.Key), block_index, iblock) - } - } - } - c := int64(BitCount(bitmap)) - - var dumb = COUNTERMASK - COUNTER_KEY := int64(dumb) - - StoreBlock(db, id, shard_key, COUNTER_KEY, 0, c) - return nil -} - -func StoreBlock(db *gocql.Session, id int64, shard_key int32, chunk int64, block_index int32, block int64) error { - if err := db.Query(`INSERT INTO bitmap (bitmap_id, shard_id, ChunkKey, BlockIndex,block) VALUES (?,?, ?,?,?);`, id, shard_key, chunk, block_index, block).Exec(); err != nil { - log.Println(err) - log.Println("INSERT ", id, chunk, block_index) - return err - } - return nil -} diff --git a/index/storage.go b/index/storage.go new file mode 100644 index 000000000..979539358 --- /dev/null +++ b/index/storage.go @@ -0,0 +1,7 @@ +package index + +type Storage interface { + Fetch(bitmap_id uint64, shard int32) IBitmap + Store(id int64, shard_key int32, bitmap *Bitmap) error + StoreBlock(id int64, shard_key int32, chunk int64, block_index int32, block int64) error +} diff --git a/index/storage_cass.go b/index/storage_cass.go new file mode 100644 index 000000000..5accf6026 --- /dev/null +++ b/index/storage_cass.go @@ -0,0 +1,91 @@ +package index + +// #cgo CFLAGS:-mpopcnt + +import ( + "log" + "tux21b.org/v1/gocql" +) + +type CassandraStorage struct{ + db *gocql.Session +} +func NewCassStorage() Storage{ + obj := new(CassandraStorage) + cluster := gocql.NewCluster("127.0.0.1") + cluster.Keyspace = "hotbox" + cluster.Consistency = gocql.Quorum + //cluster.ProtoVersion = 1 + // cluster.CQLVersion = "3.0.0" + session := cluster.CreateSession() + if err := session.Query("USE hotbox").Exec(); err != nil { + } + obj.db = session + return obj +} + +func (c *CassandraStorage)Fetch( bitmap_id uint64, shard int32) IBitmap { + var dumb = COUNTERMASK + last_key := int64(dumb) + marker := int64(dumb) + var id = int64(bitmap_id) + + var ( + chunk *Chunk + chunk_key, block int64 + block_index uint32 + s8 uint8 + ) + log.Println("FETCHING ", bitmap_id, shard) + + bitmap := CreateRBBitmap() + iter := c.db.Query("SELECT Chunkkey,BlockIndex,block FROM bitmap WHERE bitmap_id=? AND shard_id=? ", id, shard).Iter() + count := int64(0) + for iter.Scan(&chunk_key, &block_index, &block) { + s8 = uint8(block_index) + if chunk_key != marker { + if chunk_key != last_key { + chunk = &Chunk{uint64(chunk_key), BlockArray{}} + bitmap.AddChunk(chunk) + } + chunk.Value.Block[s8] = uint64(block) + + } else { + count = block + } + last_key = chunk_key + + } + bitmap.SetCount(uint64(count)) + return bitmap +} + +func (c *CassandraStorage) Store( id int64, shard_key int32, bitmap *Bitmap) error { + for i := bitmap.Min(); !i.Limit(); i = i.Next() { + var chunk = i.Item() + for idx, block := range chunk.Value.Block { + block_index := int32(idx) + iblock := int64(block) + if iblock != 0 { + c.StoreBlock(id, shard_key, int64(chunk.Key), block_index, iblock) + } + } + } + cnt := int64(BitCount(bitmap)) + + var dumb = COUNTERMASK + COUNTER_KEY := int64(dumb) + + c.StoreBlock(id, shard_key, COUNTER_KEY, 0, cnt) + return nil +} + +func (c *CassandraStorage)StoreBlock(id int64, shard_key int32, chunk int64, block_index int32, block int64) error { + + if err := c.db.Query(`INSERT INTO bitmap (bitmap_id, shard_id, ChunkKey, BlockIndex,block) VALUES (?,?, ?,?,?);`, id, shard_key, chunk, block_index, block).Exec(); err != nil { + log.Println(err) + log.Println("INSERT ", id, chunk, block_index) + return err + } + return nil +} diff --git a/index/storage_mem.go b/index/storage_mem.go new file mode 100644 index 000000000..ba000496c --- /dev/null +++ b/index/storage_mem.go @@ -0,0 +1,47 @@ +package index + +// #cgo CFLAGS:-mpopcnt + +import ( +// "log" + "fmt" +) + +type MemoryStorage struct{ + db map[string]*Bitmap + +} +func NewMemoryStorage() Storage{ + obj := new(MemoryStorage) + obj.db = make(map[string]*Bitmap) + + return obj +} + +func (c *MemoryStorage)Fetch( bitmap_id uint64, shard_key int32) IBitmap { + key := fmt.Sprintf("%d:%d",bitmap_id,shard_key) + bitmap,found := c.db[key] + if !found{ + bitmap = CreateRBBitmap().(*Bitmap) + c.db[key]=bitmap + } + return bitmap +} + +func (c *MemoryStorage) Store( bitmap_id int64, shard_key int32, bitmap *Bitmap) error { + key := fmt.Sprintf("%d:%d",bitmap_id,shard_key) + c.db[key]= bitmap + return nil +} + +func (c *MemoryStorage)StoreBlock(bitmap_id int64, shard_key int32, chunk_key int64, block_index int32, block int64) error { + bm := c.Fetch(uint64(bitmap_id),shard_key) + node := GetChunk(bm,uint64(chunk_key)) + if node == nil{ + node = &Chunk{uint64(chunk_key), BlockArray{}} + bm.AddChunk(node) + } + node.Value.Block[block_index]=uint64(block) + + return nil +}