mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-06 19:07:50 +00:00
added new logging technique
This commit is contained in:
parent
377d850534
commit
170fe96661
12 changed files with 85 additions and 5 deletions
|
|
@ -298,6 +298,7 @@ func (self *WebService) HandleQuery(w http.ResponseWriter, r *http.Request) {
|
|||
}
|
||||
_, bits := r.Form["bits"]
|
||||
|
||||
log.Info("PQL:", database_name, pql)
|
||||
results, err := self.service.Executor.RunPQL(database_name, pql)
|
||||
if err != nil {
|
||||
log.Warn("PQL Exec Error:", err.Error(), database_name, pql)
|
||||
|
|
|
|||
|
|
@ -54,7 +54,7 @@ func (self *RemoteSetBit) Request() {
|
|||
QueryId: random_id,
|
||||
DestProcessId: *process,
|
||||
}
|
||||
wait := len(request)
|
||||
wait := len(request) * 10
|
||||
if wait < 10 {
|
||||
wait = 10
|
||||
}
|
||||
|
|
|
|||
|
|
@ -71,13 +71,15 @@ func (self *Service) PrepareLogging() {
|
|||
log.SetOutput(f)
|
||||
*/
|
||||
fname := fmt.Sprintf("%s/%s.%s", base_path, self.name, self.Id)
|
||||
//<seelog minlevel="debug" maxlevel="error">
|
||||
log_level := config.GetStringDefault("log_level", "info")
|
||||
prod_config := fmt.Sprintf(`
|
||||
<seelog>
|
||||
<seelog minlevel="%s">
|
||||
<outputs>
|
||||
<rollingfile type="size" filename="%s" maxsize="524288000" maxrolls="4" />
|
||||
</outputs>
|
||||
</seelog>
|
||||
`, fname)
|
||||
`, log_level, fname)
|
||||
logger, _ := log.LoggerFromConfigAsBytes([]byte(prod_config))
|
||||
log.ReplaceLogger(logger)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -30,10 +30,12 @@ func (self *Dispatch) Run() {
|
|||
message := self.service.Transport.Receive()
|
||||
switch data := message.Data.(type) {
|
||||
case core.BatchRequest:
|
||||
log.Trace("Dispatch.Run BatchRequest")
|
||||
response := db.Message{Data: core.BatchResponse{Id: data.Id}}
|
||||
self.service.Index.LoadBitmap(data.Fragment_id, data.Bitmap_id, data.Compressed_bitmap, data.Filter)
|
||||
self.service.Transport.Send(&response, data.Source)
|
||||
case core.BitsRequest:
|
||||
log.Trace("Dispatch.Run BitsRequest")
|
||||
var results []core.SBResult
|
||||
result := false
|
||||
|
||||
|
|
@ -50,13 +52,17 @@ func (self *Dispatch) Run() {
|
|||
response := db.Message{Data: core.BitsResponse{Id: &data.QueryId, Items: results}}
|
||||
self.service.Transport.Send(&response, &data.ReturnProcessId)
|
||||
case core.PingRequest:
|
||||
log.Trace("Dispatch.Run Ping")
|
||||
pong := db.Message{Data: core.PongRequest{Id: data.Id}}
|
||||
self.service.Transport.Send(&pong, data.Source)
|
||||
case db.HoldResult:
|
||||
log.Trace("Dispatch.Run HoldResult")
|
||||
self.service.Hold.Set(data.ResultId(), data.ResultData(), 30)
|
||||
case query.PortableQueryStep:
|
||||
log.Trace("Dispatch.Run PortableQueryStep")
|
||||
go self.service.Executor.NewJob(message)
|
||||
case core.TopFill:
|
||||
log.Trace("Dispatch.Run TopFill")
|
||||
go self.service.TopFillHandler(message)
|
||||
case core.BitsResponse:
|
||||
self.service.Hold.Set(data.ResultId(), data.ResultData(), 30)
|
||||
|
|
|
|||
|
|
@ -17,15 +17,16 @@ type Executor struct {
|
|||
}
|
||||
|
||||
func (self *Executor) Init() error {
|
||||
log.Warn("Starting Executor")
|
||||
log.Trace("Executor.Init()")
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *Executor) Close() {
|
||||
log.Warn("Shutting down Executor")
|
||||
log.Trace("Executor.Close()")
|
||||
}
|
||||
|
||||
func (self *Executor) NewJob(job *db.Message) {
|
||||
log.Trace("NewJob", job)
|
||||
switch job.Data.(type) {
|
||||
case query.CountQueryStep:
|
||||
self.service.CountQueryStepHandler(job)
|
||||
|
|
@ -73,6 +74,7 @@ func (self *Executor) RunQueryTest(database_name string, pql string) string {
|
|||
}
|
||||
|
||||
func (self *Executor) runQuery(database *db.Database, qry *query.Query) error {
|
||||
log.Trace("Executor.runQuery", database, qry)
|
||||
process, err := self.service.GetProcess()
|
||||
if err != nil {
|
||||
return err
|
||||
|
|
@ -109,6 +111,7 @@ func (self *Executor) runQuery(database *db.Database, qry *query.Query) error {
|
|||
}
|
||||
|
||||
func (self *Executor) RunPQL(database_name string, pql string) (interface{}, error) {
|
||||
log.Trace("Executor.RunPQL", database_name, pql)
|
||||
database := self.service.Cluster.GetOrCreateDatabase(database_name)
|
||||
|
||||
// see if the outer query function is a custom query
|
||||
|
|
@ -190,5 +193,6 @@ func (self *Executor) Run() {
|
|||
}
|
||||
|
||||
func NewExecutor(service *core.Service) *Executor {
|
||||
log.Trace("NewExector")
|
||||
return &Executor{service, make(chan *db.Message)}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package hold
|
|||
|
||||
import (
|
||||
"errors"
|
||||
log "github.com/cihub/seelog"
|
||||
. "pilosa/util"
|
||||
"time"
|
||||
)
|
||||
|
|
@ -24,11 +25,13 @@ type Holder struct {
|
|||
//var Hold Holder
|
||||
|
||||
func (self *Holder) DelChan(id *GUID) {
|
||||
log.Trace("Holder.DelChan", id)
|
||||
req := delhold{id}
|
||||
self.delchan <- req
|
||||
}
|
||||
|
||||
func (self *Holder) GetChan(id *GUID) holdchan {
|
||||
log.Trace("Holder.GetChan", id)
|
||||
reply := make(chan holdchan)
|
||||
req := gethold{id, reply}
|
||||
self.getchan <- req
|
||||
|
|
@ -36,6 +39,7 @@ func (self *Holder) GetChan(id *GUID) holdchan {
|
|||
}
|
||||
|
||||
func (self *Holder) Get(id *GUID, timeout int) (interface{}, error) {
|
||||
log.Trace("Holder.Get", id, timeout)
|
||||
ch := self.GetChan(id)
|
||||
select {
|
||||
case val := <-ch:
|
||||
|
|
@ -47,6 +51,7 @@ func (self *Holder) Get(id *GUID, timeout int) (interface{}, error) {
|
|||
}
|
||||
|
||||
func (self *Holder) Set(id *GUID, value interface{}, timeout int) {
|
||||
log.Trace("Holder.Set", id, value, timeout)
|
||||
ch := self.GetChan(id)
|
||||
go func() {
|
||||
select {
|
||||
|
|
|
|||
|
|
@ -104,6 +104,7 @@ func (c *CassandraStorage) Fetch(bitmap_id uint64, db string, frame string, slic
|
|||
}
|
||||
delta := time.Since(start)
|
||||
util.SendTimer("cassandra_storage_Fetch", delta.Nanoseconds())
|
||||
util.SendInc("cassandra_storage_Read")
|
||||
bitmap.SetCount(uint64(count))
|
||||
return bitmap, uint64(filter)
|
||||
}
|
||||
|
|
@ -144,6 +145,7 @@ func (self *CassandraStorage) EndBatch() {
|
|||
}
|
||||
delta := time.Since(start)
|
||||
util.SendTimer("cassandra_storage_EndBatch", delta.Nanoseconds())
|
||||
util.SendInc("cassandra_storage_Write")
|
||||
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ package query
|
|||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
log "github.com/cihub/seelog"
|
||||
"strings"
|
||||
"unicode"
|
||||
"unicode/utf8"
|
||||
|
|
@ -44,6 +45,7 @@ type Lexer struct {
|
|||
}
|
||||
|
||||
func (lexer *Lexer) emit(typ int) {
|
||||
log.Trace("Lexer.emit", typ)
|
||||
if lexer.start < lexer.pos {
|
||||
lexer.ch <- Token{lexer.text[lexer.start:lexer.pos], typ}
|
||||
}
|
||||
|
|
@ -51,6 +53,7 @@ func (lexer *Lexer) emit(typ int) {
|
|||
}
|
||||
|
||||
func (lexer *Lexer) acceptUntil(chars string, consume bool) (rune, error) {
|
||||
log.Trace("Lexer.acceptUntil", chars, consume)
|
||||
start_pos := lexer.pos
|
||||
|
||||
for {
|
||||
|
|
@ -77,6 +80,7 @@ func (lexer *Lexer) acceptUntil(chars string, consume bool) (rune, error) {
|
|||
}
|
||||
|
||||
func (lexer *Lexer) acceptRun(valid string) {
|
||||
log.Trace("Lexer.acceptRun", valid)
|
||||
for strings.IndexRune(valid, lexer.next()) >= 0 {
|
||||
}
|
||||
lexer.backup()
|
||||
|
|
@ -84,6 +88,7 @@ func (lexer *Lexer) acceptRun(valid string) {
|
|||
|
||||
// next returns the next rune in the input.
|
||||
func (lexer *Lexer) next() (runey rune) {
|
||||
log.Trace("Lexer.next", runey)
|
||||
if lexer.pos >= len(lexer.text) {
|
||||
lexer.width = 0
|
||||
return 0
|
||||
|
|
@ -95,18 +100,21 @@ func (lexer *Lexer) next() (runey rune) {
|
|||
|
||||
// ignore skips over the pending input before this point.
|
||||
func (lexer *Lexer) ignore() {
|
||||
log.Trace("Lexer.ignore")
|
||||
lexer.start = lexer.pos
|
||||
}
|
||||
|
||||
// backup steps back one rune.
|
||||
// Can be called only once per call of next.
|
||||
func (lexer *Lexer) backup() {
|
||||
log.Trace("Lexer.backup")
|
||||
lexer.pos -= lexer.width
|
||||
}
|
||||
|
||||
// peek returns but does not consume
|
||||
// the next rune in the input.
|
||||
func (lexer *Lexer) peek() rune {
|
||||
log.Trace("Lexer.peek")
|
||||
for {
|
||||
next_rune := lexer.next()
|
||||
// ignore spaces
|
||||
|
|
@ -119,6 +127,7 @@ func (lexer *Lexer) peek() rune {
|
|||
}
|
||||
|
||||
func stateError(err error) func(lexer *Lexer) statefn {
|
||||
log.Trace("stateError", err)
|
||||
return func(lexer *Lexer) statefn {
|
||||
lexer.ch <- Token{err.Error(), TYPE_ERROR}
|
||||
close(lexer.ch)
|
||||
|
|
@ -127,6 +136,7 @@ func stateError(err error) func(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateFunc(lexer *Lexer) statefn {
|
||||
log.Trace("stateFunc", lexer)
|
||||
_, err := lexer.acceptUntil("(", true)
|
||||
if err != nil {
|
||||
return stateError(err)
|
||||
|
|
@ -136,6 +146,7 @@ func stateFunc(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateLP(lexer *Lexer) statefn {
|
||||
log.Trace("stateLP", lexer)
|
||||
lexer.pos += 1
|
||||
lexer.emit(TYPE_LP)
|
||||
// handle multiple LPs
|
||||
|
|
@ -146,6 +157,7 @@ func stateLP(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateLB(lexer *Lexer) statefn {
|
||||
log.Trace("stateLB", lexer)
|
||||
lexer.acceptUntil("[", true)
|
||||
lexer.next()
|
||||
lexer.emit(TYPE_LB)
|
||||
|
|
@ -177,6 +189,7 @@ func stateLB(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateRB(lexer *Lexer) statefn {
|
||||
log.Trace("stateRB", lexer)
|
||||
lexer.pos += 1
|
||||
lexer.emit(TYPE_RB)
|
||||
|
||||
|
|
@ -191,6 +204,7 @@ func stateRB(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateArgs(lexer *Lexer) statefn {
|
||||
log.Trace("stateArgs", lexer)
|
||||
r, err := lexer.acceptUntil("(),=[]", false)
|
||||
if err != nil {
|
||||
return stateError(err)
|
||||
|
|
@ -215,6 +229,7 @@ func stateArgs(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateKeyword(lexer *Lexer) statefn {
|
||||
log.Trace("stateKeyword", lexer)
|
||||
_, err := lexer.acceptUntil("=", true)
|
||||
if err != nil {
|
||||
return stateError(err)
|
||||
|
|
@ -224,6 +239,7 @@ func stateKeyword(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateEquals(lexer *Lexer) statefn {
|
||||
log.Trace("stateEquals", lexer)
|
||||
e := lexer.next()
|
||||
if e != '=' {
|
||||
return stateError(errors.New("Expecting '='!"))
|
||||
|
|
@ -233,6 +249,7 @@ func stateEquals(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateValue(lexer *Lexer) statefn {
|
||||
log.Trace("stateValue", lexer)
|
||||
r, err := lexer.acceptUntil("(),[", false)
|
||||
if err != nil {
|
||||
return stateError(err)
|
||||
|
|
@ -259,6 +276,7 @@ func stateValue(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateRP(lexer *Lexer) statefn {
|
||||
log.Trace("stateRP", lexer)
|
||||
lexer.pos += 1
|
||||
lexer.emit(TYPE_RP)
|
||||
|
||||
|
|
@ -275,23 +293,27 @@ func stateRP(lexer *Lexer) statefn {
|
|||
}
|
||||
|
||||
func stateRPComma(lexer *Lexer) statefn {
|
||||
log.Trace("stateRPComma", lexer)
|
||||
lexer.pos += 1
|
||||
lexer.emit(TYPE_COMMA)
|
||||
return stateValue
|
||||
}
|
||||
|
||||
func stateComma(lexer *Lexer) statefn {
|
||||
log.Trace("stateComma", lexer)
|
||||
lexer.pos += 1
|
||||
lexer.emit(TYPE_COMMA)
|
||||
return stateArgs
|
||||
}
|
||||
|
||||
func stateEOF(lexer *Lexer) statefn {
|
||||
log.Trace("stateEOF", lexer)
|
||||
close(lexer.ch)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (lexer *Lexer) Lex() (tokens []Token, err error) {
|
||||
log.Trace("Lexer.Lex", tokens, err)
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
var ok bool
|
||||
|
|
@ -321,6 +343,7 @@ func (lexer *Lexer) Lex() (tokens []Token, err error) {
|
|||
}
|
||||
|
||||
func Lex(input string) ([]Token, error) {
|
||||
log.Trace("Lex", input)
|
||||
lexer := Lexer{input, 0, 0, 0, TYPE_FUNC, make(chan Token)}
|
||||
return lexer.Lex()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -40,6 +40,7 @@ bm := db.Bitmap{bitmap_id, frame_type, filter}
|
|||
return []QueryInput{&bm}, uint64(profile_id), 0
|
||||
*/
|
||||
func (self *QueryParser) next() *Token {
|
||||
log.Trace("QueryParser.next")
|
||||
self.pos += 1
|
||||
if self.pos > len(self.tokens) {
|
||||
return nil
|
||||
|
|
@ -48,16 +49,19 @@ func (self *QueryParser) next() *Token {
|
|||
}
|
||||
|
||||
func (self *QueryParser) peek() *Token {
|
||||
log.Trace("QueryParser.peek")
|
||||
token := self.next()
|
||||
self.backup()
|
||||
return token
|
||||
}
|
||||
|
||||
func (self *QueryParser) backup() {
|
||||
log.Trace("QueryParser.backup")
|
||||
self.pos -= 1
|
||||
}
|
||||
|
||||
func (self *QueryParser) Parse() (query *Query, err error) {
|
||||
log.Trace("QueryParser.Parse", query, err)
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
var ok bool
|
||||
|
|
@ -341,6 +345,7 @@ ArgLoop:
|
|||
}
|
||||
|
||||
func Parse(tokens []Token) (*Query, error) {
|
||||
log.Trace("Parse", tokens)
|
||||
parser := QueryParser{tokens, 0}
|
||||
return parser.Parse()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -66,17 +66,22 @@ type BaseQueryStep struct {
|
|||
}
|
||||
|
||||
func (self *BaseQueryStep) GetId() *util.GUID {
|
||||
log.Trace("BaseQueryStep", *self.Id)
|
||||
return self.Id
|
||||
}
|
||||
func (self *BaseQueryStep) GetLocation() *db.Location {
|
||||
log.Trace("BaseQueryStep.GetLocation", *self.Location)
|
||||
return self.Location
|
||||
}
|
||||
|
||||
func (self *BaseQueryStep) LocIsDest() bool {
|
||||
log.Trace("BaseQueryStep.LocIsDest")
|
||||
if self.Location.ProcessId == self.Destination.ProcessId &&
|
||||
self.Location.FragmentId == self.Destination.FragmentId {
|
||||
log.Trace("BaseQueryStep.LocIsDest Return true")
|
||||
return true
|
||||
}
|
||||
log.Trace("BaseQueryStep.LocIsDest Return false")
|
||||
return false
|
||||
}
|
||||
|
||||
|
|
@ -86,10 +91,13 @@ type BaseQueryResult struct {
|
|||
}
|
||||
|
||||
func (self *BaseQueryResult) ResultId() *util.GUID {
|
||||
log.Trace("BaseQueryStep.ResultId", self)
|
||||
|
||||
return self.Id
|
||||
}
|
||||
|
||||
func (self *BaseQueryResult) ResultData() interface{} {
|
||||
log.Trace("BaseQueryStep.ResultData", self)
|
||||
return self.Data
|
||||
}
|
||||
|
||||
|
|
@ -113,6 +121,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, error) {
|
||||
log.Trace("CountQueryTree.getLocation", d, qt)
|
||||
return qt.subquery.getLocation(d)
|
||||
}
|
||||
|
||||
|
|
@ -143,6 +152,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, error) {
|
||||
log.Trace("TopNQueryTree.getLocation", d, qt)
|
||||
var err error
|
||||
if qt.location == nil {
|
||||
frame := d.GetOrCreateFrame(qt.Frame)
|
||||
|
|
@ -178,6 +188,7 @@ type UnionQueryTree struct {
|
|||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *UnionQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("UnionQueryTree.getLocation", d, qt)
|
||||
var err error
|
||||
if qt.location == nil {
|
||||
subqueryLength := len(qt.subqueries)
|
||||
|
|
@ -210,6 +221,7 @@ type IntersectQueryTree struct {
|
|||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *IntersectQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("IntersectQueryTree.getLocation", d, qt)
|
||||
var err error
|
||||
if qt.location == nil {
|
||||
subqueryLength := len(qt.subqueries)
|
||||
|
|
@ -242,6 +254,7 @@ type DifferenceQueryTree struct {
|
|||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *DifferenceQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("DifferenceQueryTree.getLocation", d, qt)
|
||||
var err error
|
||||
if qt.location == nil {
|
||||
subqueryLength := len(qt.subqueries)
|
||||
|
|
@ -278,11 +291,13 @@ type CatQueryTree struct {
|
|||
}
|
||||
|
||||
func (qt *CatQueryTree) Append(subtree QueryTree) {
|
||||
log.Trace("CatQueryTree.Append", qt, subtree)
|
||||
qt.subqueries = append(qt.subqueries, subtree)
|
||||
}
|
||||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *CatQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("CatQueryTree.getLocation", d)
|
||||
var err error
|
||||
if qt.location == nil {
|
||||
subqueryLength := len(qt.subqueries)
|
||||
|
|
@ -305,6 +320,7 @@ type AllQueryTree struct {
|
|||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *AllQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("AllQueryTree.getLocation", d)
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
|
|
@ -329,6 +345,7 @@ type GetQueryTree struct {
|
|||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *GetQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("GetQueryTree.getLocation", d)
|
||||
slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create)
|
||||
fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap)
|
||||
if err != nil {
|
||||
|
|
@ -359,6 +376,7 @@ type SetQueryTree struct {
|
|||
|
||||
// Uses consistent hashing function to select node containing data for GET operation
|
||||
func (qt *SetQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("SetQueryTree.getLocation", d)
|
||||
// check here for supported frames
|
||||
if !d.IsValidFrame(qt.bitmap.FrameType) {
|
||||
return nil, NewInvalidFrame(d.Name, qt.bitmap.FrameType)
|
||||
|
|
@ -420,6 +438,7 @@ type QueryTree interface {
|
|||
}
|
||||
|
||||
func validateRange(Args map[string]interface{}) error {
|
||||
log.Trace("validateRange", Args)
|
||||
_, ok := Args["id"]
|
||||
if !ok {
|
||||
return errors.New("missing bitmap id")
|
||||
|
|
@ -441,6 +460,7 @@ func validateRange(Args map[string]interface{}) error {
|
|||
|
||||
// Builds QueryTree object from Query. Pass slice=-1 to perform operation on all slices
|
||||
func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error) {
|
||||
log.Trace("QueryPlanner.buildTree", query, slice)
|
||||
var tree QueryTree
|
||||
|
||||
// handle SET operation regardless of the slice
|
||||
|
|
@ -583,6 +603,7 @@ func (self *QueryPlanner) buildTree(query *Query, slice int) (QueryTree, error)
|
|||
|
||||
// Produces flattened QueryPlan from QueryTree input
|
||||
func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Location) (*QueryPlan, error) {
|
||||
log.Trace("QueryPlanner.flatten", self, qt, id, location)
|
||||
plan := QueryPlan{}
|
||||
if cat, ok := qt.(*CatQueryTree); ok {
|
||||
inputs := make([]*util.GUID, len(cat.subqueries))
|
||||
|
|
@ -750,6 +771,7 @@ func (self *QueryPlanner) flatten(qt QueryTree, id *util.GUID, location *db.Loca
|
|||
|
||||
// Transforms Query into QueryTree and flattens to QueryPlan object
|
||||
func (self *QueryPlanner) Plan(query *Query, id *util.GUID, destination *db.Location) (*QueryPlan, error) {
|
||||
log.Trace("QueryPlanner.Plan", self, query, id, destination)
|
||||
queryTree, err := self.buildTree(query, -1)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -809,6 +831,7 @@ type RangeQueryTree struct {
|
|||
}
|
||||
|
||||
func (qt *RangeQueryTree) getLocation(d *db.Database) (*db.Location, error) {
|
||||
log.Trace("RangeQueryTree", qt, d)
|
||||
slice := d.GetOrCreateSlice(qt.slice) // TODO: this should probably be just GetSlice (no create)
|
||||
fragment, err := d.GetFragmentForBitmap(slice, qt.bitmap)
|
||||
if err != nil {
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package query
|
||||
|
||||
import (
|
||||
log "github.com/cihub/seelog"
|
||||
"pilosa/db"
|
||||
"pilosa/util"
|
||||
"strings"
|
||||
|
|
@ -28,6 +29,7 @@ type Query struct {
|
|||
}
|
||||
|
||||
func QueryPlanForPQL(database *db.Database, pql string, destination *db.Location) (*QueryPlan, error) {
|
||||
log.Trace("QueryPlanFOrPQL", database, pql, destination)
|
||||
tokens, err := Lex(pql)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -36,6 +38,7 @@ func QueryPlanForPQL(database *db.Database, pql string, destination *db.Location
|
|||
}
|
||||
|
||||
func QueryForPQL(pql string) (*Query, error) {
|
||||
log.Trace("QueryForPQL", pql)
|
||||
tokens, err := Lex(pql)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -44,6 +47,7 @@ func QueryForPQL(pql string) (*Query, error) {
|
|||
}
|
||||
|
||||
func QueryForTokens(tokens []Token) (*Query, error) {
|
||||
log.Trace("QueryForTokens", tokens)
|
||||
query, err := Parse(tokens)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -52,6 +56,7 @@ func QueryForTokens(tokens []Token) (*Query, error) {
|
|||
}
|
||||
|
||||
func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.Location) (*QueryPlan, error) {
|
||||
log.Trace("QueryPlanForTokens", database, tokens, destination)
|
||||
query, err := QueryForTokens(tokens)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -60,6 +65,7 @@ func QueryPlanForTokens(database *db.Database, tokens []Token, destination *db.L
|
|||
}
|
||||
|
||||
func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Location) (*QueryPlan, error) {
|
||||
log.Trace("QueryPlanForQuery", database, query, destination)
|
||||
query_planner := QueryPlanner{Database: database, Query: query}
|
||||
id := util.RandomUUID()
|
||||
query_plan, err := query_planner.Plan(query, &id, destination)
|
||||
|
|
@ -70,6 +76,7 @@ func QueryPlanForQuery(database *db.Database, query *Query, destination *db.Loca
|
|||
}
|
||||
|
||||
func TokensToFilterStrings(tokens []Token) (string, []string) {
|
||||
log.Trace("TokensToFilterStrings", tokens)
|
||||
var whole []string
|
||||
var filter string
|
||||
var filters []string
|
||||
|
|
|
|||
|
|
@ -153,12 +153,14 @@ func (self *TcpTransport) Close() {
|
|||
}
|
||||
|
||||
func (self *TcpTransport) Send(message *db.Message, host *GUID) {
|
||||
log.Trace("TcpTransport.Send", message, host)
|
||||
envelope := db.Envelope{message, host}
|
||||
notify.Post("outbox", &envelope)
|
||||
self.outbox <- envelope
|
||||
}
|
||||
|
||||
func (self *TcpTransport) Receive() *db.Message {
|
||||
log.Trace("TcpTransport.Receive")
|
||||
message := <-self.inbox
|
||||
notify.Post("inbox", message)
|
||||
return message
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue