initial compile of hole filling

This commit is contained in:
Todd Gruben 2014-08-08 19:33:23 +00:00
parent 345689bbe9
commit d7be05b216
16 changed files with 751 additions and 70 deletions

14
core/core_test.go Normal file
View file

@ -0,0 +1,14 @@
package core
import (
"testing"
. "github.com/smartystreets/goconvey/convey"
)
func TestCompile(t *testing.T) {
Convey("has asm", t, func() {
So(canCompile(), ShouldEqual, true)
})
}

View file

@ -4,6 +4,7 @@ import (
"pilosa/db"
"pilosa/index"
"pilosa/query"
"pilosa/util"
"sort"
"github.com/davecgh/go-spew/spew"
@ -30,39 +31,6 @@ func (self *Service) CountQueryStepHandler(msg *db.Message) {
self.Transport.Send(&result_message, qs.Destination.ProcessId)
}
func (self *Service) TopNQueryStepHandler(msg *db.Message) {
//spew.Dump("TOP-N QUERYSTEP")
qs := msg.Data.(query.TopNQueryStep)
//spew.Dump(qs)
//need categories in qs I just added so it would compile
var categoryleaves []uint64
input := qs.Input
value, _ := self.Hold.Get(input, 10)
var bh index.BitmapHandle
switch val := value.(type) {
case index.BitmapHandle:
bh = val
case []byte:
bh, _ = self.Index.FromBytes(qs.Location.FragmentId, val)
}
categoryleaves = qs.Filters
/*
spew.Dump(bh, qs.N*2, categoryleaves)
result_message := db.Message{Data: "foobar"}
self.Transport.Send(&result_message, qs.Destination.ProcessId)
*/
topn, err := self.Index.TopN(qs.Location.FragmentId, bh, qs.N*2, categoryleaves)
if err != nil {
spew.Dump(err)
}
result_message := db.Message{Data: query.TopNQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: topn}}}
self.Transport.Send(&result_message, qs.Destination.ProcessId)
}
func (self *Service) UnionQueryStepHandler(msg *db.Message) {
//spew.Dump("UNION QUERYSTEP")
qs := msg.Data.(query.UnionQueryStep)
@ -174,10 +142,28 @@ func (self *Service) CatQueryStepHandler(msg *db.Message) {
var handles []index.BitmapHandle
return_type := "bitmap-handles"
var sum uint64
merge_map := map[uint64]uint64{}
merge_map := make(map[uint64]uint64)
slice_map := make(map[uint64]map[util.SUUID]struct{})
all_slice := make(map[util.SUUID]struct {
process util.GUID
handle index.BitmapHandle
})
// either create a list of bitmap handles to cat (i.e. union), or sum the integer values
part := make(chan interface{})
num_parts := len(qs.Inputs)
for _, input := range qs.Inputs {
value, _ := self.Hold.Get(input, 10)
go func(id *util.GUID, part chan interface{}) {
value, _ := self.Hold.Get(id, 10)
part <- value
}(input, part)
}
//for _, input := range qs.Inputs {
for i := 0; i < num_parts; i++ {
value := <-part
switch val := value.(type) {
case index.BitmapHandle:
handles = append(handles, val)
@ -188,15 +174,32 @@ func (self *Service) CatQueryStepHandler(msg *db.Message) {
//spew.Dump(val)
return_type = "sum"
sum += val
case []index.Pair:
//spew.Dump(val)
case TopNPackage:
return_type = "pair-list"
for _, pair := range val {
var e struct{}
for _, pair := range val.Pairs {
//merge_map[pair.Key] += pair.Count
merge_map[pair.Key] += pair.Count
mm, ok := slice_map[pair.Key]
if !ok {
mm = make(map[util.SUUID]struct{})
slice_map[pair.Key] = mm
}
mm[val.FragmentId] = e
}
all_slice[val.FragmentId] = struct {
process util.GUID
handle index.BitmapHandle
}{val.ProcessId, val.HBitmap}
}
}
tasks := BuildTask(merge_map, slice_map, all_slice)
FetchMissing(tasks, self)
for k, v := range GatherResults(tasks, self) {
merge_map[k] += v
}
// either return the sum, or return the compressed bitmap resulting from the cat (union)
var result interface{}
if return_type == "sum" {

191
core/topn.go Normal file
View file

@ -0,0 +1,191 @@
package core
import (
"log"
"pilosa/db"
"pilosa/index"
"pilosa/query"
"pilosa/util"
"github.com/davecgh/go-spew/spew"
)
type Task struct {
processid util.GUID
f map[util.SUUID]index.FillArgs
hold_id util.GUID
}
type TopFill struct {
Args []index.FillArgs
ReturnProcessId util.GUID
QueryId util.GUID
DestProcessId util.GUID
}
//portable query step
func (self *TopFill) GetId() *util.GUID {
return &self.QueryId
}
func (self *TopFill) GetLocation() *db.Location {
return &db.Location{&self.DestProcessId, 0} //this message is a broadcast to many fragments so i'm choosing fragmentzero
}
//
type hole struct {
process util.GUID
handle index.BitmapHandle
fragment util.SUUID
}
func missing(fids map[util.SUUID]struct{}, all map[util.SUUID]struct {
process util.GUID
handle index.BitmapHandle
}) []hole {
results := make([]hole, 0, 0)
for k, v := range all {
_, ok := fids[k]
if ok {
results = append(results, hole{v.process, v.handle, k})
}
}
return results
}
func newtask(p util.GUID) *Task {
result := new(Task)
result.processid = p
result.f = make(map[util.SUUID]index.FillArgs)
result.hold_id = util.RandomUUID()
return result
}
func (t *Task) Add(frag util.SUUID, bitmap_id uint64, handle index.BitmapHandle) {
fa, ok := t.f[frag]
if ok {
fa = index.FillArgs{frag, handle, make([]uint64, 0, 0)}
}
fa.Bitmaps = append(fa.Bitmaps, bitmap_id)
t.f[frag] = fa
}
func BuildTask(merge_map map[uint64]uint64,
slice_map map[uint64]map[util.SUUID]struct{},
total_fragments map[util.SUUID]struct {
process util.GUID
handle index.BitmapHandle
}) map[util.GUID]*Task {
tasks := make(map[util.GUID]*Task)
for bitmap_id, _ := range merge_map { //for all brands
//for fragment_id, reported_fragments := range slice_map[bitmap_id] { //find missing fragments
reporting_fragments := slice_map[bitmap_id]
//id slice ==> SUUID,BitmapHandle
for _, p := range missing(reporting_fragments, total_fragments) {
task, ok := tasks[p.process]
if ok {
task = newtask(p.process)
tasks[p.process] = task
}
task.Add(p.fragment, bitmap_id, p.handle)
}
//}
}
return tasks
}
func (self *Service) TopFillHandler(msg *db.Message) { //in order for this to get executed it needs to be a portable query step
topfill := msg.Data.(TopFill)
topn, err := self.Index.TopFillBatch(topfill.Args)
if err != nil {
log.Println("TopFileHandler:", err)
}
sendmsg := new(db.Message)
sendmsg.Data = query.BaseQueryResult{Id: &topfill.QueryId, Data: topn}
self.Transport.Send(sendmsg, &topfill.ReturnProcessId)
}
func SendRequest(process_id util.GUID, t *Task, service *Service) {
args := make([]index.FillArgs, len(t.f), len(t.f))
for _, v := range t.f {
args = append(args, v)
}
msg := new(db.Message)
p, _ := service.GetProcess()
msg.Data = TopFill{args, p.Id(), t.hold_id, process_id}
service.Transport.Send(msg, &process_id)
}
func FetchMissing(tasks map[util.GUID]*Task, service *Service) {
for k, v := range tasks {
go SendRequest(k, v, service)
}
}
func GatherResults(tasks map[util.GUID]*Task, service *Service) map[uint64]uint64 {
results := make(map[uint64]uint64)
answers := make(chan []index.Pair)
for _, task := range tasks {
go func(id util.GUID) {
value, _ := service.Hold.Get(&id, 10) //eiher need to be the frame process or the handler process?
answers <- value.([]index.Pair)
}(task.hold_id)
}
for i := 0; i < len(tasks); i++ {
batch := <-answers
for _, pair := range batch {
results[pair.Key] += pair.Count
}
}
close(answers)
return results
}
type TopNPackage struct {
ProcessId util.GUID
FragmentId util.SUUID
Pairs []index.Pair
HBitmap index.BitmapHandle
}
func (self *Service) TopNQueryStepHandler(msg *db.Message) {
//spew.Dump("TOP-N QUERYSTEP")
qs := msg.Data.(query.TopNQueryStep)
//spew.Dump(qs)
//need categories in qs I just added so it would compile
var categoryleaves []uint64
input := qs.Input
value, _ := self.Hold.Get(input, 10)
var bh index.BitmapHandle
switch val := value.(type) {
case index.BitmapHandle:
bh = val
case []byte:
bh, _ = self.Index.FromBytes(qs.Location.FragmentId, val)
}
categoryleaves = qs.Filters
/*
spew.Dump(bh, qs.N*2, categoryleaves)
result_message := db.Message{Data: "foobar"}
self.Transport.Send(&result_message, qs.Destination.ProcessId)
*/
topn, err := self.Index.TopN(qs.Location.FragmentId, bh, qs.N*2, categoryleaves)
topnPackage := TopNPackage{*qs.Location.ProcessId, qs.Location.FragmentId, topn, bh}
if err != nil {
spew.Dump(err)
}
result_message := db.Message{Data: query.TopNQueryResult{&query.BaseQueryResult{Id: qs.Id, Data: topnPackage}}}
self.Transport.Send(&result_message, qs.Destination.ProcessId)
}
func canCompile() bool {
return true
}

View file

@ -23,8 +23,8 @@ func TestTopology(t *testing.T) {
spew.Dump(fragment_id)
database.GetOrCreateFragment(frame, slice, fragment_id)
spew.Dump(database)
spew.Dump("DONE")
// spew.Dump(database)
// spew.Dump("DONE")
})
}

View file

@ -44,6 +44,10 @@ func (self *Executor) NewJob(job *db.Message) {
self.service.GetQueryStepHandler(job)
case query.SetQueryStep:
self.service.SetQueryStepHandler(job)
// case query.MaskQueryStep:
// self.service.MaskQueryStepHandler(job)
// case query.RangeQueryStep:
// self.service.RangeQueryStepHandler(job)
default:
fmt.Println("unknown")
}
@ -99,7 +103,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", "difference", "count", "top-n"}
reserved_functions := stringSlice{"get", "set", "union", "intersect", "difference", "count", "top-n", "mask", "range"}
tokens, err := query.Lex(pql)
if err != nil {
return nil, err

View file

@ -67,7 +67,6 @@ func popcount(i uint64)uint64{
//x:= uint64(val)
return uint64(val)
}
*/
func popcount(x uint64) (n uint64) {
// bit population count, see
// http://graphics.stanford.edu/~seander/bithacks.html#CountBitsSetParallel
@ -78,6 +77,7 @@ func popcount(x uint64) (n uint64) {
x *= 0x0101010101010101
return uint64(x >> 56)
}
*/
type BlockArray struct {
Block [32]uint64
@ -89,6 +89,7 @@ func (s *BlockArray) bitcount() uint64 {
sum += popcount(b)
}
return sum
// return popcntSlice(s.Block)
}
func BlockArray_union(a *BlockArray, b *BlockArray) BlockArray {
var o = BlockArray{}

View file

@ -291,3 +291,23 @@ func (self *CmdMask) Execute(f *Fragment) Calculation {
}
return f.AllocHandle(result)
}
type CmdTopFill struct {
*Responder
args FillArgs
}
func NewTopFill(a FillArgs) *CmdTopFill {
return &CmdTopFill{NewResponder("TopFill"), a}
}
func (self *CmdTopFill) Execute(f *Fragment) Calculation {
result := make([]Pair, len(self.args.Bitmaps))
for _, v := range self.args.Bitmaps {
a := f.NewHandle(v)
res := f.intersect([]BitmapHandle{self.args.Handle, a})
bm, _ := f.getBitmap(res)
result = append(result, Pair{v, BitCount(bm)})
}
return result
}

View file

@ -39,6 +39,11 @@ func NewFragmentContainer() *FragmentContainer {
}
type BitmapHandle uint64
type FillArgs struct {
Frag_id util.SUUID
Handle BitmapHandle
Bitmaps []uint64
}
func init() {
var vh BitmapHandle
@ -103,7 +108,6 @@ func (self *FragmentContainer) Intersect(frag_id util.SUUID, bh []BitmapHandle)
}
return 0, errors.New("Invalid Bitmap Handle")
}
func (self *FragmentContainer) Union(frag_id util.SUUID, bh []BitmapHandle) (BitmapHandle, error) {
if fragment, found := self.GetFragment(frag_id); found {
request := NewUnion(bh)
@ -169,6 +173,35 @@ func (self *FragmentContainer) TopN(frag_id util.SUUID, bh BitmapHandle, n int,
}
return nil, errors.New(fmt.Sprintf("Fragment not found:%s", util.SUUID_to_Hex(frag_id)))
}
func (self *FragmentContainer) TopFillBatch(args []FillArgs) ([]Pair, error) {
//should probaly make this concurrent but then all hell breaks lose
results := make(map[uint64]uint64)
for _, v := range args {
items, _ := self.TopFillFragment(v)
if len(args) == 1 {
return items, nil
}
for _, pair := range items {
results[pair.Key] += pair.Count
}
}
ret_val := make([]Pair, len(results))
for k, v := range results {
ret_val = append(ret_val, Pair{k, v})
}
return ret_val, nil
}
func (self *FragmentContainer) TopFillFragment(arg FillArgs) ([]Pair, error) {
if fragment, found := self.GetFragment(arg.Frag_id); found {
request := NewTopFill(arg)
fragment.requestChan <- request
result := request.Response()
util.SendTimer("fragmant_container_TopFillFragment", result.exec_time.Nanoseconds())
return result.answer.([]Pair), nil
}
return nil, errors.New("Invalid Bitmap Handle")
}
func (self *FragmentContainer) GetList(frag_id util.SUUID, bitmap_id []uint64) ([]BitmapHandle, error) {
if fragment, found := self.GetFragment(frag_id); found {

53
index/popcnt.go Normal file
View file

@ -0,0 +1,53 @@
package index
// bit population count, take from
// https://code.google.com/p/go/issues/detail?id=4988#c11
// credit: https://code.google.com/u/arnehormann/
func popcount(x uint64) (n uint64) {
x -= (x >> 1) & 0x5555555555555555
x = (x>>2)&0x3333333333333333 + x&0x3333333333333333
x += x >> 4
x &= 0x0f0f0f0f0f0f0f0f
x *= 0x0101010101010101
return x >> 56
}
func popcntSliceGo(s []uint64) uint64 {
cnt := uint64(0)
for _, x := range s {
cnt += popcount(x)
}
return cnt
}
func popcntMaskSliceGo(s, m []uint64) uint64 {
cnt := uint64(0)
for i := range s {
cnt += popcount(s[i] &^ m[i])
}
return cnt
}
func popcntAndSliceGo(s, m []uint64) uint64 {
cnt := uint64(0)
for i := range s {
cnt += popcount(s[i] & m[i])
}
return cnt
}
func popcntOrSliceGo(s, m []uint64) uint64 {
cnt := uint64(0)
for i := range s {
cnt += popcount(s[i] | m[i])
}
return cnt
}
func popcntXorSliceGo(s, m []uint64) uint64 {
cnt := uint64(0)
for i := range s {
cnt += popcount(s[i] ^ m[i])
}
return cnt
}

102
index/popcnt_amd64.s Normal file
View file

@ -0,0 +1,102 @@
TEXT ·hasAsm(SB),4,$0
MOVQ $1, AX
CPUID
SHRQ $23, CX
ANDQ $1, CX
MOVB CX, ret+0(FP)
RET
#define POPCNTQ_DX_DX BYTE $0xf3; BYTE $0x48; BYTE $0x0f; BYTE $0xb8; BYTE $0xd2
TEXT ·popcntSliceAsm(SB),4,$0-32
XORQ AX, AX
MOVQ s+0(FP), SI
MOVQ s+8(FP), CX
TESTQ CX, CX
JZ popcntSliceEnd
popcntSliceLoop:
BYTE $0xf3; BYTE $0x48; BYTE $0x0f; BYTE $0xb8; BYTE $0x16 // POPCNTQ (SI), DX
ADDQ DX, AX
ADDQ $8, SI
LOOP popcntSliceLoop
popcntSliceEnd:
MOVQ AX, ret+24(FP)
RET
TEXT ·popcntMaskSliceAsm(SB),4,$0-56
XORQ AX, AX
MOVQ s+0(FP), SI
MOVQ s+8(FP), CX
TESTQ CX, CX
JZ popcntMaskSliceEnd
MOVQ m+24(FP), DI
popcntMaskSliceLoop:
MOVQ (DI), DX
NOTQ DX
ANDQ (SI), DX
POPCNTQ_DX_DX
ADDQ DX, AX
ADDQ $8, SI
ADDQ $8, DI
LOOP popcntMaskSliceLoop
popcntMaskSliceEnd:
MOVQ AX, ret+48(FP)
RET
TEXT ·popcntAndSliceAsm(SB),4,$0-56
XORQ AX, AX
MOVQ s+0(FP), SI
MOVQ s+8(FP), CX
TESTQ CX, CX
JZ popcntAndSliceEnd
MOVQ m+24(FP), DI
popcntAndSliceLoop:
MOVQ (DI), DX
ANDQ (SI), DX
POPCNTQ_DX_DX
ADDQ DX, AX
ADDQ $8, SI
ADDQ $8, DI
LOOP popcntAndSliceLoop
popcntAndSliceEnd:
MOVQ AX, ret+48(FP)
RET
TEXT ·popcntOrSliceAsm(SB),4,$0-56
XORQ AX, AX
MOVQ s+0(FP), SI
MOVQ s+8(FP), CX
TESTQ CX, CX
JZ popcntOrSliceEnd
MOVQ m+24(FP), DI
popcntOrSliceLoop:
MOVQ (DI), DX
ORQ (SI), DX
POPCNTQ_DX_DX
ADDQ DX, AX
ADDQ $8, SI
ADDQ $8, DI
LOOP popcntOrSliceLoop
popcntOrSliceEnd:
MOVQ AX, ret+48(FP)
RET
TEXT ·popcntXorSliceAsm(SB),4,$0-56
XORQ AX, AX
MOVQ s+0(FP), SI
MOVQ s+8(FP), CX
TESTQ CX, CX
JZ popcntXorSliceEnd
MOVQ m+24(FP), DI
popcntXorSliceLoop:
MOVQ (DI), DX
XORQ (SI), DX
POPCNTQ_DX_DX
ADDQ DX, AX
ADDQ $8, SI
ADDQ $8, DI
LOOP popcntXorSliceLoop
popcntXorSliceEnd:
MOVQ AX, ret+48(FP)
RET

64
index/popcnt_asm.go Normal file
View file

@ -0,0 +1,64 @@
// +build amd64
package index
//go:noescape
func hasAsm() bool
var useAsm = hasAsm()
//go:noescape
func popcntSliceAsm(s []uint64) uint64
//go:noescape
func popcntMaskSliceAsm(s, m []uint64) uint64
//go:noescape
func popcntAndSliceAsm(s, m []uint64) uint64
//go:noescape
func popcntOrSliceAsm(s, m []uint64) uint64
//go:noescape
func popcntXorSliceAsm(s, m []uint64) uint64
func popcntSlice(s []uint64) uint64 {
if useAsm {
return popcntSliceAsm(s)
}
return popcntSliceGo(s)
}
func popcntMaskSlice(s, m []uint64) uint64 {
if useAsm {
return popcntMaskSliceAsm(s, m)
}
return popcntMaskSliceGo(s, m)
}
func popcntAndSlice(s, m []uint64) uint64 {
if useAsm {
return popcntAndSliceAsm(s, m)
}
return popcntAndSliceGo(s, m)
}
func popcntOrSlice(s, m []uint64) uint64 {
if useAsm {
return popcntOrSliceAsm(s, m)
}
return popcntOrSliceGo(s, m)
}
func popcntXorSlice(s, m []uint64) uint64 {
if useAsm {
return popcntXorSliceAsm(s, m)
}
return popcntXorSliceGo(s, m)
}

23
index/popcnt_generic.go Normal file
View file

@ -0,0 +1,23 @@
// +build !amd64
package index
func popcntSlice(s []uint64) uint64 {
return popcntSliceGo(s)
}
func popcntMaskSlice(s, m []uint64) uint64 {
return popcntMaskSliceGo(s, m)
}
func popcntAndSlice(s, m []uint64) uint64 {
return popcntAndSliceGo(s, m)
}
func popcntOrSlice(s, m []uint64) uint64 {
return popcntOrSliceGo(s, m)
}
func popcntXorSlice(s, m []uint64) uint64 {
return popcntXorSliceGo(s, m)
}

View file

@ -2,6 +2,7 @@ package index
import (
"fmt"
"log"
"testing"
"time"
@ -9,6 +10,30 @@ import (
. "github.com/smartystreets/goconvey/convey"
)
func getTime(id uint64, s string) {
const shortForm = "2006-01-02 15:04"
t1, _ := time.Parse(shortForm, s)
for i, v := range GetTimeIds(uint64(id), t1, YMDH) {
log.Println(i, v, s)
}
log.Println()
}
func TestDemo(t *testing.T) {
Convey("Test ID", t, func() {
getTime(uint64(666), "2014-01-01 00:00")
getTime(uint64(666), "2014-02-01 00:00")
getTime(uint64(666), "2014-03-01 00:00")
getTime(uint64(666), "2014-04-01 00:00")
getTime(uint64(666), "2014-04-02 00:00")
getTime(uint64(666), "2014-04-03 00:00")
getTime(uint64(666), "2014-04-04 00:00")
getTime(uint64(666), "2014-04-05 00:00")
getTime(uint64(666), "2014-05-01 00:00")
getTime(uint64(666), "2014-06-01 00:00")
So(1, ShouldEqual, 1)
})
}
func TestTimeFrame(t *testing.T) {
//print get_YMD_id(2014,3,28,1234)

View file

@ -6,6 +6,7 @@ import (
"log"
"pilosa/util"
"strconv"
"time"
"github.com/davecgh/go-spew/spew"
)
@ -81,6 +82,7 @@ func (self *QueryParser) Parse() (query *Query, err error) {
if token.Type != TYPE_LP {
return nil, fmt.Errorf("Expected '(', found token %v.", token)
}
const shortForm = "2006-01-02 15:04"
ArgLoop:
for {
@ -99,6 +101,46 @@ ArgLoop:
query.Subqueries = append(query.Subqueries, *subquery)
case TYPE_VALUE:
switch query.Operation {
case "range":
switch len(query.Args) {
case 0:
i, err := strconv.ParseUint(token.Text, 10, 64)
if err != nil {
return nil, fmt.Errorf("Expecting integer id! (%v)", err)
}
query.Args["id"] = i
case 1:
query.Args["frame"] = token.Text
case 2:
t, err := time.Parse(shortForm, token.Text)
if err != nil {
return nil, fmt.Errorf("Expecting integer DateTime (%v)", err)
}
query.Args["start"] = t
case 3:
t, err := time.Parse(shortForm, token.Text)
if err != nil {
return nil, fmt.Errorf("Expecting integer DateTime (%v)", err)
}
query.Args["end"] = t
}
case "mask":
switch len(query.Args) {
case 0:
query.Args["frame"] = token.Text
case 1:
i, err := strconv.ParseUint(token.Text, 10, 64)
if err != nil {
return nil, fmt.Errorf("Expecting integer id! (%v)", err)
}
query.Args["start"] = i
case 2:
i, err := strconv.ParseUint(token.Text, 10, 64)
if err != nil {
return nil, fmt.Errorf("Expecting integer id! (%v)", err)
}
query.Args["end"] = i
}
case "get":
switch len(query.Args) {
case 0:
@ -227,6 +269,12 @@ ArgLoop:
if query.Operation == "difference" {
return nil, fmt.Errorf("No Args Given")
}
if query.Operation == "range" {
return nil, fmt.Errorf("No Args Given")
}
if query.Operation == "mask" {
return nil, fmt.Errorf("No Args Given")
}
}
return query, nil
}

View file

@ -8,6 +8,7 @@ import (
"math/rand"
"pilosa/db"
"pilosa/util"
"time"
"github.com/davecgh/go-spew/spew"
)
@ -146,7 +147,7 @@ func (qt *TopNQueryTree) getLocation(d *db.Database) (*db.Location, error) {
slice := d.GetOrCreateSlice(qt.Slice)
fragment, err := d.GetFragmentForFrameSlice(frame, slice)
if err != nil {
log.Println("GetFragmentForFrameSliceFailed GetQueryTree", frame, slice)
log.Println("GetFragmentForFrameSliceFailed TopNQueryTree", frame, slice)
return nil, err
}
qt.location = fragment.GetLocation()
@ -389,7 +390,7 @@ type QueryTree interface {
}
// Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices
func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
var tree QueryTree
// handle SET operation regardless of the slice
@ -407,13 +408,13 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
n = n_.(int)
}
tree = &CatQueryTree{N: n}
numSlices, err := qp.Database.NumSlices()
numSlices, err := self.Database.NumSlices()
if err != nil {
return nil, err
}
for slice := 0; slice < numSlices; slice++ {
//for slice := 0; slice < 3; slice++ {
subtree, err := qp.buildTree(query, slice)
subtree, err := self.buildTree(query, slice)
if err != nil {
return nil, err
}
@ -425,7 +426,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
tree = &GetQueryTree{&db.Bitmap{query.Args["id"].(uint64), query.Args["frame"].(string), 0}, slice}
return tree, nil
} else if query.Operation == "count" {
subquery, err := qp.buildTree(&query.Subqueries[0], slice)
subquery, err := self.buildTree(&query.Subqueries[0], slice)
if err != nil {
return nil, err
}
@ -447,7 +448,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
filters = filters_.([]uint64)
}
subquery, err := qp.buildTree(&query.Subqueries[0], slice)
subquery, err := self.buildTree(&query.Subqueries[0], slice)
if err != nil {
return nil, err
}
@ -456,7 +457,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
subqueries := make([]QueryTree, len(query.Subqueries))
var err error
for i, query := range query.Subqueries {
subqueries[i], err = qp.buildTree(&query, slice)
subqueries[i], err = self.buildTree(&query, slice)
if err != nil {
return nil, err
}
@ -466,7 +467,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
subqueries := make([]QueryTree, len(query.Subqueries))
var err error
for i, query := range query.Subqueries {
subqueries[i], err = qp.buildTree(&query, slice)
subqueries[i], err = self.buildTree(&query, slice)
if err != nil {
return nil, err
}
@ -476,7 +477,7 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
subqueries := make([]QueryTree, len(query.Subqueries))
var err error
for i, query := range query.Subqueries {
subqueries[i], err = qp.buildTree(&query, slice)
subqueries[i], err = self.buildTree(&query, slice)
if err != nil {
return nil, err
}
@ -492,11 +493,11 @@ func (qp *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
}
// Produces flattened QueryPlan from QueryTree input
func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) {
func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) {
plan := QueryPlan{}
if cat, ok := qt.(*CatQueryTree); ok {
inputs := make([]*util.GUID, len(cat.subqueries))
loc, err := cat.getLocation(qp.Database)
loc, err := cat.getLocation(self.Database)
if err != nil {
return nil, err
}
@ -504,7 +505,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
for index, subq := range cat.subqueries {
sub_id := util.RandomUUID()
step.Inputs[index] = &sub_id
subq_steps, err := qp.flatten(subq, &sub_id, loc)
subq_steps, err := self.flatten(subq, &sub_id, loc)
if err != nil {
return nil, err
}
@ -513,7 +514,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
plan = append(plan, step)
} else if union, ok := qt.(*UnionQueryTree); ok {
inputs := make([]*util.GUID, len(union.subqueries))
loc, err := union.getLocation(qp.Database)
loc, err := union.getLocation(self.Database)
if err != nil {
return nil, err
}
@ -521,7 +522,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
for index, subq := range union.subqueries {
sub_id := util.RandomUUID()
step.Inputs[index] = &sub_id
subq_steps, err := qp.flatten(subq, &sub_id, loc)
subq_steps, err := self.flatten(subq, &sub_id, loc)
if err != nil {
return nil, err
}
@ -530,7 +531,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
plan = append(plan, step)
} else if intersect, ok := qt.(*IntersectQueryTree); ok {
inputs := make([]*util.GUID, len(intersect.subqueries))
loc, err := intersect.getLocation(qp.Database)
loc, err := intersect.getLocation(self.Database)
if err != nil {
return nil, err
}
@ -538,7 +539,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
for index, subq := range intersect.subqueries {
sub_id := util.RandomUUID()
step.Inputs[index] = &sub_id
subq_steps, err := qp.flatten(subq, &sub_id, loc)
subq_steps, err := self.flatten(subq, &sub_id, loc)
if err != nil {
return nil, err
}
@ -547,7 +548,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
plan = append(plan, step)
} else if difference, ok := qt.(*DifferenceQueryTree); ok {
inputs := make([]*util.GUID, len(difference.subqueries))
loc, err := difference.getLocation(qp.Database)
loc, err := difference.getLocation(self.Database)
if err != nil {
return nil, err
}
@ -555,7 +556,7 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
for index, subq := range difference.subqueries {
sub_id := util.RandomUUID()
step.Inputs[index] = &sub_id
subq_steps, err := qp.flatten(subq, &sub_id, loc)
subq_steps, err := self.flatten(subq, &sub_id, loc)
if err != nil {
return nil, err
}
@ -563,15 +564,33 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
}
plan = append(plan, step)
} else if get, ok := qt.(*GetQueryTree); ok {
loc, err := get.getLocation(qp.Database)
loc, err := get.getLocation(self.Database)
if err != nil {
return nil, err
}
step := GetQueryStep{&BaseQueryStep{id, "get", loc, location}, get.bitmap, get.slice}
plan := QueryPlan{step}
return &plan, nil
/*
} else if mask, ok := qt.(*MaskQueryTree); ok {
loc, err := mask.getLocation(self.Database)
if err != nil {
return nil, err
}
step := MaskQueryStep{&BaseQueryStep{id, "mask", loc, location}, mask.start, mask.end}
plan := QueryPlan{step}
return &plan, nil
} else if rang, ok := qt.(*RangeQueryTree); ok {
loc, err := rang.getLocation(self.Database)
if err != nil {
return nil, err
}
step := RangeQueryStep{&BaseQueryStep{id, "range", loc, location}, rang.bitmap, rang.start, rang.end}
plan := QueryPlan{step}
return &plan, nil
*/
} else if set, ok := qt.(*SetQueryTree); ok {
loc, err := set.getLocation(qp.Database)
loc, err := set.getLocation(self.Database)
if err != nil {
return nil, err
}
@ -580,12 +599,12 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
return &plan, nil
} else if cnt, ok := qt.(*CountQueryTree); ok {
sub_id := util.RandomUUID()
loc, err := cnt.getLocation(qp.Database)
loc, err := cnt.getLocation(self.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)
subq_steps, err := self.flatten(cnt.subquery, &sub_id, loc)
if err != nil {
return nil, err
}
@ -593,12 +612,12 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
plan = append(plan, step)
} else if topn, ok := qt.(*TopNQueryTree); ok {
sub_id := util.RandomUUID()
loc, err := topn.getLocation(qp.Database)
loc, err := topn.getLocation(self.Database)
if err != nil {
return nil, err
}
step := &TopNQueryStep{&BaseQueryStep{id, "top-n", loc, location}, &sub_id, topn.Filters, topn.N, topn.Frame}
subq_steps, err := qp.flatten(topn.subquery, &sub_id, loc)
subq_steps, err := self.flatten(topn.subquery, &sub_id, loc)
if err != nil {
return nil, err
}
@ -609,10 +628,71 @@ func (qp *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Locati
}
// Transforms Query into QueryTree and flattens to QueryPlan object
func (qp *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) {
queryTree, err := qp.buildTree(query, -1)
func (self *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) {
queryTree, err := self.buildTree(query, -1)
if err != nil {
return nil, err
}
return qp.flatten(queryTree, query.Id, destination)
return self.flatten(queryTree, query.Id, destination)
}
///////////////////////////////////////////////////////////////////////////////////////////////////
//MASK
///////////////////////////////////////////////////////////////////////////////////////////////////
type MaskQueryStep struct {
*BaseQueryStep
start, end uint64
}
type MaskQueryResult struct {
*BaseQueryResult
}
// QueryTree for Mask queries
type MaskQueryTree struct {
start, end uint64
bitmap *db.Bitmap
}
// Uses consistent hashing function to select node containing data for GET operation
/*
func (qt *MaskQueryTree) 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 {
log.Println("GetFragmenForBitmapFailed GetQueryTree", slice)
return nil, err
}
return fragment.GetLocation(), nil
}
*/
///////////////////////////////////////////////////////////////////////////////////////////////////
//Range
///////////////////////////////////////////////////////////////////////////////////////////////////
type RangeQueryStep struct {
*BaseQueryStep
bitmap *db.Bitmap
start, end time.Time
}
type RangeQueryResult struct {
*BaseQueryResult
}
// QueryTree for Mask queries
type RangeQueryTree struct {
}
// Uses consistent hashing function to select node containing data for GET operation
/*
func (self *RangeQueryTree) 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, self.bitmap)
if err != nil {
log.Println("GetFragmenForBitmapFailed GetQueryTree", slice)
return nil, err
}
return fragment.GetLocation(), nil
}
*/

View file

@ -108,3 +108,23 @@ func ParseGUID(input string) (GUID, error) {
}
return u, nil
}
func In(val int, list []int) bool {
for _, v := range list {
if val == v {
return true
}
}
return false
}
func Difference(a, b []int) []int {
results := make([]int, 0, len(a))
for _, v := range a {
if !In(v, b) {
results = append(results, v)
}
}
return results
}