Merge pull request #1327 from alanbernstein/sql-history

Include SQL string in query-history
This commit is contained in:
Alan Bernstein 2021-01-20 19:05:36 -06:00 • committed by GitHub
commit b273f3ba60
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
8 changed files with 53 additions and 26 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

@ -37,6 +37,7 @@ import (
"github.com/pilosa/pilosa/v2/encoding/proto"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/pql"
pb "github.com/pilosa/pilosa/v2/proto"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
)
@ -1502,6 +1503,16 @@ func TestQueryHistory(t *testing.T) {
test.Do(t, "POST", cmd.URL()+"/index/i0", "")
test.Do(t, "POST", cmd.URL()+"/index/i0/field/f0", "")
gh := server.NewGRPCHandler(cmd.API)
_, err = gh.QuerySQLUnary(context.Background(), &pb.QuerySQLRequest{
Sql: `select * from i0`,
})
if err != nil {
t.Fatalf("QuerySQLUnary failed: %v", err)
}
test.Do(t, "POST", cmd.URL()+"/index/i0/query", "Set(0, f0=0)")
test.Do(t, "POST", cmd.URL()+"/index/i0/query", "Set(3000000, f0=0)")
test.Do(t, "POST", cmd.URL()+"/index/i0/query", "TopN(f0)")
@ -1511,7 +1522,7 @@ func TestQueryHistory(t *testing.T) {
t.Fatalf("unexpected status code: %d", w.Code)
}
ret := make([]pilosa.PastQueryStatus, 3)
ret := make([]pilosa.PastQueryStatus, 4)
b, err := ioutil.ReadAll(w.Body)
if err != nil {
t.Fatalf("reading: %v", err)
@ -1522,10 +1533,10 @@ func TestQueryHistory(t *testing.T) {
}
// verify result length
if len(ret) != 3 {
if len(ret) != 4 {
// each set query executes on both nodes once
// topn query gets added to history on node0 once, node1 twice
t.Fatalf("expected list of length 3, got %d", len(ret))
t.Fatalf("expected list of length 4, got %d\n%+v", len(ret), ret)
}
// verify sort order
@ -1543,8 +1554,14 @@ func TestQueryHistory(t *testing.T) {
if ret[0].Node != cluster.GetNode(0).Server.NodeID() {
t.Fatalf("response value for 'Node' was '%s', expected '%s'", ret[0].Node, cluster.GetNode(0).Server.NodeID())
}
if ret[0].Query != "TopN(f0)" {
t.Fatalf("response value for 'Query' was '%s', expected 'TopN(f0)'", ret[0].Query)
if ret[3].PQL != "Extract(All(),Rows(f0))" {
t.Fatalf("response value for 'PQL' was '%s', expected 'Extract(All(),Rows(f0))'", ret[0].PQL)
}
if ret[3].SQL != "select * from i0" {
t.Fatalf("response value for 'SQL' was '%s', expected 'select * from i0'", ret[0].SQL)
}
if ret[0].PQL != "TopN(f0)" {
t.Fatalf("response value for 'PQL' was '%s', expected 'TopN(f0)'", ret[0].PQL)
}
}

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,12 @@ 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")
}
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

View file

@ -47,11 +47,11 @@ func TestRingBuffer(t *testing.T) {
}
for n, q := range buffer.slice() {
if q.query != tests[k].queries[n] {
t.Fatalf("test[%d], buffer[%d] expected querystring '%s', found '%s'", k, n, tests[k].queries[n], q.query)
if q.PQL != tests[k].queries[n] {
t.Fatalf("test[%d], buffer[%d] expected querystring '%s', found '%s'", k, n, tests[k].queries[n], q.PQL)
}
}
buffer.add(pastQuery{query: fmt.Sprintf("%d", k)})
buffer.add(pastQuery{PQL: fmt.Sprintf("%d", k)})
}
}
@ -63,13 +63,13 @@ func TestQueryTracker(t *testing.T) {
t.Fatalf("expected no active queries; found %v", queries)
}
qs := tracker.Start("test query", "node0", "i", time.Now())
qs := tracker.Start("test query", "test SQL", "node0", "i", time.Now())
var queries []ActiveQueryStatus
for len(queries) < 1 {
queries = tracker.ActiveQueries()
}
if len(queries) > 1 || queries[0].Query != "test query" {
if len(queries) > 1 || queries[0].PQL != "test query" {
t.Fatalf("unexpected queries: %v", queries)
}