mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-10 12:57:54 +00:00
Union/Intersect Query & Result support
This commit is contained in:
parent
71c782ae6a
commit
b2c9b1860d
3 changed files with 207 additions and 8 deletions
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
141
query/planner.go
141
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}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue