Pass SQL query string from mapper to tracker

This commit is contained in:
Alan Bernstein 2021-01-13 11:02:54 -06:00
parent 06241d1f84
commit a0d6c253d1
6 changed files with 27 additions and 16 deletions

2
api.go
View file

@ -161,7 +161,7 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er
}
if !req.Remote {
defer api.tracker.Finish(api.tracker.Start(req.Query, api.server.nodeID, req.Index, start))
defer api.tracker.Finish(api.tracker.Start(req.Query, req.SQLQuery, api.server.nodeID, req.Index, start))
}
// TODO can we get rid of exec options and pass the QueryRequest directly to executor?
execOpts := &execOptions{

View file

@ -29,6 +29,9 @@ type QueryRequest struct {
// The query string to parse and execute.
Query string
// The SQL source query, if applicable.
SQLQuery string
// The shards to include in the query execution.
// If empty, all shards are included.
Shards []uint64

View file

@ -1182,7 +1182,7 @@ func (h *Handler) handleGetActiveQueries(w http.ResponseWriter, r *http.Request)
}
}
for i, q := range queries {
_, err := fmt.Fprintf(w, "%*s%q\n", -(maxlen + 2), durations[i], q.Query)
_, err := fmt.Fprintf(w, "%*s%q\n", -(maxlen + 2), durations[i], q.PQL)
if err != nil {
h.logger.Printf("sending GetActiveQueries response: %s", err)
return

View file

@ -40,6 +40,7 @@ type MappedSQL struct {
Statement sqlparser.Statement
Mask QueryMask
Tables []string
SQL string
}
// Mapper is responsible for mapping a SQL query to structure representation
@ -109,5 +110,6 @@ func (m *Mapper) MapSQL(sql string) (*MappedSQL, error) {
Statement: stmt,
Mask: qm,
Tables: tableNames,
SQL: sql,
}, nil
}

View file

@ -43,6 +43,7 @@ func NewSelectHandler(api *pilosa.API) *SelectHandler {
// Handle executes mapped SQL
func (s *SelectHandler) Handle(ctx context.Context, mapped *MappedSQL) (pproto.ToRowser, error) {
stmt, ok := mapped.Statement.(*sqlparser.Select)
if !ok {
return nil, fmt.Errorf("statement is not type select: %T", mapped.Statement)
}
@ -50,7 +51,7 @@ func (s *SelectHandler) Handle(ctx context.Context, mapped *MappedSQL) (pproto.T
if err != nil {
return nil, errors.Wrap(err, "mapping select")
}
return s.execMappingResult(ctx, mr)
return s.execMappingResult(ctx, mr, mapped.SQL)
}
func (s *SelectHandler) mapSelect(ctx context.Context, selectStmt *sqlparser.Select, qm QueryMask) (*MappingResult, error) {
@ -74,12 +75,13 @@ func (s *SelectHandler) mapSelect(ctx context.Context, selectStmt *sqlparser.Sel
return mr, nil
}
func (s *SelectHandler) execMappingResult(ctx context.Context, mr *MappingResult) (pproto.ToRowser, error) {
func (s *SelectHandler) execMappingResult(ctx context.Context, mr *MappingResult, sql string) (pproto.ToRowser, error) {
if mr.Query == "" {
return nil, errors.New("no pql query created")
}
fmt.Printf("execMappingResult: %+v\n", sql)
resp, err := s.api.Query(ctx, &pilosa.QueryRequest{Index: mr.IndexName, Query: mr.Query})
resp, err := s.api.Query(ctx, &pilosa.QueryRequest{Index: mr.IndexName, Query: mr.Query, SQLQuery: sql})
if err != nil {
return nil, errors.Wrap(err, "doing pql query")
}

View file

@ -21,14 +21,16 @@ import (
)
type ActiveQueryStatus struct {
Query string `json:"query"`
PQL string `json:"PQL"`
SQL string `json:"SQL,omitempty"`
Node string `json:"node"`
Index string `json:"index"`
Age time.Duration `json:"age"`
}
type PastQueryStatus struct {
Query string `json:"query"`
PQL string `json:"PQL"`
SQL string `json:"SQL,omitempty"`
Node string `json:"nodeID"`
Index string `json:"index"`
Start time.Time `json:"start"`
@ -36,14 +38,16 @@ type PastQueryStatus struct {
}
type activeQuery struct {
query string
PQL string
SQL string
node string
index string
started time.Time
}
type pastQuery struct {
query string
PQL string
SQL string
node string
index string
started time.Time
@ -123,7 +127,7 @@ func newQueryTracker(historyLength int) *queryTracker {
select {
case update := <-updates:
if update.end {
pq := pastQuery{update.q.query, update.q.node, update.q.index, update.q.started, update.endTime.Sub(update.q.started)}
pq := pastQuery{update.q.PQL, update.q.SQL, update.q.node, update.q.index, update.q.started, update.endTime.Sub(update.q.started)}
tracker.history.add(pq)
delete(activeQueries, update.q)
} else {
@ -146,8 +150,8 @@ func newQueryTracker(historyLength int) *queryTracker {
return tracker
}
func (t *queryTracker) Start(query, nodeID, index string, start time.Time) *activeQuery {
q := &activeQuery{query, nodeID, index, start}
func (t *queryTracker) Start(pql, sql, nodeID, index string, start time.Time) *activeQuery {
q := &activeQuery{pql, sql, nodeID, index, start}
t.updates <- queryStatusUpdate{q, false, time.Time{}}
return q
}
@ -166,9 +170,9 @@ func (t *queryTracker) ActiveQueries() []ActiveQueryStatus {
return true
case queries[i].started.After(queries[j].started):
return false
case queries[i].query < queries[j].query:
case queries[i].PQL < queries[j].PQL:
return true
case queries[i].query > queries[j].query:
case queries[i].PQL > queries[j].PQL:
return false
default:
return false
@ -177,7 +181,7 @@ func (t *queryTracker) ActiveQueries() []ActiveQueryStatus {
now := time.Now()
out := make([]ActiveQueryStatus, len(queries))
for i, v := range queries {
out[i] = ActiveQueryStatus{v.query, v.node, v.index, now.Sub(v.started)}
out[i] = ActiveQueryStatus{v.PQL, v.SQL, v.node, v.index, now.Sub(v.started)}
}
return out
}
@ -186,7 +190,7 @@ func (t *queryTracker) PastQueries() []PastQueryStatus {
queries := t.history.slice()
out := make([]PastQueryStatus, len(queries))
for i, v := range queries {
out[i] = PastQueryStatus{v.query, v.node, v.index, v.started, v.runtime}
out[i] = PastQueryStatus{v.PQL, v.SQL, v.node, v.index, v.started, v.runtime}
}
return out