mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
added bytes interface to fragment; added needed files for kv storage
This commit is contained in:
parent
637c758858
commit
3c85487f04
5 changed files with 218 additions and 0 deletions
|
|
@ -41,6 +41,7 @@ type Command interface {
|
|||
Execute(*Fragment) Calculation
|
||||
GetResponder() *Responder
|
||||
}
|
||||
|
||||
type CmdGet struct {
|
||||
meta *Responder
|
||||
bitmap_id uint64
|
||||
|
|
@ -125,3 +126,38 @@ func (cmd *CmdSetBit) Execute(f *Fragment) Calculation {
|
|||
func (cmd *CmdSetBit) GetResponder() *Responder {
|
||||
return cmd.meta
|
||||
}
|
||||
|
||||
type CmdGetBytes struct {
|
||||
meta *Responder
|
||||
bitmap BitmapHandle
|
||||
}
|
||||
|
||||
func NewGetBytes(bh BitmapHandle) *CmdGetBytes {
|
||||
return &CmdGetBytes{NewResponder("GetBytes"), bh}
|
||||
}
|
||||
|
||||
func (cmd *CmdGetBytes) GetResponder() *Responder {
|
||||
return cmd.meta
|
||||
}
|
||||
func (cmd *CmdGetBytes) Execute(f *Fragment) Calculation {
|
||||
bm, _ := f.getBitmap(cmd.bitmap)
|
||||
return bm.ToBytes()
|
||||
}
|
||||
|
||||
type CmdFromBytes struct {
|
||||
meta *Responder
|
||||
bytes []byte
|
||||
}
|
||||
|
||||
func NewFromBytes(bytes []byte) *CmdFromBytes {
|
||||
return &CmdFromBytes{NewResponder("FromBytes"), bytes}
|
||||
}
|
||||
|
||||
func (cmd *CmdFromBytes) GetResponder() *Responder {
|
||||
return cmd.meta
|
||||
}
|
||||
func (cmd *CmdFromBytes) Execute(f *Fragment) Calculation {
|
||||
result := NewBitmap()
|
||||
result.FromBytes(cmd.bytes)
|
||||
return f.AllocHandle(result)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -58,6 +58,24 @@ func (self *FragmentContainer) Count(frag_id SUUID, bitmap BitmapHandle) (uint64
|
|||
return 0, errors.New("Invalid Bitmap Handle")
|
||||
}
|
||||
|
||||
func (self *FragmentContainer) GetBytes(frag_id SUUID, bh BitmapHandle) ([]byte, error) {
|
||||
if fragment, found := self.GetFragment(frag_id); found {
|
||||
request := NewGetBytes(bh)
|
||||
fragment.requestChan <- request
|
||||
return request.GetResponder().Response().answer.([]byte), nil
|
||||
}
|
||||
return nil, errors.New("Invalid Bitmap Handle")
|
||||
}
|
||||
|
||||
func (self *FragmentContainer) FromBytes(frag_id SUUID, bytes []byte) (BitmapHandle, error) {
|
||||
if fragment, found := self.GetFragment(frag_id); found {
|
||||
request := NewFromBytes(bytes)
|
||||
fragment.requestChan <- request
|
||||
return request.GetResponder().Response().answer.(BitmapHandle), nil
|
||||
}
|
||||
return 0, errors.New("Invalid Bitmap Handle")
|
||||
}
|
||||
|
||||
func (self *FragmentContainer) SetBit(frag_id SUUID, bitmap_id uint64, pos uint64) (bool, error) {
|
||||
if fragment, found := self.GetFragment(frag_id); found {
|
||||
request := NewSetBit(bitmap_id, pos)
|
||||
|
|
|
|||
|
|
@ -52,4 +52,15 @@ func TestServer(t *testing.T) {
|
|||
num, _ = dummy.Count(id, result)
|
||||
So(num, ShouldEqual, 0)
|
||||
})
|
||||
Convey("Bytes", t, func() {
|
||||
bi1 := uint64(1234)
|
||||
bh1, _ := dummy.Get(id, bi1)
|
||||
before, _ := dummy.Count(id, bh1)
|
||||
|
||||
bytes, _ := dummy.GetBytes(id, bh1)
|
||||
bh2, _ := dummy.FromBytes(id, bytes)
|
||||
|
||||
after, _ := dummy.Count(id, bh2)
|
||||
So(before, ShouldEqual, after)
|
||||
})
|
||||
}
|
||||
|
|
|
|||
124
index/storage_kvlite.go
Normal file
124
index/storage_kvlite.go
Normal file
|
|
@ -0,0 +1,124 @@
|
|||
package index
|
||||
|
||||
// #cgo CFLAGS:-mpopcnt
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"os"
|
||||
|
||||
"github.com/steveyen/gkvlite"
|
||||
)
|
||||
|
||||
type KVStorage struct {
|
||||
db *gkvlite.Store
|
||||
cc *gkvlite.Collection
|
||||
}
|
||||
|
||||
func NewKVStorage(path string, slice int, db string) (Storage, error) {
|
||||
// log.Println("Hello")
|
||||
obj := new(KVStorage)
|
||||
|
||||
f, err := os.Create(fmt.Sprintf("%sKV.%s.%d", path, db, slice))
|
||||
s, err := gkvlite.NewStore(f)
|
||||
obj.cc = s.SetCollection(db, nil)
|
||||
|
||||
obj.db = s
|
||||
|
||||
return obj, err
|
||||
}
|
||||
|
||||
func (self *KVStorage) Fetch(bitmap_id uint64, db string, slice int) IBitmap {
|
||||
|
||||
var (
|
||||
chunk *Chunk
|
||||
)
|
||||
last_key := COUNTER_KEY
|
||||
marker := COUNTER_KEY
|
||||
count := uint64(0)
|
||||
|
||||
key, _ := toKeyBytes(int64(bitmap_id), slice, 0, 0)
|
||||
|
||||
bitmap := CreateRBBitmap()
|
||||
|
||||
self.cc.VisitItemsAscend(key, true, func(i *gkvlite.Item) bool {
|
||||
bid, _, chunk_key, block_index, _ := fromKeyBytes(i.Key)
|
||||
if bid != int64(bitmap_id) {
|
||||
return false
|
||||
}
|
||||
block, _ := binary.Uvarint(i.Val) //just need to cast as a block
|
||||
if chunk_key != marker {
|
||||
if chunk_key != last_key {
|
||||
chunk = &Chunk{uint64(chunk_key), BlockArray{}}
|
||||
bitmap.AddChunk(chunk)
|
||||
}
|
||||
chunk.Value.Block[block_index] = block
|
||||
|
||||
} else {
|
||||
count = block
|
||||
}
|
||||
last_key = chunk_key
|
||||
return true
|
||||
})
|
||||
bitmap.SetCount(uint64(count))
|
||||
return bitmap
|
||||
}
|
||||
|
||||
func (self *KVStorage) Store(bitmap_id int64, db string, slice int, 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 {
|
||||
self.StoreBlock(bitmap_id, db, slice, int64(chunk.Key), block_index, iblock)
|
||||
}
|
||||
}
|
||||
}
|
||||
cnt := int64(BitCount(bitmap))
|
||||
|
||||
//var dumb = COUNTERMASK
|
||||
//COUNTER_KEY := int64(dumb)
|
||||
|
||||
self.StoreBlock(bitmap_id, db, slice, COUNTER_KEY, 0, cnt)
|
||||
self.db.Flush()
|
||||
return nil
|
||||
|
||||
}
|
||||
func Bytes(src int64) []byte {
|
||||
buf := new(bytes.Buffer)
|
||||
err := binary.Write(buf, binary.LittleEndian, src)
|
||||
if err != nil {
|
||||
fmt.Println("binary.Write failed:", err)
|
||||
}
|
||||
return buf.Bytes()
|
||||
}
|
||||
func fromKeyBytes(src []byte) (bitmap_id int64, slice int, chunk_key int64, block_index int32, err error) {
|
||||
buf := bytes.NewReader(src)
|
||||
err = binary.Read(buf, binary.LittleEndian, &bitmap_id)
|
||||
err = binary.Read(buf, binary.LittleEndian, &slice)
|
||||
err = binary.Read(buf, binary.LittleEndian, &chunk_key)
|
||||
err = binary.Read(buf, binary.LittleEndian, &block_index)
|
||||
return
|
||||
}
|
||||
|
||||
func toKeyBytes(bitmap_id int64, slice int, chunk_key int64, block_index int32) ([]byte, error) {
|
||||
buf := new(bytes.Buffer)
|
||||
err := binary.Write(buf, binary.LittleEndian, bitmap_id)
|
||||
err = binary.Write(buf, binary.LittleEndian, slice)
|
||||
err = binary.Write(buf, binary.LittleEndian, chunk_key)
|
||||
err = binary.Write(buf, binary.LittleEndian, block_index)
|
||||
|
||||
return buf.Bytes(), err
|
||||
}
|
||||
|
||||
func (self *KVStorage) StoreBlock(bitmap_id int64, db string, slice int, chunk_key int64, block_index int32, block int64) error {
|
||||
//key := fmt.Sprintf("%d%d:%d:%d", bitmap_id, slice, chunk_key, block_index)
|
||||
key, _ := toKeyBytes(bitmap_id, slice, chunk_key, block_index)
|
||||
value := Bytes(block)
|
||||
|
||||
self.cc.Set(key, value)
|
||||
|
||||
return nil
|
||||
}
|
||||
29
index/storage_test.go
Normal file
29
index/storage_test.go
Normal file
|
|
@ -0,0 +1,29 @@
|
|||
package index
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
// "io/ioutil"
|
||||
// "time"
|
||||
. "github.com/smartystreets/goconvey/convey"
|
||||
)
|
||||
|
||||
func TestStorage(t *testing.T) {
|
||||
db := "db"
|
||||
slice := 0
|
||||
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, 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)
|
||||
|
||||
})
|
||||
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue