From 266f772c17de59361a6110c6fa4ee7c7fe34b1f9 Mon Sep 17 00:00:00 2001 From: travisturner Date: Tue, 7 Jan 2014 12:23:01 -0600 Subject: [PATCH] got some basic cat/count/bitmap-handles working --- executor/executor.go | 52 +++++++++++++++++++++++++++++++++++++++----- query/planner.go | 7 ++++++ transport/tcp.go | 3 +-- 3 files changed, 54 insertions(+), 8 deletions(-) diff --git a/executor/executor.go b/executor/executor.go index e633b025e..2503887f5 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -34,17 +34,48 @@ func (self *Executor) NewJob(job *db.Message) { qs := job.Data.(query.CatQueryStep) fmt.Println("CAT QUERYSTEP") spew.Dump(qs) + var bhs []index.BitmapHandle + use_sum := false + sum := 0 for _, input := range qs.Inputs { spew.Dump(input) - bhi, err := self.service.Hold.Get(input, 10) - bh := bhi.(index.BitmapHandle) - // TODO: git rid of this count, need the cat to do a sum() or a true cat() - count, err := self.service.Index.Count(qs.Location.FragmentId, bh) + bhi, _ := self.service.Hold.Get(input, 10) + spew.Dump(bhi) + switch bhi.(type) { + case index.BitmapHandle: + spew.Dump("BH") + bhs = append(bhs, bhi.(index.BitmapHandle)) + case []byte: + spew.Dump("RAW") + bh, _ := self.service.Index.FromBytes(qs.Location.FragmentId, bhi.([]byte)) + bhs = append(bhs, bh) + case int: + use_sum = true + sum += bhi.(int) + } + } + + if use_sum { + spew.Dump("SUM", sum) + } else { + + spew.Dump("BHS", bhs) + + unionbh, err := self.service.Index.Union(qs.Location.FragmentId, bhs) if err != nil { spew.Dump(err) } - spew.Dump("COUNT", count) + spew.Dump("UNION BH", unionbh) + + /* + count, err := self.service.Index.Count(qs.Location.FragmentId, unionbh) + if err != nil { + spew.Dump(err) + } + spew.Dump("FINAL COUNT", count) + */ } + case query.GetQueryStep: qs := job.Data.(query.GetQueryStep) @@ -58,7 +89,16 @@ func (self *Executor) NewJob(job *db.Message) { } //spew.Dump("COUNT", count) // push results to the map - self.service.Hold.Set(qs.Id, bh, 10) + + if qs.LocIsDest() { + self.service.Hold.Set(qs.Id, bh, 10) + } else { + bm, err := self.service.Index.GetBytes(qs.Location.FragmentId, bh) + if err != nil { + spew.Dump(err) + } + self.service.Hold.Set(qs.Id, bm, 10) + } case query.SetQueryStep: qs := job.Data.(query.SetQueryStep) diff --git a/query/planner.go b/query/planner.go index 3c13560d0..d5e6439a6 100644 --- a/query/planner.go +++ b/query/planner.go @@ -34,6 +34,13 @@ type GetQueryStep struct { Destination *db.Location } +func (qs GetQueryStep) LocIsDest() bool { + if qs.Location.ProcessId == qs.Destination.ProcessId && qs.Location.FragmentId == qs.Destination.FragmentId { + return true + } + return false +} + type SetQueryStep struct { Id *uuid.UUID Operation string diff --git a/transport/tcp.go b/transport/tcp.go index 3b9c0b409..c8f242571 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -8,7 +8,6 @@ import ( "pilosa/config" "pilosa/core" "pilosa/db" - "time" "tux21b.org/v1/gocql/uuid" ) @@ -46,7 +45,7 @@ func (self *TcpTransport) Run() { log.Println(err.Error()) return } - time.Sleep(1 * time.Second) + //time.Sleep(1 * time.Second) } } }