From 6230d820a1c3a9fdecfe0cde36465f0a3d06ee69 Mon Sep 17 00:00:00 2001 From: travisturner Date: Tue, 8 Apr 2014 23:32:01 -0500 Subject: [PATCH] bubble query-planner errors up to final Hold.Set (remove panics) --- executor/executor.go | 1 + query/planner.go | 159 ++++++++++++++++++++++++++++++------------ query/planner_test.go | 20 ++++-- query/query.go | 5 +- transport/tcp.go | 1 - 5 files changed, 133 insertions(+), 53 deletions(-) diff --git a/executor/executor.go b/executor/executor.go index 9278299db..7711f66d4 100644 --- a/executor/executor.go +++ b/executor/executor.go @@ -73,6 +73,7 @@ func (self *Executor) runQuery(database *db.Database, qry *query.Query) error { query_plan, err := query.QueryPlanForQuery(database, qry, &destination) if err != nil { + self.service.Hold.Set(qry.Id, err, 30) return err } // loop over the query steps and send to Transport diff --git a/query/planner.go b/query/planner.go index 5e4303708..1736f4a7c 100644 --- a/query/planner.go +++ b/query/planner.go @@ -70,7 +70,7 @@ type CountQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *CountQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *CountQueryTree) getLocation(d *db.Database) (*db.Location, error) { return qt.subquery.getLocation(d) } @@ -97,7 +97,7 @@ type TopNQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *TopNQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *TopNQueryTree) getLocation(d *db.Database) (*db.Location, error) { return qt.subquery.getLocation(d) } @@ -120,16 +120,17 @@ type UnionQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *UnionQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *UnionQueryTree) 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 = subquery.getLocation(d) + qt.location, err = subquery.getLocation(d) } } - return qt.location + return qt.location, err } /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -151,16 +152,17 @@ type IntersectQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *IntersectQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *IntersectQueryTree) 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 = subquery.getLocation(d) + qt.location, err = subquery.getLocation(d) } } - return qt.location + return qt.location, err } /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -184,16 +186,17 @@ type CatQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *CatQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *CatQueryTree) 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 = subquery.getLocation(d) + qt.location, err = subquery.getLocation(d) } } - return qt.location + return qt.location, err } /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -216,13 +219,13 @@ type GetQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *GetQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *GetQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - panic(err) + return nil, err } - return fragment.GetLocation() + return fragment.GetLocation(), nil } /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -245,13 +248,13 @@ type SetQueryTree struct { } // Uses consistent hashing function to select node containing data for GET operation -func (qt *SetQueryTree) getLocation(d *db.Database) *db.Location { +func (qt *SetQueryTree) getLocation(d *db.Database) (*db.Location, error) { slice, err := d.GetSliceForProfile(qt.profile_id) fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap) if err != nil { - panic(err) + return nil, err } - return fragment.GetLocation() + return fragment.GetLocation(), nil } /////////////////////////////////////////////////////////////////////////////////////////////////// @@ -285,18 +288,18 @@ type QueryPlanner struct { Query *Query } type QueryTree interface { - getLocation(d *db.Database) *db.Location + getLocation(d *db.Database) (*db.Location, error) } // Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices -func (qp *QueryPlanner) buildTree(query *Query, slice int) QueryTree { +func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) { var tree QueryTree // handle SET operation regardless of the slice if query.Operation == "set" { tree = &SetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), query.Args["filter"].(uint64)}, query.Args["profile_id"].(uint64)} - return tree + return tree, nil } // handle the remaining operations, taking slice into consideration @@ -309,20 +312,26 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) QueryTree { tree = &CatQueryTree{N: n} numSlices, err := qp.Database.NumSlices() if err != nil { - panic(err) + return nil, err } for slice := 0; slice < numSlices; slice++ { //for slice := 0; slice < 3; slice++ { - subtree := qp.buildTree(query, slice) + subtree, err := qp.buildTree(query, slice) + if err != nil { + return nil, err + } composite := tree.(*CatQueryTree) composite.subqueries = append(composite.subqueries, subtree) } } else { if query.Operation == "get" { tree = &GetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), 0}, slice} - return tree + return tree, nil } else if query.Operation == "count" { - subquery := qp.buildTree(&query.Subqueries[0], slice) + subquery, err := qp.buildTree(&query.Subqueries[0], slice) + if err != nil { + return nil, err + } tree = &CountQueryTree{subquery: subquery} } else if query.Operation == "top-n" { var n int @@ -336,86 +345,144 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) QueryTree { filters = filters_.([]uint64) } - subquery := qp.buildTree(&query.Subqueries[0], slice) + subquery, err := qp.buildTree(&query.Subqueries[0], slice) + if err != nil { + return nil, err + } tree = &TopNQueryTree{subquery: subquery, Filters: filters, N: n} } else if query.Operation == "union" { subqueries := make([]QueryTree, len(query.Subqueries)) + var err error for i, query := range query.Subqueries { - subqueries[i] = qp.buildTree(&query, slice) + subqueries[i], err = qp.buildTree(&query, slice) + if err != nil { + return nil, err + } } tree = &UnionQueryTree{subqueries: subqueries} } else if query.Operation == "intersect" { subqueries := make([]QueryTree, len(query.Subqueries)) + var err error for i, query := range query.Subqueries { - subqueries[i] = qp.buildTree(&query, slice) + subqueries[i], err = qp.buildTree(&query, slice) + if err != nil { + return nil, err + } } tree = &IntersectQueryTree{subqueries: subqueries} } else { + //TODO return error gracefully panic("invalid operation") } } - return tree + return tree, nil } // Produces flattened QueryPlan from QueryTree input -func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Location) *QueryPlan { +func (qp *QueryPlanner) flatten(qt QueryTree, id *uuid.UUID, location *db.Location) (*QueryPlan, error) { plan := QueryPlan{} if cat, ok := qt.(*CatQueryTree); ok { inputs := make([]*uuid.UUID, len(cat.subqueries)) - step := CatQueryStep{&BaseQueryStep{id, "cat", cat.getLocation(qp.Database), location}, inputs, cat.N} + loc, err := cat.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := CatQueryStep{&BaseQueryStep{id, "cat", loc, location}, inputs, cat.N} for index, subq := range cat.subqueries { sub_id := uuid.RandomUUID() step.Inputs[index] = &sub_id - subq_steps := qp.flatten(subq, &sub_id, cat.getLocation(qp.Database)) + 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 union, ok := qt.(*UnionQueryTree); ok { inputs := make([]*uuid.UUID, len(union.subqueries)) - step := UnionQueryStep{&BaseQueryStep{id, "union", union.getLocation(qp.Database), location}, inputs} + loc, err := union.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := UnionQueryStep{&BaseQueryStep{id, "union", loc, location}, inputs} 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)) + 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 intersect, ok := qt.(*IntersectQueryTree); ok { inputs := make([]*uuid.UUID, len(intersect.subqueries)) - step := IntersectQueryStep{&BaseQueryStep{id, "intersect", intersect.getLocation(qp.Database), location}, inputs} + loc, err := intersect.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := IntersectQueryStep{&BaseQueryStep{id, "intersect", loc, location}, inputs} 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)) + 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 { - step := GetQueryStep{&BaseQueryStep{id, "get", get.getLocation(qp.Database), location}, get.bitmap, get.slice} + loc, err := get.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := GetQueryStep{&BaseQueryStep{id, "get", loc, location}, get.bitmap, get.slice} plan := QueryPlan{step} - return &plan + return &plan, nil } else if set, ok := qt.(*SetQueryTree); ok { - step := SetQueryStep{&BaseQueryStep{id, "set", set.getLocation(qp.Database), location}, set.bitmap, set.profile_id} + loc, err := set.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := SetQueryStep{&BaseQueryStep{id, "set", loc, location}, set.bitmap, set.profile_id} plan := QueryPlan{step} - return &plan + return &plan, nil } else if cnt, ok := qt.(*CountQueryTree); ok { sub_id := uuid.RandomUUID() - step := &CountQueryStep{&BaseQueryStep{id, "count", cnt.getLocation(qp.Database), location}, &sub_id} - subq_steps := qp.flatten(cnt.subquery, &sub_id, cnt.getLocation(qp.Database)) + loc, err := cnt.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := &CountQueryStep{&BaseQueryStep{id, "count", loc, location}, &sub_id} + subq_steps, err := qp.flatten(cnt.subquery, &sub_id, loc) + if err != nil { + return nil, err + } plan = append(plan, *subq_steps...) plan = append(plan, step) } else if topn, ok := qt.(*TopNQueryTree); ok { sub_id := uuid.RandomUUID() - step := &TopNQueryStep{&BaseQueryStep{id, "top-n", topn.getLocation(qp.Database), location}, &sub_id, topn.Filters, topn.N} - subq_steps := qp.flatten(topn.subquery, &sub_id, topn.getLocation(qp.Database)) + loc, err := topn.getLocation(qp.Database) + if err != nil { + return nil, err + } + step := &TopNQueryStep{&BaseQueryStep{id, "top-n", loc, location}, &sub_id, topn.Filters, topn.N} + subq_steps, err := qp.flatten(topn.subquery, &sub_id, loc) + if err != nil { + return nil, err + } plan = append(plan, *subq_steps...) plan = append(plan, step) } - return &plan + return &plan, nil } // Transforms Query into QueryTree and flattens to QueryPlan object -func (qp *QueryPlanner) Plan(query *Query, id *uuid.UUID, destination *db.Location) *QueryPlan { - queryTree := qp.buildTree(query, -1) +func (qp *QueryPlanner) Plan(query *Query, id *uuid.UUID, destination *db.Location) (*QueryPlan, error) { + queryTree, err := qp.buildTree(query, -1) + if err != nil { + return nil, err + } return qp.flatten(queryTree, query.Id, destination) } diff --git a/query/planner_test.go b/query/planner_test.go index 66503d0e1..2e720c07f 100644 --- a/query/planner_test.go +++ b/query/planner_test.go @@ -50,8 +50,10 @@ func TestQueryPlanner(t *testing.T) { destination := fragment1.GetLocation() id := uuid.RandomUUID() - qp := *qplanner.Plan(&query, &id, destination) + qpp, err := qplanner.Plan(&query, &id, destination) + qp := *qpp + So(err, ShouldEqual, nil) So(len(qp), ShouldEqual, 7) So(qp[0].(GetQueryStep).Operation, ShouldEqual, "get") So(qp[0].(GetQueryStep).Slice, ShouldEqual, 0) @@ -92,7 +94,9 @@ func TestQueryPlanner(t *testing.T) { destination := fragment1.GetLocation() id := uuid.RandomUUID() - qp := *qplanner.Plan(query, &id, destination) + qpp, err := qplanner.Plan(query, &id, destination) + qp := *qpp + So(err, ShouldEqual, nil) So(len(qp), ShouldEqual, 3) So(qp[0].(GetQueryStep).Operation, ShouldEqual, "get") @@ -118,7 +122,9 @@ func TestQueryPlanner(t *testing.T) { destination := fragment1.GetLocation() id := uuid.RandomUUID() - qp := *qplanner.Plan(query, &id, destination) + qpp, err := qplanner.Plan(query, &id, destination) + qp := *qpp + So(err, ShouldEqual, nil) So(len(qp), ShouldEqual, 7) So(qp[0].(GetQueryStep).Operation, ShouldEqual, "get") @@ -159,7 +165,9 @@ func TestQueryPlanner(t *testing.T) { destination := fragment1.GetLocation() id := uuid.RandomUUID() - qp := *qplanner.Plan(query, &id, destination) + qpp, err := qplanner.Plan(query, &id, destination) + qp := *qpp + So(err, ShouldEqual, nil) So(len(qp), ShouldEqual, 1) So(qp[0].(SetQueryStep).Operation, ShouldEqual, "set") So(qp[0].(SetQueryStep).ProfileId, ShouldEqual, 100) @@ -175,7 +183,9 @@ func TestQueryPlanner(t *testing.T) { destination := fragment1.GetLocation() id := uuid.RandomUUID() - qp := *qplanner.Plan(query, &id, destination) + qpp, err := qplanner.Plan(query, &id, destination) + qp := *qpp + So(err, ShouldEqual, nil) So(len(qp), ShouldEqual, 5) So(qp[0].(GetQueryStep).Operation, ShouldEqual, "get") So(*(qp[0].(GetQueryStep).Bitmap), ShouldResemble, db.Bitmap{10, "general", 0}) diff --git a/query/query.go b/query/query.go index a521b307c..76e560b6b 100644 --- a/query/query.go +++ b/query/query.go @@ -63,7 +63,10 @@ func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.L func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Location) (*QueryPlan, error) { query_planner := QueryPlanner{Database: database, Query: query} id := uuid.RandomUUID() - query_plan := query_planner.Plan(query, &id, destination) + query_plan, err := query_planner.Plan(query, &id, destination) + if err != nil { + return nil, err + } return query_plan, nil } diff --git a/transport/tcp.go b/transport/tcp.go index 6e7032cc9..be0e94f89 100644 --- a/transport/tcp.go +++ b/transport/tcp.go @@ -34,7 +34,6 @@ func init() { func (self *connection) manage() { BeginManageConnection: for { - log.Println("manage", self) if self.conn == nil { process, err := self.transport.service.ProcessMap.GetProcess(self.process) if err != nil {