added storage interface and factored out cass dependancy

This commit is contained in:
Todd Gruben 2013-10-28 12:48:16 -05:00
parent ba4d107c3c
commit 77da1e0f15
6 changed files with 155 additions and 117 deletions

View file

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

View file

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

View file

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

7
index/storage.go Normal file
View file

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

91
index/storage_cass.go Normal file
View file

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

47
index/storage_mem.go Normal file
View file

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