changed cassanadra storage to handle client timestamp bug

This commit is contained in:
Todd Gruben 2014-10-17 20:15:16 +00:00
parent 938995fe36
commit a9133b4b10
6 changed files with 63 additions and 42 deletions

View file

@ -441,7 +441,9 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) {
encoder := json.NewEncoder(w)
err = encoder.Encode(results)
if err != nil {
log.Fatal("Error encoding stats")
log.Println("Error encoding stats")
http.Error(w, "Error econding stats", http.StatusMethodNotAllowed)
return
}
}
@ -478,7 +480,8 @@ func (self *WebService) HandleStats(w http.ResponseWriter, r *http.Request) {
err := encoder.Encode(m)
if err != nil {
log.Fatal("Error encoding stats")
log.Println("Error encoding stats")
http.Error(w, "Error econding stats", http.StatusMethodNotAllowed)
}
}

View file

@ -1,13 +1,14 @@
package index
import (
"fmt"
// "fmt"
"testing"
"time"
// "time"
. "github.com/smartystreets/goconvey/convey"
// . "github.com/smartystreets/goconvey/convey"
)
/*
func TestBitmaps(t *testing.T) {
Convey("function BitCount should equal method bm.Count()", t, func() {
bm := CreateRBBitmap()
@ -115,3 +116,11 @@ func TestBitmaps(t *testing.T) {
})
}
*/
func BenchmarkSetBit(b *testing.B) {
// run the Fib function b.N times
a := CreateRBBitmap()
for n := 0; n < b.N; n++ {
SetBit(a, uint64(n))
}
}

View file

@ -394,10 +394,17 @@ func (self *Brand) Load(requestChan chan Command, f *Fragment) {
//log.Println("Bad mojo")
}
time.Sleep(time.Duration(rand.Intn(15)) * time.Second) //trying to avoid mass cassandra hit
counter := 0
for _, k := range keys {
request := NewLoadRequest(k)
requestChan <- request
request.Response()
counter++
if counter > 100 {
backoff := (counter / 100) * 100
time.Sleep(time.Duration(backoff) * time.Millisecond) //trying to avoid mass cassandra hit
}
}
}

View file

@ -287,7 +287,7 @@ func (self *FragmentContainer) AddFragment(db string, frame string, slice int, i
loader := make(chan Command)
self.fragments[id] = f
go f.ServeFragment(loader)
//go f.Load(loader)
go f.Load(loader)
}
}
@ -467,19 +467,19 @@ func (self *Fragment) processCommand(req Command) {
func (self *Fragment) ServeFragment(loadChan chan Command) {
for {
select {
// case req := <-self.requestChan:
// self.processCommand(req)
// default:
// select {
case req := <-self.requestChan:
self.processCommand(req)
case req := <-loadChan:
self.processCommand(req)
case wg := <-self.exit:
log.Println("Fragment Shutdown")
self.Persist()
wg.Done()
// }
default:
select {
case req := <-self.requestChan:
self.processCommand(req)
case req := <-loadChan:
self.processCommand(req)
case wg := <-self.exit:
log.Println("Fragment Shutdown")
self.Persist()
wg.Done()
}
}
}
}

View file

@ -52,7 +52,7 @@ func NewCassStorage() Storage {
}
obj.db = session
obj.stmt = `INSERT INTO bitmap ( bitmap_id, db, frame, slice , filter, ChunkKey, BlockIndex, block) VALUES (?,?,?,?,?,?,?,?);`
obj.stmt = `INSERT INTO bitmap ( bitmap_id, db, frame, slice , filter, ChunkKey, BlockIndex, block) VALUES (?,?,?,?,?,?,?,?) USING timestamp ?;`
obj.batch = nil
obj.batch_time = time.Now()
obj.batch_counter = 0
@ -168,7 +168,7 @@ func (self *CassandraStorage) StoreBlock(id int64, db string, frame string, slic
self.BeginBatch()
}
start := time.Now()
self.batch.Query(self.stmt, id, db, frame, slice, int(filter), chunk, block_index, block)
self.batch.Query(self.stmt, id, db, frame, slice, int(filter), chunk, block_index, block, start.UnixNano())
delta := time.Since(start)
util.SendTimer("cassandra_storage_StoreBlock", delta.Nanoseconds())
return nil

View file

@ -17,7 +17,7 @@ func TestStorage(t *testing.T) {
frame := "main"
slice := 0
filter := 10
bitmap_id := uint64(1234)
bitmap_id := uint64(999999)
/* Convey("KV ", t, func() {
storage, _ := NewKVStorage("/tmp/", 0, db)
bm := storage.Fetch(bitmap_id, db, slice)
@ -44,38 +44,40 @@ func TestStorage(t *testing.T) {
fmt.Println("FETCH")
bm, _ := storage.Fetch(bitmap_id, db, frame, slice)
SetBit(bm, 0)
SetBit(bm, 1)
SetBit(bm, 2)
for i := uint64(0); i < 256; i++ {
SetBit(bm, i)
}
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)
So(BitCount(bm), ShouldEqual, 256)
})
}
Convey("leveldb", t, func() {
storage := NewLevelDBStorage("./basic/one")
/*
Convey("leveldb", t, func() {
storage := NewLevelDBStorage("./basic/one")
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()
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()
})
})
*/
}