From b2c9b1860d66323691b2c6768e18c9867e0bef29 Mon Sep 17 00:00:00 2001 From: travisturner Date: Wed, 8 Jan 2014 12:11:02 -0600 Subject: [PATCH] Union/Intersect Query & Result support --- dispatch/dispatch.go | 6 +- executor/executor.go | 68 +++++++++++++++++++++ query/planner.go | 141 +++++++++++++++++++++++++++++++++++++++++-- 3 files changed, 207 insertions(+), 8 deletions(-) diff --git a/dispatch/dispatch.go b/dispatch/dispatch.go index 96c8106bb..0e0a3cb2e 100644 --- a/dispatch/dispatch.go +++ b/dispatch/dispatch.go @@ -24,18 +24,16 @@ func (self *Dispatch) Run() { log.Println("Dispatch Run...") for { message := self.service.Transport.Receive() - //spew.Dump("Processing ", message) switch data := message.Data.(type) { case core.PingRequest: pong := db.Message{Data: core.PongRequest{Id: data.Id}} self.service.Transport.Send(&pong, data.Source) case db.HoldResult: self.service.Hold.Set(data.ResultId(), data.ResultData(), 30) - case query.CatQueryStep, query.GetQueryStep, query.SetQueryStep, query.CountQueryStep: - //fmt.Println("CAT/GET/SET QUERYSTEP") + case query.CatQueryStep, query.GetQueryStep, query.SetQueryStep, query.CountQueryStep, query.UnionQueryStep, query.IntersectQueryStep: go self.service.Executor.NewJob(message) default: - //log.Println("Unprocessed message", data) + log.Println("Unprocessed message", data) } } } diff --git a/executor/executor.go b/executor/executor.go index 05c3e1b51..1e46a328d 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -52,6 +52,74 @@ func (self *Executor) NewJob(job *db.Message) { result_message := db.Message{Data: query.CountQueryResult{Id: qs.Id, Data: count}} self.service.Transport.Send(&result_message, qs.Destination.ProcessId) + case query.UnionQueryStep: + qs := job.Data.(query.UnionQueryStep) + spew.Dump("UNION QUERYSTEP") + var handles []index.BitmapHandle + // create a list of bitmap handles + for _, input := range qs.Inputs { + value, _ := self.service.Hold.Get(input, 10) + switch val := value.(type) { + case index.BitmapHandle: + handles = append(handles, val) + case []byte: + bh, _ := self.service.Index.FromBytes(qs.Location.FragmentId, val) + handles = append(handles, bh) + } + } + + bh, err := self.service.Index.Union(qs.Location.FragmentId, handles) + if err != nil { + spew.Dump(err) + } + + var result interface{} + if qs.LocIsDest() { + result = bh + } else { + bm, err := self.service.Index.GetBytes(qs.Location.FragmentId, bh) + if err != nil { + spew.Dump(err) + } + result = bm + } + result_message := db.Message{Data: query.UnionQueryResult{Id: qs.Id, Data: result}} + self.service.Transport.Send(&result_message, qs.Destination.ProcessId) + + case query.IntersectQueryStep: + qs := job.Data.(query.IntersectQueryStep) + spew.Dump("INTERSECT QUERYSTEP") + var handles []index.BitmapHandle + // create a list of bitmap handles + for _, input := range qs.Inputs { + value, _ := self.service.Hold.Get(input, 10) + switch val := value.(type) { + case index.BitmapHandle: + handles = append(handles, val) + case []byte: + bh, _ := self.service.Index.FromBytes(qs.Location.FragmentId, val) + handles = append(handles, bh) + } + } + + bh, err := self.service.Index.Intersect(qs.Location.FragmentId, handles) + if err != nil { + spew.Dump(err) + } + + var result interface{} + if qs.LocIsDest() { + result = bh + } else { + bm, err := self.service.Index.GetBytes(qs.Location.FragmentId, bh) + if err != nil { + spew.Dump(err) + } + result = bm + } + result_message := db.Message{Data: query.IntersectQueryResult{Id: qs.Id, Data: result}} + self.service.Transport.Send(&result_message, qs.Destination.ProcessId) + case query.CatQueryStep: qs := job.Data.(query.CatQueryStep) spew.Dump("CAT QUERYSTEP") diff --git a/query/planner.go b/query/planner.go index 6582c3623..333fa03a3 100644 --- a/query/planner.go +++ b/query/planner.go @@ -48,9 +48,8 @@ func (self CountQueryResult) ResultData() interface{} { // QueryTree for COUNT queries type CountQueryTree struct { - operation string - subquery QueryTree - location *db.Location + subquery QueryTree + location *db.Location } // Uses consistent hashing function to select node containing data for GET operation @@ -58,6 +57,106 @@ func (qt *CountQueryTree) getLocation(d *db.Database) *db.Location { return qt.subquery.getLocation(d) } +/////////////////////////////////////////////////////////////////////////////////////////////////// +// UNION +/////////////////////////////////////////////////////////////////////////////////////////////////// +type UnionQueryStep struct { + Id *uuid.UUID + Operation string + Inputs []*uuid.UUID + Location *db.Location + Destination *db.Location +} + +func (qs UnionQueryStep) LocIsDest() bool { + if qs.Location.ProcessId == qs.Destination.ProcessId && qs.Location.FragmentId == qs.Destination.FragmentId { + return true + } + return false +} + +type UnionQueryResult struct { + Id *uuid.UUID + Data interface{} +} + +func (self UnionQueryResult) ResultId() *uuid.UUID { + return self.Id +} + +func (self UnionQueryResult) ResultData() interface{} { + return self.Data +} + +// QueryTree for UNION queries +type UnionQueryTree struct { + subqueries []QueryTree + location *db.Location +} + +// Uses consistent hashing function to select node containing data for GET operation +func (qt *UnionQueryTree) getLocation(d *db.Database) *db.Location { + if qt.location == nil { + subqueryLength := len(qt.subqueries) + if subqueryLength > 0 { + locationIndex := rand.Intn(subqueryLength) + subquery := qt.subqueries[locationIndex] + qt.location = subquery.getLocation(d) + } + } + return qt.location +} + +/////////////////////////////////////////////////////////////////////////////////////////////////// +// INTERSECT +/////////////////////////////////////////////////////////////////////////////////////////////////// +type IntersectQueryStep struct { + Id *uuid.UUID + Operation string + Inputs []*uuid.UUID + Location *db.Location + Destination *db.Location +} + +func (qs IntersectQueryStep) LocIsDest() bool { + if qs.Location.ProcessId == qs.Destination.ProcessId && qs.Location.FragmentId == qs.Destination.FragmentId { + return true + } + return false +} + +type IntersectQueryResult struct { + Id *uuid.UUID + Data interface{} +} + +func (self IntersectQueryResult) ResultId() *uuid.UUID { + return self.Id +} + +func (self IntersectQueryResult) ResultData() interface{} { + return self.Data +} + +// QueryTree for UNION queries +type IntersectQueryTree struct { + subqueries []QueryTree + location *db.Location +} + +// Uses consistent hashing function to select node containing data for GET operation +func (qt *IntersectQueryTree) getLocation(d *db.Database) *db.Location { + if qt.location == nil { + subqueryLength := len(qt.subqueries) + if subqueryLength > 0 { + locationIndex := rand.Intn(subqueryLength) + subquery := qt.subqueries[locationIndex] + qt.location = subquery.getLocation(d) + } + } + return qt.location +} + /////////////////////////////////////////////////////////////////////////////////////////////////// // CAT /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -180,6 +279,8 @@ func (qt *SetQueryTree) getLocation(d *db.Database) *db.Location { func init() { gob.Register(CatQueryResult{}) + gob.Register(UnionQueryResult{}) + gob.Register(IntersectQueryResult{}) gob.Register(GetQueryResult{}) gob.Register(CountQueryResult{}) } @@ -245,7 +346,19 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) QueryTree { return tree } else if query.Operation == "count" { subquery := qp.buildTree(query.Inputs[0].(*Query), slice) - tree = &CountQueryTree{operation: query.Operation, subquery: subquery} + tree = &CountQueryTree{subquery: subquery} + } else if query.Operation == "union" { + subqueries := make([]QueryTree, len(query.Inputs)) + for i, input := range query.Inputs { + subqueries[i] = qp.buildTree(input.(*Query), slice) + } + tree = &UnionQueryTree{subqueries: subqueries} + } else if query.Operation == "intersect" { + subqueries := make([]QueryTree, len(query.Inputs)) + for i, input := range query.Inputs { + subqueries[i] = qp.buildTree(input.(*Query), slice) + } + tree = &IntersectQueryTree{subqueries: subqueries} } else { subqueries := make([]QueryTree, len(query.Inputs)) for i, input := range query.Inputs { @@ -280,6 +393,26 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Locati plan = append(plan, *subq_steps...) } plan = append(plan, step) + } else if union, ok := qt.(*UnionQueryTree); ok { + inputs := make([]*uuid.UUID, len(union.subqueries)) + step := UnionQueryStep{id, "union", inputs, union.getLocation(qp.Database), location} + for index, subq := range union.subqueries { + sub_id := uuid.RandomUUID() + step.Inputs[index] = &sub_id + subq_steps := qp.flatten(subq, &sub_id, union.getLocation(qp.Database)) + plan = append(plan, *subq_steps...) + } + plan = append(plan, step) + } else if intersect, ok := qt.(*IntersectQueryTree); ok { + inputs := make([]*uuid.UUID, len(intersect.subqueries)) + step := IntersectQueryStep{id, "intersect", inputs, intersect.getLocation(qp.Database), location} + for index, subq := range intersect.subqueries { + sub_id := uuid.RandomUUID() + step.Inputs[index] = &sub_id + subq_steps := qp.flatten(subq, &sub_id, intersect.getLocation(qp.Database)) + plan = append(plan, *subq_steps...) + } + plan = append(plan, step) } else if get, ok := qt.(*GetQueryTree); ok { step := GetQueryStep{id, "get", get.bitmap, get.slice, get.getLocation(qp.Database), location} plan := QueryPlan{step}