From a043490996953349c27c0afce822a593e08ce6f2 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Fri, 3 Jan 2020 14:12:39 -0700 Subject: [PATCH] add multi-shard translation --- executor.go | 284 ++++++++++++++++++++------------------ executor_internal_test.go | 4 +- index.go | 1 - 3 files changed, 154 insertions(+), 135 deletions(-) diff --git a/executor.go b/executor.go index 3f9a48d67..125f46358 100644 --- a/executor.go +++ b/executor.go @@ -3526,164 +3526,196 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu } } -func (e *executor) translateCalls(ctx context.Context, index string, idx *Index, calls []*pql.Call) (err error) { +func (e *executor) translateCalls(ctx context.Context, defaultIndexName string, defaultIdx *Index, calls []*pql.Call) (err error) { span, _ := tracing.StartSpanFromContext(ctx, "Executor.translateCalls") defer span.Finish() - // TODO(BBJ): Handle cross-index boundaries. - keyMap := make(map[string]uint64) - if idx.Keys() { - // Collect all index keys. - keySet := make(map[string]struct{}) - for i := range calls { - if err := e.collectCallIndexKeys(index, idx, calls[i], keySet); err != nil { - return err - } - } - - if keyMap, err = e.Cluster.translateIndexKeySet(ctx, index, keySet); err != nil { + // Generate a list of all used + indexNameMap := make(map[string]struct{}) + for i := range calls { + if err := e.collectCallIndexNameMap(ctx, defaultIndexName, calls[i], indexNameMap); err != nil { return err } } - // Translate calls. - for i := range calls { - // Possibly change to another index for translation, if this - // call crosses index boundaries. - newIdxName := calls[i].CallIndex() - var newIdx *Index - if newIdxName == "" || newIdxName == index { - newIdxName = index - newIdx = idx - } else { - newIdx = idx.holder.indexes[newIdxName] - if newIdx == nil { - return fmt.Errorf("unknown index %q specified in cross-index call", newIdxName) + // Perform a separate batch translation for each separate index used. + for indexName := range indexNameMap { + // Determine the target index name. + if indexName == "" { + indexName = defaultIndexName + } + isDefaultIndex := indexName == defaultIndexName + + // Determine the target index. + idx := defaultIdx + if !isDefaultIndex { + idx = idx.holder.indexes[indexName] + if idx == nil { + return fmt.Errorf("unknown index %q specified in cross-index call", indexName) } } - if err := e.translateCall(newIdxName, newIdx, calls[i], keyMap); err != nil { + + // Collect all index keys & bulk translate them. + keyMap := make(map[string]uint64) + if idx.Keys() { + keySet := make(map[string]struct{}) + for i := range calls { + if err := e.collectCallIndexKeys(indexName, idx, isDefaultIndex, calls[i], keySet); err != nil { + return err + } + } + + if keyMap, err = e.Cluster.translateIndexKeySet(ctx, indexName, keySet); err != nil { + return err + } + } + + // Translate calls. + for i := range calls { + if err := e.translateCall(indexName, idx, isDefaultIndex, calls[i], keyMap); err != nil { + return err + } + } + } + return nil +} + +func (e *executor) collectCallIndexNameMap(ctx context.Context, defaultIndexName string, c *pql.Call, m map[string]struct{}) error { + callIndex := c.CallIndex() + if callIndex == "" { + callIndex = defaultIndexName + } + m[callIndex] = struct{}{} + + if c.Name == "GroupBy" { + if filter, ok, err := c.CallArg("filter"); ok { + if err != nil { + return errors.Wrap(err, "getting filter call") + } + err := e.collectCallIndexNameMap(ctx, defaultIndexName, filter, m) + if err != nil { + return errors.Wrap(err, "collecting filter call index name") + } + } + } + + for _, child := range c.Children { + if err := e.collectCallIndexNameMap(ctx, defaultIndexName, child, m); err != nil { return err } } return nil } -func (e *executor) collectCallIndexKeys(index string, idx *Index, c *pql.Call, keySet map[string]struct{}) error { +func (e *executor) collectCallIndexKeys(index string, idx *Index, isDefaultIndex bool, c *pql.Call, keySet map[string]struct{}) error { // Handle group by separately. if c.Name == "GroupBy" { for _, child := range c.Children { - if err := e.collectCallIndexKeys(index, idx, child, keySet); err != nil { + if err := e.collectCallIndexKeys(index, idx, isDefaultIndex, child, keySet); err != nil { return errors.Wrapf(err, "translating %s", child) } } - if filter, ok, err := c.CallArg("filter"); ok { - if err != nil { - return errors.Wrap(err, "getting filter call") - } - err = e.collectCallIndexKeys(index, idx, filter, keySet) - if err != nil { - return errors.Wrap(err, "translating filter call") + if callIndex := c.CallIndex(); callIndex == index || (callIndex == "" && isDefaultIndex) { + if filter, ok, err := c.CallArg("filter"); ok { + if err != nil { + return errors.Wrap(err, "getting filter call") + } + err = e.collectCallIndexKeys(index, idx, isDefaultIndex, filter, keySet) + if err != nil { + return errors.Wrap(err, "translating filter call") + } } } return nil } - colKey, _, _ := c.TranslateInfo(columnLabel, rowLabel) - if c.Args[colKey] != nil && isString(c.Args[colKey]) { - if value := callArgString(c, colKey); value != "" { - keySet[value] = struct{}{} + if callIndex := c.CallIndex(); callIndex == index || (callIndex == "" && isDefaultIndex) { + colKey, _, _ := c.TranslateInfo(columnLabel, rowLabel) + if c.Args[colKey] != nil && isString(c.Args[colKey]) { + if value := callArgString(c, colKey); value != "" { + keySet[value] = struct{}{} + } } } return nil } -func (e *executor) translateCall(index string, idx *Index, c *pql.Call, keyMap map[string]uint64) error { +func (e *executor) translateCall(indexName string, idx *Index, isDefaultIndex bool, c *pql.Call, keyMap map[string]uint64) error { if c.Name == "GroupBy" { - return errors.Wrap(e.translateGroupByCall(index, idx, c, keyMap), "translating GroupBy") + return errors.Wrap(e.translateGroupByCall(indexName, idx, isDefaultIndex, c, keyMap), "translating GroupBy") } // Translate column key. - colKey, rowKey, fieldName := c.TranslateInfo(columnLabel, rowLabel) - if idx.Keys() { - if c.Args[colKey] != nil && !isString(c.Args[colKey]) { - if !isValidID(c.Args[colKey]) { - return errors.Errorf("column value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[colKey]) - } - } else if value := callArgString(c, colKey); value != "" { - c.Args[colKey] = keyMap[value] - } - } else { - if isString(c.Args[colKey]) { - return errors.New("string 'col' value not allowed unless index 'keys' option enabled") - } - } - - // Translate row key, if field is specified & key exists. - if fieldName != "" { - field := idx.Field(fieldName) - if field == nil { - // Instead of returning ErrFieldNotFound here, - // we just return, and don't attempt the translation. - // The assumption is that the non-existent field - // will raise an error downstream when it's used. - return nil - } - - // Bool field keys do not use the translator because there - // are only two possible values. Instead, they are handled - // directly. - if field.Type() == FieldTypeBool { - // TODO: This code block doesn't make sense for a `Rows()` - // queries on a `bool` field. Need to review this better, - // include it in tests, and probably back-port it to Pilosa. - if c.Name != "Rows" { - boolVal, err := callArgBool(c, rowKey) - if err != nil { - return errors.Wrap(err, "getting bool key") + if callIndex := c.CallIndex(); callIndex == indexName || (callIndex == "" && isDefaultIndex) { + colKey, rowKey, fieldName := c.TranslateInfo(columnLabel, rowLabel) + if idx.Keys() { + if c.Args[colKey] != nil && !isString(c.Args[colKey]) { + if !isValidID(c.Args[colKey]) { + return errors.Errorf("column value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[colKey]) } - rowID := falseRowID - if boolVal { - rowID = trueRowID - } - c.Args[rowKey] = rowID - } - } else if field.keys() { - if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) { - // allow passing row id directly (this can come in handy, but make sure it is a valid row id) - if !isValidID(c.Args[rowKey]) { - return errors.Errorf("row value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[rowKey]) - } - } else if value := callArgString(c, rowKey); value != "" { - id, err := field.TranslateStore().TranslateKey(value) - if err != nil { - return err - } - c.Args[rowKey] = id + } else if value := callArgString(c, colKey); value != "" { + c.Args[colKey] = keyMap[value] } } else { - if isString(c.Args[rowKey]) { - return errors.New("string 'row' value not allowed unless field 'keys' option enabled") + if isString(c.Args[colKey]) { + return errors.New("string 'col' value not allowed unless index 'keys' option enabled") + } + } + + // Translate row key, if field is specified & key exists. + if fieldName != "" { + field := idx.Field(fieldName) + if field == nil { + // Instead of returning ErrFieldNotFound here, + // we just return, and don't attempt the translation. + // The assumption is that the non-existent field + // will raise an error downstream when it's used. + return nil + } + + // Bool field keys do not use the translator because there + // are only two possible values. Instead, they are handled + // directly. + if field.Type() == FieldTypeBool { + // TODO: This code block doesn't make sense for a `Rows()` + // queries on a `bool` field. Need to review this better, + // include it in tests, and probably back-port it to Pilosa. + if c.Name != "Rows" { + boolVal, err := callArgBool(c, rowKey) + if err != nil { + return errors.Wrap(err, "getting bool key") + } + rowID := falseRowID + if boolVal { + rowID = trueRowID + } + c.Args[rowKey] = rowID + } + } else if field.keys() { + if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) { + // allow passing row id directly (this can come in handy, but make sure it is a valid row id) + if !isValidID(c.Args[rowKey]) { + return errors.Errorf("row value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[rowKey]) + } + } else if value := callArgString(c, rowKey); value != "" { + id, err := field.TranslateStore().TranslateKey(value) + if err != nil { + return err + } + c.Args[rowKey] = id + } + } else { + if isString(c.Args[rowKey]) { + return errors.New("string 'row' value not allowed unless field 'keys' option enabled") + } } } } // Translate child calls. for _, child := range c.Children { - // Possibly change to another index for translation, if this - // call crosses index boundaries. - newIdxName := child.CallIndex() - var newIdx *Index - if newIdxName == "" || newIdxName == index { - newIdxName = index - newIdx = idx - } else { - newIdx = idx.holder.indexes[newIdxName] - if newIdx == nil { - return fmt.Errorf("unknown index %q specified in cross-index call", newIdxName) - } - } - if err := e.translateCall(newIdxName, newIdx, child, keyMap); err != nil { + if err := e.translateCall(indexName, idx, isDefaultIndex, child, keyMap); err != nil { return err } } @@ -3691,34 +3723,22 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call, keyMap m return nil } -func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call, keyMap map[string]uint64) error { +func (e *executor) translateGroupByCall(index string, idx *Index, isDefaultIndex bool, c *pql.Call, keyMap map[string]uint64) error { if c.Name != "GroupBy" { panic("translateGroupByCall called with '" + c.Name + "'") } for _, child := range c.Children { - if err := e.translateCall(index, idx, child, keyMap); err != nil { + if err := e.translateCall(index, idx, isDefaultIndex, child, keyMap); err != nil { return errors.Wrapf(err, "translating %s", child) } } - if filter, ok, err := c.CallArg("filter"); ok { - if err != nil { - return errors.Wrap(err, "getting filter call") - } - err = e.translateCall(index, idx, filter, keyMap) - if err != nil { - return errors.Wrap(err, "translating filter call") - } - } - - if aggregate, ok, err := c.CallArg("aggregate"); ok { - if err != nil { - return errors.Wrap(err, "getting aggregate call") - } - err = e.translateCall(index, idx, aggregate, keyMap) - if err != nil { - return errors.Wrap(err, "translating aggregate call") + for _, arg := range c.Args { + if arg, ok := arg.(*pql.Call); ok { + if err := e.translateCall(index, idx, isDefaultIndex, arg, keyMap); err != nil { + return errors.Wrap(err, "translating group by arg") + } } } diff --git a/executor_internal_test.go b/executor_internal_test.go index b4487dca6..0c1ccd4ba 100644 --- a/executor_internal_test.go +++ b/executor_internal_test.go @@ -61,7 +61,7 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) { t.Fatalf("parsing query: %v", err) } c := query.Calls[0] - err = e.translateGroupByCall("i", idx, c, make(map[string]uint64)) + err = e.translateGroupByCall("i", idx, true, c, make(map[string]uint64)) if err != nil { t.Fatalf("translating call: %v", err) } @@ -125,7 +125,7 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) { t.Fatalf("parsing query: %v", err) } c := query.Calls[0] - err = e.translateGroupByCall("i", idx, c, make(map[string]uint64)) + err = e.translateGroupByCall("i", idx, true, c, make(map[string]uint64)) if err == nil { t.Fatalf("expected error, but translated call is '%s", c) } diff --git a/index.go b/index.go index a8adfc674..c54aa9814 100644 --- a/index.go +++ b/index.go @@ -164,7 +164,6 @@ func (i *Index) Open() (err error) { return errors.Wrap(err, "opening attrstore") } - // TODO(BBJ): Support non-default partition counts. i.logger.Debugf("open translate store for index: %s", i.name) for partitionID := 0; partitionID < i.partitionN; partitionID++ { store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.partitionN)