From c6e3788a20ed6f45c116c4dbc27bfcbd9e45336b Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 21 Jan 2015 20:28:07 +0000 Subject: [PATCH] first pass at remote setbit, doesn't work --- core/etcd.go | 5 +++ core/http.go | 17 +++++--- core/remotebits.go | 99 +++++++++++++++++++++++++++++++++++++++++++ dispatch/dispatch.go | 11 +++++ index/storage_cass.go | 2 +- transport/tcp.go | 3 ++ 6 files changed, 129 insertions(+), 8 deletions(-) create mode 100644 core/remotebits.go diff --git a/core/etcd.go b/core/etcd.go index 9433f5a00..a5ac92386 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -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") diff --git a/core/http.go b/core/http.go index 0cd08ab2f..47fbb9201 100644 --- a/core/http.go +++ b/core/http.go @@ -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) diff --git a/core/remotebits.go b/core/remotebits.go new file mode 100644 index 000000000..e0e9c0239 --- /dev/null +++ b/core/remotebits.go @@ -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 +} diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index a010b7cf9..43bea66c9 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -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) diff --git a/index/storage_cass.go b/index/storage_cass.go index c0b403c0a..de198c530 100644 --- a/index/storage_cass.go +++ b/index/storage_cass.go @@ -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++ } diff --git a/transport/tcp.go b/transport/tcp.go index 1af9a9556..9a85739ad 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -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)