first pass at remote setbit, doesn't work

This commit is contained in:
Todd Gruben 2015-01-21 20:28:07 +00:00
parent 7d59066c3f
commit c6e3788a20
6 changed files with 129 additions and 8 deletions

View file

@ -7,6 +7,7 @@ import (
"pilosa/config"
"pilosa/db"
"pilosa/util"
"runtime/debug"
"sort"
"strconv"
"strings"
@ -299,6 +300,10 @@ func (self *ProcessMap) AddProcess(process *db.Process) {
func (self *ProcessMap) GetProcess(id *util.GUID) (*db.Process, error) {
self.mutex.Lock()
defer self.mutex.Unlock()
if id == nil {
debug.PrintStack()
return nil, errors.New("Nil process")
}
process, ok := self.nodes[*id]
if !ok {
return nil, errors.New("No such process")

View file

@ -398,7 +398,7 @@ func bitmaps(frame string, obj JsonObject) chan uint64 {
type SBResult struct {
Bitmap_id uint64
Frame string
Filter int
Filter uint64
Profile_id uint64
Result interface{}
}
@ -450,8 +450,9 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) {
return
}
t = float64(obj["filter"].(float64))
filter := int(t)
filter := uint64(t)
result := false
remoteSetBit := NewRemoteSetBit(self.service)
for bitmap_id := range bitmaps(frame, obj) {
start := time.Now()
@ -461,14 +462,14 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) {
//no fragment
self.service.TopologyMapper.MakeFragments(dbs, db.GetSlice(profile_id))
time.Sleep(2 * time.Second)
continue
break
}
if util.Equal(frag.GetProcessId(), self.service.Id) {
// The Local Route
result, _ = self.service.Index.SetBit(frag.GetId(), bitmap_id, profile_id, uint64(filter))
result, _ = self.service.Index.SetBit(frag.GetId(), bitmap_id, profile_id, filter)
} else {
println("remote")
// The Remote Route
remoteSetBit.Add(frag, bitmap_id, profile_id, filter, frame)
}
//result, err := self.service.Executor.RunPQL(db, pql)
@ -483,9 +484,11 @@ func (self *WebService) HandleSetBit(w http.ResponseWriter, r *http.Request) {
}
bundle := SBResult{bitmap_id, frame, filter, profile_id, result}
results = append(results, bundle)
// results = append(results, db+" "+pql)
}
remoteSetBit.Request()
remoteSetBit.MergeResults(results)
}
encoder := json.NewEncoder(w)
err = encoder.Encode(results)

99
core/remotebits.go Normal file
View file

@ -0,0 +1,99 @@
package core
import (
"log"
"pilosa/db"
"pilosa/util"
)
type RemoteSetBit struct {
requests []util.GUID
cluster map[*util.GUID][]BitmapRequestItem
service *Service
}
type BitsRequest struct {
Bits []BitmapRequestItem
ReturnProcessId util.GUID
QueryId util.GUID
DestProcessId util.GUID
}
type BitmapRequestItem struct {
Fragment_id util.SUUID
Bitmap_id uint64
Profile_id uint64
Filter uint64
Frame string
}
func NewRemoteSetBit(s *Service) *RemoteSetBit {
obj := new(RemoteSetBit)
obj.cluster = make(map[*util.GUID][]BitmapRequestItem)
obj.service = s
return obj
}
func (self *RemoteSetBit) Request() {
self.requests = make([]util.GUID, len(self.cluster), len(self.cluster))
source_process, _ := self.service.GetProcess()
for process, request := range self.cluster {
random_id := util.RandomUUID()
msg := new(db.Message)
msg.Data = BitsRequest{
Bits: request,
ReturnProcessId: source_process.Id(),
QueryId: random_id,
DestProcessId: *process,
}
self.requests = append(self.requests, random_id)
self.service.Transport.Send(msg, process)
}
}
func (self *RemoteSetBit) MergeResults(local_results []SBResult) []SBResult {
answers := make(chan []SBResult)
for _, task := range self.requests {
go func(id util.GUID) {
value, err := self.service.Hold.Get(&id, 10) //eiher need to be the frame process or the handler process?
if value == nil {
log.Println("Bad RemoteSetBi Result:", err)
empty := make([]SBResult, 0, 0)
answers <- empty
} else {
answers <- value.([]SBResult)
}
}(task)
}
for i := 0; i < len(self.requests); i++ {
batch := <-answers
for _, item := range batch {
local_results = append(local_results, item)
}
}
return local_results
}
func (self *RemoteSetBit) Add(frag *db.Fragment, bitmap_id, profile_id, filter uint64, frame string) {
x, found := self.cluster[frag.GetProcessId()]
if !found {
x = make([]BitmapRequestItem, 8)
}
x = append(x, BitmapRequestItem{frag.GetId(), bitmap_id, profile_id, filter, frame})
self.cluster[frag.GetProcessId()] = x
}
type BitsResponse struct {
Id *util.GUID
Items []SBResult
}
func (self *BitsResponse) ResultId() *util.GUID {
return self.Id
}
func (self *BitsResponse) ResultData() interface{} {
return self.Items
}

View file

@ -31,6 +31,17 @@ func (self *Dispatch) Run() {
response := db.Message{Data: core.BatchResponse{Id: data.Id}}
self.service.Index.LoadBitmap(data.Fragment_id, data.Bitmap_id, data.Compressed_bitmap, data.Filter)
self.service.Transport.Send(&response, data.Source)
case core.BitsRequest:
var results []core.SBResult
for _, v := range data.Bits {
result, _ := self.service.Index.SetBit(v.Fragment_id, v.Bitmap_id, v.Profile_id, uint64(v.Filter))
//jbundle := core.SBResult{v.Bitmap_id, ''v.Frame, v.Filter, v.Profile_id, result}
bundle := core.SBResult{v.Bitmap_id, v.Frame, v.Filter, v.Profile_id, result}
results = append(results, bundle)
}
response := db.Message{Data: core.BitsResponse{Id: &data.QueryId, Items: results}}
self.service.Transport.Send(&response, &data.ReturnProcessId)
case core.PingRequest:
pong := db.Message{Data: core.PongRequest{Id: data.Id}}
self.service.Transport.Send(&pong, data.Source)

View file

@ -107,7 +107,7 @@ func (c *CassandraStorage) Fetch(bitmap_id uint64, db string, frame string, slic
}
func (self *CassandraStorage) BeginBatch() {
if self.batch == nil {
self.batch = gocql.NewBatch(gocql.LoggedBatch)
self.batch = gocql.NewBatch(gocql.UnloggedBatch)
}
self.batch_counter++
}

View file

@ -10,6 +10,7 @@ import (
"pilosa/core"
"pilosa/db"
. "pilosa/util"
"runtime/debug"
"sync"
"time"
@ -38,6 +39,7 @@ func init() {
}
func newConnection(transport *TcpTransport, conn net.Conn, proc *GUID) *connection {
println("New Connection", proc)
p := new(connection)
p.transport = transport
p.outbox = make(chan *db.Message, 2048)
@ -72,6 +74,7 @@ func (self *connection) Shutdown() {
func (self *connection) serviceConnection() {
var host_string string
if self.conn == nil {
println("Service Connection", self.process)
process, err := self.transport.service.ProcessMap.GetProcess(self.process)
if err != nil {
log.Println("transport/tcp: error getting process, retrying in 2 seconds... ", self.process, err)