From 62a9961a15eb3b364cbc24aed1235d29f14df6f9 Mon Sep 17 00:00:00 2001 From: travisturner Date: Fri, 25 Jul 2014 16:36:02 -0500 Subject: [PATCH] add parser/planner support for the index.Difference method --- core/etcd.go | 2 +- core/query.go | 35 +++++++++++++++++++++++++ db/topology.go | 2 +- executor/executor.go | 4 ++- query/planner.go | 61 ++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 101 insertions(+), 3 deletions(-) diff --git a/core/etcd.go b/core/etcd.go index 4332ed425..145dcefc4 100644 --- a/core/etcd.go +++ b/core/etcd.go @@ -118,7 +118,7 @@ func (self *TopologyMapper) MakeFragments(db string, slice_int int) error { id_string := self.service.Id.String() m[id_string] = 0 //at least have one process if none created for _, dbs := range self.service.Cluster.GetDatabases() { - for _, fsi := range dbs.GetFramSliceIntersects() { + for _, fsi := range dbs.GetFrameSliceIntersects() { for _, fragment := range fsi.GetFragments() { process := fragment.GetProcess().Id().String() if len(process) > 1 { diff --git a/core/query.go b/core/query.go index d5af82c9f..09eb4b732 100644 --- a/core/query.go +++ b/core/query.go @@ -133,6 +133,41 @@ func (self *Service) IntersectQueryStepHandler(msg *db.Message) { self.Transport.Send(&result_message, qs.Destination.ProcessId) } +func (self *Service) DifferenceQueryStepHandler(msg *db.Message) { + //spew.Dump("DIFFERENCE QUERYSTEP") + qs := msg.Data.(query.DifferenceQueryStep) + var handles []index.BitmapHandle + // create a list of bitmap handles + for _, input := range qs.Inputs { + value, _ := self.Hold.Get(input, 10) + switch val := value.(type) { + case index.BitmapHandle: + handles = append(handles, val) + case []byte: + bh, _ := self.Index.FromBytes(qs.Location.FragmentId, val) + handles = append(handles, bh) + } + } + + bh, err := self.Index.Difference(qs.Location.FragmentId, handles) + if err != nil { + spew.Dump(err) + } + + var result interface{} + if qs.LocIsDest() { + result = bh + } else { + bm, err := self.Index.GetBytes(qs.Location.FragmentId, bh) + if err != nil { + spew.Dump(err) + } + result = bm + } + result_message := db.Message{Data: query.DifferenceQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: result}}} + self.Transport.Send(&result_message, qs.Destination.ProcessId) +} + func (self *Service) CatQueryStepHandler(msg *db.Message) { qs := msg.Data.(query.CatQueryStep) //spew.Dump("CAT QUERYSTEP") diff --git a/db/topology.go b/db/topology.go index bb3b25542..824ee550a 100644 --- a/db/topology.go +++ b/db/topology.go @@ -128,7 +128,7 @@ type Database struct { mutex sync.Mutex } -func (self *Database) GetFramSliceIntersects() []*FrameSliceIntersect { +func (self *Database) GetFrameSliceIntersects() []*FrameSliceIntersect { return self.frame_slice_intersects } diff --git a/executor/executor.go b/executor/executor.go index f297510ea..a24923778 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -36,6 +36,8 @@ func (self *Executor) NewJob(job *db.Message) { self.service.UnionQueryStepHandler(job) case query.IntersectQueryStep: self.service.IntersectQueryStepHandler(job) + case query.DifferenceQueryStep: + self.service.DifferenceQueryStepHandler(job) case query.CatQueryStep: self.service.CatQueryStepHandler(job) case query.GetQueryStep: @@ -97,7 +99,7 @@ func (self *Executor) RunPQL(database_name string, pql string) (interface{}, err database := self.service.Cluster.GetOrCreateDatabase(database_name) // see if the outer query function is a custom query - reserved_functions := stringSlice{"get", "set", "union", "intersect", "count", "top-n"} + reserved_functions := stringSlice{"get", "set", "union", "intersect", "difference", "count", "top-n"} tokens, err := query.Lex(pql) if err != nil { return nil, err diff --git a/query/planner.go b/query/planner.go index b08829e5f..4ef76306a 100644 --- a/query/planner.go +++ b/query/planner.go @@ -219,6 +219,38 @@ func (qt *IntersectQueryTree) getLocation(d *db.Database) (*db.Location, error) return qt.location, err } +/////////////////////////////////////////////////////////////////////////////////////////////////// +// DIFFERENCE +/////////////////////////////////////////////////////////////////////////////////////////////////// +type DifferenceQueryStep struct { + *BaseQueryStep + Inputs []*util.GUID +} + +type DifferenceQueryResult struct { + *BaseQueryResult +} + +// QueryTree for UNION queries +type DifferenceQueryTree struct { + subqueries []QueryTree + location *db.Location +} + +// Uses consistent hashing function to select node containing data for GET operation +func (qt *DifferenceQueryTree) getLocation(d *db.Database) (*db.Location, error) { + var err error + if qt.location == nil { + subqueryLength := len(qt.subqueries) + if subqueryLength > 0 { + locationIndex := rand.Intn(subqueryLength) + subquery := qt.subqueries[locationIndex] + qt.location, err = subquery.getLocation(d) + } + } + return qt.location, err +} + /////////////////////////////////////////////////////////////////////////////////////////////////// // CAT /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -328,6 +360,7 @@ func init() { gob.Register(CatQueryResult{}) gob.Register(UnionQueryResult{}) gob.Register(IntersectQueryResult{}) + gob.Register(DifferenceQueryResult{}) gob.Register(CountQueryResult{}) gob.Register(TopNQueryResult{}) @@ -336,6 +369,7 @@ func init() { gob.Register(CatQueryStep{}) gob.Register(UnionQueryStep{}) gob.Register(IntersectQueryStep{}) + gob.Register(DifferenceQueryStep{}) gob.Register(CountQueryStep{}) gob.Register(TopNQueryStep{}) } @@ -438,6 +472,16 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { } } tree = &IntersectQueryTree{subqueries: subqueries} + } else if query.Operation == "difference" { + subqueries := make([]QueryTree, len(query.Subqueries)) + var err error + for i, query := range query.Subqueries { + subqueries[i], err = qp.buildTree(&query, slice) + if err != nil { + return nil, err + } + } + tree = &DifferenceQueryTree{subqueries: subqueries} } else { //TODO return error gracefully log.Println(spew.Sdump(query)) @@ -501,6 +545,23 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati plan = append(plan, *subq_steps...) } plan = append(plan, step) + } else if difference, ok := qt.(*DifferenceQueryTree); ok { + inputs := make([]*util.GUID, len(difference.subqueries)) + loc, err := difference.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := DifferenceQueryStep{&BaseQueryStep{id, "difference", loc, location}, inputs} + for index, subq := range difference.subqueries { + sub_id := util.RandomUUID() + step.Inputs[index] = &sub_id + subq_steps, err := qp.flatten(subq, &sub_id, loc) + if err != nil { + return nil, err + } + plan = append(plan, *subq_steps...) + } + plan = append(plan, step) } else if get, ok := qt.(*GetQueryTree); ok { loc, err := get.getLocation(qp.Database) if err != nil {