From a0d6c253d1a6f875cacd62e91a9903c0fc77b159 Mon Sep 17 00:00:00 2001 From: Alan Bernstein Date: Wed, 13 Jan 2021 11:02:54 -0600 Subject: [PATCH] Pass SQL query string from mapper to tracker --- api.go | 2 +- handler.go | 3 +++ http/handler.go | 2 +- sql/mapper.go | 2 ++ sql/select.go | 8 +++++--- tracker.go | 26 +++++++++++++++----------- 6 files changed, 27 insertions(+), 16 deletions(-) diff --git a/api.go b/api.go index 9140153a1..2bd8eb2e4 100644 --- a/api.go +++ b/api.go @@ -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{ diff --git a/handler.go b/handler.go index d8e370d56..b50e6e946 100644 --- a/handler.go +++ b/handler.go @@ -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 diff --git a/http/handler.go b/http/handler.go index bc2fbb1ab..a6dfd5828 100644 --- a/http/handler.go +++ b/http/handler.go @@ -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 diff --git a/sql/mapper.go b/sql/mapper.go index a43dba06d..eb6c23f0c 100644 --- a/sql/mapper.go +++ b/sql/mapper.go @@ -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 } diff --git a/sql/select.go b/sql/select.go index b54eee0b4..7e1418b14 100644 --- a/sql/select.go +++ b/sql/select.go @@ -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") } diff --git a/tracker.go b/tracker.go index 359850005..bec0cf4db 100644 --- a/tracker.go +++ b/tracker.go @@ -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