add multi-shard translation

This commit is contained in:
Ben Johnson 2020-01-03 14:12:39 -07:00
parent bdfdeb1291
commit a043490996
3 changed files with 154 additions and 135 deletions

View file

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

View file

@ -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)
}

View file

@ -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)