adding nodes arg to options

This commit is contained in:
jacob 2023-03-22 05:21:29 +00:00
parent be4f365eaf
commit 5383daf026
4 changed files with 174 additions and 101 deletions

View file

@ -146,7 +146,7 @@ func (e *executor) executeApply(ctx context.Context, qcx *Qcx, index string, c *
reduceFn, tablerFn = IvyReduce(ivyReduce, ",", opt)
}
_, err = e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
_, err = e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nil)
if err != nil {
return nil, err
}

View file

@ -88,7 +88,7 @@ func (e *executor) executeArrow(ctx context.Context, qcx *Qcx, index string, c *
return nil
}
_, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
_, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nil)
if err != nil {
return nil, err
}

View file

@ -245,7 +245,7 @@ func (e *executor) Execute(ctx context.Context, tableKeyer dax.TableKeyer, q *pq
}
defer qcx.Abort()
results, err := e.execute(ctx, qcx, index, q, shards, opt)
results, err := e.execute(ctx, qcx, index, q, shards, opt, nil)
if err != nil {
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
@ -361,7 +361,7 @@ func safeCopy(resp QueryResponse) (out QueryResponse) {
// handlePreCalls traverses the call tree looking for calls that need
// precomputed values (e.g. Distinct, UnionRows, ConstRow...).
func (e *executor) handlePreCalls(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) error {
func (e *executor) handlePreCalls(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) error {
if c.Name == "Precomputed" {
idx := c.Args["valueidx"].(int64)
if idx >= 0 && idx < int64(len(opt.EmbeddedData)) {
@ -401,7 +401,7 @@ func (e *executor) handlePreCalls(ctx context.Context, qcx *Qcx, index string, c
// we need to recompute shards, then
shards = nil
}
if err := e.handlePreCallChildren(ctx, qcx, index, c, shards, opt); err != nil {
if err := e.handlePreCallChildren(ctx, qcx, index, c, shards, opt, nodes); err != nil {
return err
}
// child calls already handled, no precall for this, so we're done
@ -417,7 +417,7 @@ func (e *executor) handlePreCalls(ctx context.Context, qcx *Qcx, index string, c
// We set c to look like a normal call, and actually execute it:
c.Type = pql.PrecallNone
// possibly override call index.
v, err := e.executeCall(ctx, qcx, index, c, shards, opt)
v, err := e.executeCall(ctx, qcx, index, c, shards, opt, nodes)
if err != nil {
return err
}
@ -460,12 +460,12 @@ func (e *executor) dumpPrecomputedCalls(ctx context.Context, c *pql.Call) {
}
// handlePreCallChildren handles any pre-calls in the children of a given call.
func (e *executor) handlePreCallChildren(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) error {
func (e *executor) handlePreCallChildren(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) error {
for i := range c.Children {
if err := ctx.Err(); err != nil {
return err
}
if err := e.handlePreCalls(ctx, qcx, index, c.Children[i], shards, opt); err != nil {
if err := e.handlePreCalls(ctx, qcx, index, c.Children[i], shards, opt, nodes); err != nil {
return err
}
}
@ -479,7 +479,7 @@ func (e *executor) handlePreCallChildren(ctx context.Context, qcx *Qcx, index st
if err := ctx.Err(); err != nil {
return err
}
if err := e.handlePreCalls(ctx, qcx, index, call, shards, opt); err != nil {
if err := e.handlePreCalls(ctx, qcx, index, call, shards, opt, nodes); err != nil {
return err
}
}
@ -487,7 +487,7 @@ func (e *executor) handlePreCallChildren(ctx context.Context, qcx *Qcx, index st
return nil
}
func (e *executor) execute(ctx context.Context, qcx *Qcx, index string, q *pql.Query, shards []uint64, opt *ExecOptions) ([]interface{}, error) {
func (e *executor) execute(ctx context.Context, qcx *Qcx, index string, q *pql.Query, shards []uint64, opt *ExecOptions, nodes []string) ([]interface{}, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.execute")
defer span.Finish()
@ -569,17 +569,35 @@ func (e *executor) execute(ctx context.Context, qcx *Qcx, index string, q *pql.Q
if call.Name == "Count" {
// Handle count specially, skipping the level directly underneath it.
for _, child := range call.Children {
err := e.handlePreCallChildren(ctx, qcx, index, child, shards, opt)
err := e.handlePreCallChildren(ctx, qcx, index, child, shards, opt, nodes)
if err != nil {
return nil, err
}
}
} else {
err := e.handlePreCallChildren(ctx, qcx, index, call, shards, opt)
// In order to respect shard and node limitations of the Options
// call, we need to set them before Pre Calling Children
// in addition to in executeOptionsCall()
if call.Name == "Options" {
newShards, err := e.shardsFromOptionsCall(call, shards)
if err != nil {
return nil, err
}
shards = newShards
newNodes, err := e.nodesFromOptionsCall(call)
if err != nil {
return nil, err
}
nodes = newNodes
}
err := e.handlePreCallChildren(ctx, qcx, index, call, shards, opt, nodes)
if err != nil {
return nil, err
}
}
var v interface{}
var err error
// Top-level calls don't need to precompute cross-index things,
@ -589,9 +607,9 @@ func (e *executor) execute(ctx context.Context, qcx *Qcx, index string, q *pql.Q
// we don't need this logic in executeCall.
newIndex := call.CallIndex()
if newIndex != "" && newIndex != index {
v, err = e.executeCall(ctx, qcx, newIndex, call, nil, opt)
v, err = e.executeCall(ctx, qcx, newIndex, call, nil, opt, nodes)
} else {
v, err = e.executeCall(ctx, qcx, index, call, shards, opt)
v, err = e.executeCall(ctx, qcx, index, call, shards, opt, nodes)
}
if err != nil {
return nil, err
@ -676,7 +694,7 @@ func (e *executor) preprocessQuery(ctx context.Context, qcx *Qcx, index string,
}
// executeCall executes a call.
func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) {
func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (interface{}, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeCall")
defer span.Finish()
@ -722,23 +740,23 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
switch c.Name {
case "Sum":
statFn(CounterQuerySumTotal)
res, err := e.executeSum(ctx, qcx, index, c, shards, opt)
res, err := e.executeSum(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeSum")
case "Min":
statFn(CounterQueryMinTotal)
res, err := e.executeMin(ctx, qcx, index, c, shards, opt)
res, err := e.executeMin(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeMin")
case "Max":
statFn(CounterQueryMaxTotal)
res, err := e.executeMax(ctx, qcx, index, c, shards, opt)
res, err := e.executeMax(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeMax")
case "MinRow":
statFn(CounterQueryMinRowTotal)
res, err := e.executeMinRow(ctx, qcx, index, c, shards, opt)
res, err := e.executeMinRow(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeMinRow")
case "MaxRow":
statFn(CounterQueryMaxRowTotal)
res, err := e.executeMaxRow(ctx, qcx, index, c, shards, opt)
res, err := e.executeMaxRow(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeMaxRow")
case "Clear":
statFn(CounterQueryClearTotal)
@ -746,19 +764,19 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
return res, errors.Wrap(err, "executeClearBit")
case "ClearRow":
statFn(CounterQueryClearRowTotal)
res, err := e.executeClearRow(ctx, qcx, index, c, shards, opt)
res, err := e.executeClearRow(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeClearRow")
case "Distinct":
statFn(CounterQueryDistinctTotal)
res, err := e.executeDistinct(ctx, qcx, index, c, shards, opt)
res, err := e.executeDistinct(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeDistinct")
case "Store":
statFn(CounterQueryStoreTotal)
res, err := e.executeSetRow(ctx, qcx, index, c, shards, opt)
res, err := e.executeSetRow(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeSetRow")
case "Count":
statFn(CounterQueryCountTotal)
res, err := e.executeCount(ctx, qcx, index, c, shards, opt)
res, err := e.executeCount(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeCount")
case "Set":
statFn(CounterQuerySetTotal)
@ -766,15 +784,15 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
return res, errors.Wrap(err, "executeSet")
case "TopK":
statFn(CounterQueryTopKTotal)
res, err := e.executeTopK(ctx, qcx, index, c, shards, opt)
res, err := e.executeTopK(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeTopK")
case "TopN":
statFn(CounterQueryTopNTotal)
res, err := e.executeTopN(ctx, qcx, index, c, shards, opt)
res, err := e.executeTopN(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeTopN")
case "Rows":
statFn(CounterQueryRowsTotal)
res, err := e.executeRows(ctx, qcx, index, c, shards, opt)
res, err := e.executeRows(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeRows")
case "ExternalLookup":
statFn(CounterQueryExternalLookupTotal)
@ -782,11 +800,11 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
return res, errors.Wrap(err, "executeExternalLookup")
case "Extract":
statFn(CounterQueryExtractTotal)
res, err := e.executeExtract(ctx, qcx, index, c, shards, opt)
res, err := e.executeExtract(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeExtract")
case "GroupBy":
statFn(CounterQueryGroupByTotal)
res, err := e.executeGroupBy(ctx, qcx, index, c, shards, opt)
res, err := e.executeGroupBy(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeGroupBy")
case "Options":
statFn(CounterQueryOptionsTotal)
@ -794,11 +812,11 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
return res, errors.Wrap(err, "executeOptionsCall")
case "IncludesColumn":
statFn(CounterQueryIncludesColumnTotal)
res, err := e.executeIncludesColumnCall(ctx, qcx, index, c, shards, opt)
res, err := e.executeIncludesColumnCall(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeIncludesColumnCall")
case "FieldValue":
statFn(CounterQueryFieldValueTotal)
res, err := e.executeFieldValueCall(ctx, qcx, index, c, shards, opt)
res, err := e.executeFieldValueCall(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeFieldValueCall")
case "Precomputed":
statFn(CounterQueryPrecomputedTotal)
@ -810,23 +828,23 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
return res, errors.Wrap(err, "executeUnionRows")
case "ConstRow":
statFn(CounterQueryConstRowTotal)
res, err := e.executeConstRow(ctx, qcx, index, c, opt)
res, err := e.executeConstRow(ctx, qcx, index, c, opt, nodes)
return res, errors.Wrap(err, "executeConstRow")
case "Limit":
statFn(CounterQueryLimitTotal)
res, err := e.executeLimitCall(ctx, qcx, index, c, shards, opt)
res, err := e.executeLimitCall(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeLimitCall")
case "Percentile":
statFn(CounterQueryPercentileTotal)
res, err := e.executePercentile(ctx, qcx, index, c, shards, opt)
res, err := e.executePercentile(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executePercentile")
case "Delete":
statFn(CounterQueryDeleteTotal)
res, err := e.executeDeleteRecords(ctx, qcx, index, c, shards, opt)
res, err := e.executeDeleteRecords(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeDelete")
case "Sort":
statFn(CounterQuerySortTotal)
res, err := e.executeSort(ctx, qcx, index, c, shards, opt)
res, err := e.executeSort(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeSort")
case "Apply":
statFn(CounterQueryApplyTotal)
@ -837,7 +855,7 @@ func (e *executor) executeCall(ctx context.Context, qcx *Qcx, index string, c *p
res, err := e.executeArrow(ctx, qcx, index, c, shards, opt)
return res, errors.Wrap(err, "executeArrow")
default: // e.g. "Row", "Union", "Intersect" or anything that returns a bitmap.
res, err := e.executeBitmapCall(ctx, qcx, index, c, shards, opt)
res, err := e.executeBitmapCall(ctx, qcx, index, c, shards, opt, nodes)
return res, errors.Wrap(err, "executeBitmapCall")
}
}
@ -880,12 +898,26 @@ func (e *executor) validateTimeCallArgs(c *pql.Call, indexName string) error {
return nil
}
func (e *executor) executeOptionsCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeOptionsCall")
defer span.Finish()
func (e *executor) nodesFromOptionsCall(c *pql.Call) ([]string, error) {
var nodes []string
if arg, ok := c.Args["nodes"]; ok {
if optNodes, ok := arg.([]interface{}); ok {
for _, n := range optNodes {
if node, ok := n.(string); ok {
fmt.Println("Trying to execute against the following node: ", n)
nodes = append(nodes, string(node))
} else {
return nil, errors.New("Options(): shards must be a list of strings")
}
}
} else {
return nil, errors.New("Options(): shards must be a list of strings")
}
}
return nodes, nil
}
optCopy := &ExecOptions{}
*optCopy = *opt
func (e *executor) shardsFromOptionsCall(c *pql.Call, shards []uint64) ([]uint64, error) {
if arg, ok := c.Args["shards"]; ok {
if optShards, ok := arg.([]interface{}); ok {
shards = []uint64{}
@ -900,11 +932,31 @@ func (e *executor) executeOptionsCall(ctx context.Context, qcx *Qcx, index strin
return nil, errors.New("Query(): shards must be a list of unsigned integers")
}
}
return e.executeCall(ctx, qcx, index, c.Children[0], shards, optCopy)
return shards, nil
}
func (e *executor) executeOptionsCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeOptionsCall")
defer span.Finish()
optCopy := &ExecOptions{}
*optCopy = *opt
shards, err := e.shardsFromOptionsCall(c, shards)
if err != nil {
return nil, err
}
nodes, err := e.nodesFromOptionsCall(c)
if err != nil {
return nil, err
}
return e.executeCall(ctx, qcx, index, c.Children[0], shards, optCopy, nodes)
}
// executeIncludesColumnCall executes an IncludesColumn() call.
func (e *executor) executeIncludesColumnCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (bool, error) {
func (e *executor) executeIncludesColumnCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (bool, error) {
// Get the shard containing the column, since that's the only
// shard that needs to execute this query.
var shard uint64
@ -932,7 +984,7 @@ func (e *executor) executeIncludesColumnCall(ctx context.Context, qcx *Qcx, inde
return other || v.(bool)
}
result, err := e.mapReduce(ctx, index, []uint64{shard}, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, []uint64{shard}, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return false, err
}
@ -940,7 +992,7 @@ func (e *executor) executeIncludesColumnCall(ctx context.Context, qcx *Qcx, inde
}
// executeFieldValueCall executes a FieldValue() call.
func (e *executor) executeFieldValueCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ ValCount, err error) {
func (e *executor) executeFieldValueCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ ValCount, err error) {
fieldName, ok := c.Args["field"].(string)
if !ok || fieldName == "" {
return ValCount{}, ErrFieldRequired
@ -984,7 +1036,7 @@ func (e *executor) executeFieldValueCall(ctx context.Context, qcx *Qcx, index st
return v
}
result, err := e.mapReduce(ctx, index, []uint64{shard}, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, []uint64{shard}, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return ValCount{}, errors.Wrap(err, "map reduce")
}
@ -1024,7 +1076,7 @@ func (e *executor) executeFieldValueCallShard(ctx context.Context, qcx *Qcx, fie
}
// executeLimitCall executes a Limit() call.
func (e *executor) executeLimitCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*Row, error) {
func (e *executor) executeLimitCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (*Row, error) {
bitmapCall := c.Children[0]
limit, hasLimit, err := c.UintArg("limit")
@ -1041,7 +1093,7 @@ func (e *executor) executeLimitCall(ctx context.Context, qcx *Qcx, index string,
}
// Execute bitmap call, storing the full result on this node.
res, err := e.executeCall(ctx, qcx, index, bitmapCall, shards, opt)
res, err := e.executeCall(ctx, qcx, index, bitmapCall, shards, opt, nodes)
if err != nil {
return nil, errors.Wrap(err, "limit map reduce")
}
@ -1116,7 +1168,7 @@ func (e *executor) executeIncludesColumnCallShard(ctx context.Context, qcx *Qcx,
}
// executeSum executes a Sum() call.
func (e *executor) executeSum(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ ValCount, err error) {
func (e *executor) executeSum(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ ValCount, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeSum")
defer span.Finish()
@ -1140,7 +1192,7 @@ func (e *executor) executeSum(ctx context.Context, qcx *Qcx, index string, c *pq
return other.Add(v.(ValCount))
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return ValCount{}, err
}
@ -1170,7 +1222,7 @@ func (e *executor) executeSum(ctx context.Context, qcx *Qcx, index string, c *pq
// executeDistinct executes a Distinct call on a field. It returns a
// SignedRow for int fields and a *Row for set/mutex/time fields.
func (e *executor) executeDistinct(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) {
func (e *executor) executeDistinct(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (interface{}, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeDistinct")
defer span.Finish()
@ -1210,7 +1262,7 @@ func (e *executor) executeDistinct(ctx context.Context, qcx *Qcx, index string,
}
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, errors.Wrap(err, "mapReduce")
}
@ -1222,7 +1274,7 @@ func (e *executor) executeDistinct(ctx context.Context, qcx *Qcx, index string,
}
// executeMin executes a Min() call.
func (e *executor) executeMin(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ ValCount, err error) {
func (e *executor) executeMin(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ ValCount, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeMin")
defer span.Finish()
@ -1245,7 +1297,7 @@ func (e *executor) executeMin(ctx context.Context, qcx *Qcx, index string, c *pq
return other.Smaller(v.(ValCount))
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return ValCount{}, err
}
@ -1258,7 +1310,7 @@ func (e *executor) executeMin(ctx context.Context, qcx *Qcx, index string, c *pq
}
// executeMax executes a Max() call.
func (e *executor) executeMax(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ ValCount, err error) {
func (e *executor) executeMax(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ ValCount, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeMax")
defer span.Finish()
@ -1281,7 +1333,7 @@ func (e *executor) executeMax(ctx context.Context, qcx *Qcx, index string, c *pq
return other.Larger(v.(ValCount))
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return ValCount{}, err
}
@ -1294,7 +1346,7 @@ func (e *executor) executeMax(ctx context.Context, qcx *Qcx, index string, c *pq
}
// executePercentile executes a Percentile() call.
func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ ValCount, err error) {
func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ ValCount, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executePercentile")
defer span.Finish()
@ -1340,7 +1392,7 @@ func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string
if filterCall != nil {
minCall.Children = append(minCall.Children, filterCall)
}
minVal, err := e.executeMin(ctx, qcx, index, minCall, shards, opt)
minVal, err := e.executeMin(ctx, qcx, index, minCall, shards, opt, nodes)
if err != nil {
return ValCount{}, errors.Wrap(err, "executing Min call for Percentile")
}
@ -1354,7 +1406,7 @@ func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string
if filterCall != nil {
maxCall.Children = append(maxCall.Children, filterCall)
}
maxVal, err := e.executeMax(ctx, qcx, index, maxCall, shards, opt)
maxVal, err := e.executeMax(ctx, qcx, index, maxCall, shards, opt, nodes)
if err != nil {
return ValCount{}, errors.Wrap(err, "executing Max call for Percentile")
}
@ -1386,7 +1438,7 @@ func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string
Op: pql.Token(pql.LT),
Value: possibleNthVal,
}
leftCountUint64, err := e.executeCount(ctx, qcx, index, countCall, shards, opt)
leftCountUint64, err := e.executeCount(ctx, qcx, index, countCall, shards, opt, nodes)
if err != nil {
return ValCount{}, errors.Wrap(err, "executing Count call L for Percentile")
}
@ -1397,7 +1449,7 @@ func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string
Op: pql.Token(pql.GT),
Value: possibleNthVal,
}
rightCountUint64, err := e.executeCount(ctx, qcx, index, countCall, shards, opt)
rightCountUint64, err := e.executeCount(ctx, qcx, index, countCall, shards, opt, nodes)
if err != nil {
return ValCount{}, errors.Wrap(err, "executing Count call R for Percentile")
}
@ -1420,7 +1472,7 @@ func (e *executor) executePercentile(ctx context.Context, qcx *Qcx, index string
}
// executeMinRow executes a MinRow() call.
func (e *executor) executeMinRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ interface{}, err error) {
func (e *executor) executeMinRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ interface{}, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeMinRow")
defer span.Finish()
@ -1455,11 +1507,11 @@ func (e *executor) executeMinRow(ctx context.Context, qcx *Qcx, index string, c
return vp
}
return e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
return e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
}
// executeMaxRow executes a MaxRow() call.
func (e *executor) executeMaxRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ interface{}, err error) {
func (e *executor) executeMaxRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ interface{}, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeMaxRow")
defer span.Finish()
@ -1494,7 +1546,7 @@ func (e *executor) executeMaxRow(ctx context.Context, qcx *Qcx, index string, c
return vp
}
return e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
return e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
}
// executePrecomputedCall pretends to execute a call that we have a precomputed value for.
@ -1510,7 +1562,7 @@ func (e *executor) executePrecomputedCall(ctx context.Context, qcx *Qcx, index s
}
// executeBitmapCall executes a call that returns a bitmap.
func (e *executor) executeBitmapCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ *Row, err error) {
func (e *executor) executeBitmapCall(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ *Row, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeBitmapCall")
span.LogKV("pqlCallName", c.Name)
defer span.Finish()
@ -1572,7 +1624,7 @@ func (e *executor) executeBitmapCall(ctx context.Context, qcx *Qcx, index string
return other
}
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, errors.Wrap(err, "map reduce")
}
@ -2158,7 +2210,7 @@ func (e *executor) executeMaxRowShard(ctx context.Context, qcx *Qcx, index strin
}, nil
}
func (e *executor) executeTopK(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) {
func (e *executor) executeTopK(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (interface{}, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeTopK")
defer span.Finish()
@ -2172,7 +2224,7 @@ func (e *executor) executeTopK(ctx context.Context, qcx *Qcx, index string, c *p
return ([]*Row)(AddBSI(x, y))
}
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, err
}
@ -2580,7 +2632,7 @@ func (f *topKFilter) fillIt(it roaring.ContainerIterator) {
// executeTopN executes a TopN() call.
// This first performs the TopN() to determine the top results and then
// requeries to retrieve the full counts for each of the top results.
func (e *executor) executeTopN(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*PairsField, error) {
func (e *executor) executeTopN(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (*PairsField, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeTopN")
defer span.Finish()
@ -2596,7 +2648,7 @@ func (e *executor) executeTopN(ctx context.Context, qcx *Qcx, index string, c *p
}
// Execute original query.
pairs, err := e.executeTopNShards(ctx, qcx, index, c, shards, opt)
pairs, err := e.executeTopNShards(ctx, qcx, index, c, shards, opt, nodes)
if err != nil {
return nil, errors.Wrap(err, "finding top results")
}
@ -2617,7 +2669,7 @@ func (e *executor) executeTopN(ctx context.Context, qcx *Qcx, index string, c *p
sort.Sort(uint64Slice(ids))
other.Args["ids"] = ids
trimmedList, err := e.executeTopNShards(ctx, qcx, index, other, shards, opt)
trimmedList, err := e.executeTopNShards(ctx, qcx, index, other, shards, opt, nodes)
if err != nil {
return nil, errors.Wrap(err, "retrieving full counts")
}
@ -2632,7 +2684,7 @@ func (e *executor) executeTopN(ctx context.Context, qcx *Qcx, index string, c *p
}, nil
}
func (e *executor) executeTopNShards(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*PairsField, error) {
func (e *executor) executeTopNShards(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (*PairsField, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeTopNShards")
defer span.Finish()
@ -2657,7 +2709,7 @@ func (e *executor) executeTopNShards(ctx context.Context, qcx *Qcx, index string
return other
}
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, err
}
@ -2977,7 +3029,7 @@ func findGroupCounts(v interface{}) []GroupCount {
return nil
}
func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*GroupCounts, error) {
func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (*GroupCounts, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeGroupBy")
defer span.Finish()
// validate call
@ -3071,7 +3123,7 @@ func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c
continue
}
r, er := e.executeRows(ctx, qcx, index, child, shards, opt)
r, er := e.executeRows(ctx, qcx, index, child, shards, opt, nodes)
if er != nil {
return nil, errors.Wrap(er, "getting rows for ")
}
@ -3122,7 +3174,7 @@ func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c
return x
}
// Get full result set.
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, errors.Wrap(err, "mapReduce")
}
@ -3181,7 +3233,7 @@ func (e *executor) executeGroupBy(ctx context.Context, qcx *Qcx, index string, c
}
opt.PreTranslated = true
aggregateCount, err := e.execute(ctx, qcx, index, &pql.Query{Calls: []*pql.Call{countDistinctIntersect}}, []uint64{}, opt)
aggregateCount, err := e.execute(ctx, qcx, index, &pql.Query{Calls: []*pql.Call{countDistinctIntersect}}, []uint64{}, opt, nodes)
if err != nil {
return nil, err
}
@ -3788,7 +3840,7 @@ func (e *executor) executeGroupByShard(ctx context.Context, qcx *Qcx, index stri
return results, nil
}
func (e *executor) executeRows(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (RowIDs, error) {
func (e *executor) executeRows(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (RowIDs, error) {
// Fetch field name from argument.
// Check "field" first for backwards compatibility.
// TODO: remove at Pilosa 2.0
@ -3841,7 +3893,7 @@ func (e *executor) executeRows(ctx context.Context, qcx *Qcx, index string, c *p
return other.Merge(v.(RowIDs), limit)
}
// Get full result set.
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, err
}
@ -4189,7 +4241,7 @@ func (e *executor) executeExternalLookup(ctx context.Context, qcx *Qcx, index st
return ExtractedTable{}, errors.Wrap(err, "parsing write argument")
}
rawArg, err := e.executeCall(ctx, qcx, index, c.Children[0], shards, opt)
rawArg, err := e.executeCall(ctx, qcx, index, c.Children[0], shards, opt, nil)
if err != nil {
return ExtractedTable{}, errors.Wrapf(err, "evaluating SQL argument call %q", c.String())
}
@ -4512,7 +4564,7 @@ func handleExtractResults(other interface{}, filter *pql.Call, opt *ExecOptions)
}
}
func (e *executor) executeExtract(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (interface{}, error) {
func (e *executor) executeExtract(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (interface{}, error) {
// Extract the column filter call.
if len(c.Children) < 1 {
return ExtractedIDMatrix{}, errors.New("missing column filter in Extract")
@ -4541,7 +4593,7 @@ func (e *executor) executeExtract(ctx context.Context, qcx *Qcx, index string, c
// Merge returned results at coordinating node.
reduceFn := makeReduceFunc(sort_desc)
// Get full result set.
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return ExtractedIDMatrix{}, err
}
@ -4659,7 +4711,7 @@ func (e *executor) executeExtractShard(ctx context.Context, qcx *Qcx, index stri
}
}
case FieldTypeTime:
// Handle a set field by listing the rows and then intersecting them with the filter.
// Handle a set field by listing the rows and then intersecting them with the filter.
timeArg := timeArgs[i]
views, err := field.viewsByTimeRange(timeArg.From, timeArg.To)
@ -5355,7 +5407,7 @@ func (e *executor) executeNotShard(ctx context.Context, qcx *Qcx, index string,
return existenceRow.Difference(row), nil
}
func (e *executor) executeConstRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, opt *ExecOptions) (res *Row, err error) {
func (e *executor) executeConstRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, opt *ExecOptions, nodes []string) (res *Row, err error) {
idx := e.Holder.Index(index)
if idx == nil {
return nil, newNotFoundError(ErrIndexNotFound, index)
@ -5412,7 +5464,7 @@ func (e *executor) executeConstRow(ctx context.Context, qcx *Qcx, index string,
return other
}
other, err := e.mapReduce(ctx, index, filteredShards, c, opt, mapFn, reduceFn)
other, err := e.mapReduce(ctx, index, filteredShards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, err
}
@ -5460,7 +5512,7 @@ func (e *executor) executeUnionRows(ctx context.Context, qcx *Qcx, index string,
}
// Execute the call.
rowsResult, err := e.executeCall(ctx, qcx, index, child, shards, opt)
rowsResult, err := e.executeCall(ctx, qcx, index, child, shards, opt, nil)
if err != nil {
return nil, err
}
@ -5528,7 +5580,7 @@ func (e *executor) executeUnionRows(ctx context.Context, qcx *Qcx, index string,
}
// Execute the generated Union() call.
return e.executeBitmapCall(ctx, qcx, index, c, shards, opt)
return e.executeBitmapCall(ctx, qcx, index, c, shards, opt, nil)
}
// executeAllCallShard executes an All() call for a local shard.
@ -5590,7 +5642,7 @@ func (e *executor) executeShiftShard(ctx context.Context, qcx *Qcx, index string
}
// executeCount executes a count() call.
func (e *executor) executeCount(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (uint64, error) {
func (e *executor) executeCount(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (uint64, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeCount")
defer span.Finish()
@ -5604,7 +5656,7 @@ func (e *executor) executeCount(ctx context.Context, qcx *Qcx, index string, c *
// If the child is distinct/similar, execute it directly here and count the result.
if child.Type == pql.PrecallGlobal {
result, err := e.executeCall(ctx, qcx, index, child, shards, opt)
result, err := e.executeCall(ctx, qcx, index, child, shards, opt, nil)
if err != nil {
return 0, err
}
@ -5636,7 +5688,7 @@ func (e *executor) executeCount(ctx context.Context, qcx *Qcx, index string, c *
return other + v.(uint64)
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return 0, err
}
@ -5727,7 +5779,7 @@ func (e *executor) executeClearBitField(ctx context.Context, qcx *Qcx, index str
}
// executeClearRow executes a ClearRow() call.
func (e *executor) executeClearRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (_ bool, err error) {
func (e *executor) executeClearRow(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (_ bool, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeClearRow")
defer span.Finish()
@ -5770,7 +5822,7 @@ func (e *executor) executeClearRow(ctx context.Context, qcx *Qcx, index string,
return pval
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return false, errors.Wrap(err, "mapreducing clearrow")
}
@ -5828,7 +5880,7 @@ func (e *executor) executeClearRowShard(ctx context.Context, qcx *Qcx, index str
// executeSetRow executes a Store() call.
func (e *executor) executeSetRow(ctx context.Context, qcx *Qcx, indexName string, c *pql.Call, shards []uint64, opt *ExecOptions) (bool, error) {
func (e *executor) executeSetRow(ctx context.Context, qcx *Qcx, indexName string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (bool, error) {
// Parse arguments.
fieldName, err := c.FieldArg()
if err != nil {
@ -5866,7 +5918,7 @@ func (e *executor) executeSetRow(ctx context.Context, qcx *Qcx, indexName string
return pval
}
result, err := e.mapReduce(ctx, indexName, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, indexName, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return false, err
}
@ -6200,7 +6252,7 @@ loop:
// mapReduce has to ensure that it never returns before any work it spawned has
// terminated. It's not enough to cancel the jobs; we have to wait for them to be
// done, or we can unmap resources they're still using.
func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64, c *pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc) (result interface{}, err error) {
func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64, c *pql.Call, opt *ExecOptions, mapFn mapFunc, reduceFn reduceFunc, nodesOpt []string) (result interface{}, err error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.mapReduce")
defer span.Finish()
@ -6233,6 +6285,25 @@ func (e *executor) mapReduce(ctx context.Context, index string, shards []uint64,
nodes = []*disco.Node{e.Cluster.nodeByID(e.Node.ID)}
}
if nodesOpt != nil {
// remove nodes that weren't included in the Options call
// under the nodes argument
var nodesToKeep []*disco.Node
for _, node := range nodes {
found := false
for _, nodeOpt := range nodesOpt {
if node.ID == nodeOpt {
found = true
continue
}
}
if found {
nodesToKeep = append(nodesToKeep, node)
}
}
nodes = nodesToKeep[:]
}
// Start mapping across all primary owners.
if err = e.mapper(ctx, eg, ch, nodes, index, shards, c, opt, e.Cluster.ReplicaN == 1, mapFn, reduceFn); err != nil {
return nil, errors.Wrap(err, "starting mapper")
@ -7940,6 +8011,7 @@ type mapFunc func(ctx context.Context, shard uint64, mopt *mapOptions) (_ interf
type mapOptions struct {
memoryAvailable *int64
nodes []string
}
type reduceFunc func(ctx context.Context, prev, v interface{}) interface{}
@ -8801,7 +8873,7 @@ func decimalToInt64(dec pql.Decimal, opt FieldOptions) int64 {
}
// executeDeleteRecords executes a delete() call.
func (e *executor) executeDeleteRecords(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (bool, error) {
func (e *executor) executeDeleteRecords(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (bool, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "executor.executeDelete")
defer span.Finish()
@ -8824,7 +8896,7 @@ func (e *executor) executeDeleteRecords(ctx context.Context, qcx *Qcx, index str
return other || v.(bool)
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return false, err
}
@ -9073,7 +9145,7 @@ func deleteKeyTranslation(ctx context.Context, idx *Index, shard uint64, records
return idx.TranslateStore(paritionID).Delete(records)
}
func (e *executor) executeSort(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions) (*SortedRow, error) {
func (e *executor) executeSort(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *ExecOptions, nodes []string) (*SortedRow, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "Executor.executeSort")
defer span.Finish()
@ -9114,7 +9186,7 @@ func (e *executor) executeSort(ctx context.Context, qcx *Qcx, index string, c *p
}
}
res, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
res, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn, nodes)
if err != nil {
return nil, errors.Wrap(err, "mapReduce")
}

View file

@ -596,6 +596,7 @@ var callInfoByFunc = map[string]callInfo{
allowUnknown: false,
prototypes: map[string]interface{}{
"shards": nil,
"nodes": nil,
},
},
"Set": {