mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
add parser/planner support for the index.Difference method
This commit is contained in:
parent
c0058c0904
commit
62a9961a15
5 changed files with 101 additions and 3 deletions
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue