Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 62 additions & 0 deletions data/test/vtgate/select_cases.txt
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,68 @@
}
}

# select with partial scatter directive
"select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ * from user"
{
"Original": "select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ * from user",
"Instructions": {
"Opcode": "SelectScatter",
"Keyspace": {
"Name": "user",
"Sharded": true
},
"Query": "select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ * from user",
"FieldQuery": "select * from user where 1 != 1",
"ScatterErrorsAsWarnings": true
}
}

# select aggregation with partial scatter directive
"select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ count(*) from user"
{
"Original": "select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ count(*) from user",
"Instructions": {
"Aggregates": [
{
"Opcode": "count",
"Col": 0
}
],
"Keys": null,
"Input": {
"Opcode": "SelectScatter",
"Keyspace": {
"Name": "user",
"Sharded": true
},
"Query": "select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ count(*) from user",
"FieldQuery": "select count(*) from user where 1 != 1",
"ScatterErrorsAsWarnings": true
}
}
}

# select limit with partial scatter directive
"select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ * from user limit 10"
{
"Original": "select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ * from user limit 10",
"Instructions": {
"Opcode": "Limit",
"Count": 10,
"Offset": null,
"Input": {
"Opcode": "SelectScatter",
"Keyspace": {
"Name": "user",
"Sharded": true
},
"Query": "select /*vt+ SCATTER_ERRORS_AS_WARNINGS=1 */ * from user limit :__upper_limit",
"FieldQuery": "select * from user where 1 != 1",
"ScatterErrorsAsWarnings": true
}
}
}

# qualified '*' expression for simple route
"select user.* from user"
{
Expand Down
10 changes: 10 additions & 0 deletions go/vt/concurrency/error_recorder.go
Original file line number Diff line number Diff line change
Expand Up @@ -135,3 +135,13 @@ func (aer *AllErrorRecorder) ErrorStrings() []string {
}
return errs
}

// GetErrors returns a reference to the internal errors array.
//
// Note that the array is not copied, so this should only be used
// once the recording is complete.
func (aer *AllErrorRecorder) GetErrors() []error {
aer.mu.Lock()
defer aer.mu.Unlock()
return aer.Errors
}
2 changes: 2 additions & 0 deletions go/vt/sqlparser/comments.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ const (
DirectiveSkipQueryPlanCache = "SKIP_QUERY_PLAN_CACHE"
// DirectiveQueryTimeout sets a query timeout in vtgate. Only supported for SELECTS.
DirectiveQueryTimeout = "QUERY_TIMEOUT_MS"
// DirectiveScatterErrorsAsWarnings enables partial success scatter select queries
DirectiveScatterErrorsAsWarnings = "SCATTER_ERRORS_AS_WARNINGS"
)

func isNonSpace(r rune) bool {
Expand Down
3 changes: 2 additions & 1 deletion go/vt/vtgate/engine/delete.go
Original file line number Diff line number Diff line change
Expand Up @@ -238,5 +238,6 @@ func (del *Delete) execDeleteByDestination(vcursor VCursor, bindVars map[string]
}
}
autocommit := (len(rss) == 1 || del.MultiShardAutocommit) && vcursor.AutocommitApproval()
return vcursor.ExecuteMultiShard(rss, queries, true /* isDML */, autocommit)
res, errs := vcursor.ExecuteMultiShard(rss, queries, true /* isDML */, autocommit)
return res, vterrors.Aggregate(errs)
}
15 changes: 12 additions & 3 deletions go/vt/vtgate/engine/fake_vcursor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ func (t noopVCursor) ExecuteAutocommit(method string, query string, bindvars map
panic("unimplemented")
}

func (t noopVCursor) ExecuteMultiShard(rss []*srvtopo.ResolvedShard, queries []*querypb.BoundQuery, isDML, autocommit bool) (*sqltypes.Result, error) {
func (t noopVCursor) ExecuteMultiShard(rss []*srvtopo.ResolvedShard, queries []*querypb.BoundQuery, isDML, autocommit bool) (*sqltypes.Result, []error) {
panic("unimplemented")
}

Expand Down Expand Up @@ -88,6 +88,10 @@ type loggingVCursor struct {
curResult int
resultErr error

// Optional errors that can be returned from nextResult() alongside the results for
// multi-shard queries
multiShardErrs []error

log []string
}

Expand All @@ -103,9 +107,14 @@ func (f *loggingVCursor) Execute(method string, query string, bindvars map[strin
return f.nextResult()
}

func (f *loggingVCursor) ExecuteMultiShard(rss []*srvtopo.ResolvedShard, queries []*querypb.BoundQuery, isDML, canAutocommit bool) (*sqltypes.Result, error) {
func (f *loggingVCursor) ExecuteMultiShard(rss []*srvtopo.ResolvedShard, queries []*querypb.BoundQuery, isDML, canAutocommit bool) (*sqltypes.Result, []error) {
f.log = append(f.log, fmt.Sprintf("ExecuteMultiShard %v%v %v", printResolvedShardQueries(rss, queries), isDML, canAutocommit))
return f.nextResult()
res, err := f.nextResult()
if err != nil {
return nil, []error{err}
}

return res, f.multiShardErrs
}

func (f *loggingVCursor) AutocommitApproval() bool {
Expand Down
6 changes: 3 additions & 3 deletions go/vt/vtgate/engine/insert.go
Original file line number Diff line number Diff line change
Expand Up @@ -220,9 +220,9 @@ func (ins *Insert) execInsertSharded(vcursor VCursor, bindVars map[string]*query
}

autocommit := (len(rss) == 1 || ins.MultiShardAutocommit) && vcursor.AutocommitApproval()
result, err := vcursor.ExecuteMultiShard(rss, queries, true /* isDML */, autocommit)
if err != nil {
return nil, vterrors.Wrap(err, "execInsertSharded")
result, errs := vcursor.ExecuteMultiShard(rss, queries, true /* isDML */, autocommit)
if errs != nil {
return nil, vterrors.Wrap(vterrors.Aggregate(errs), "execInsertSharded")
}

if insertID != 0 {
Expand Down
2 changes: 1 addition & 1 deletion go/vt/vtgate/engine/primitive.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ type VCursor interface {
AutocommitApproval() bool

// Shard-level functions.
ExecuteMultiShard(rss []*srvtopo.ResolvedShard, queries []*querypb.BoundQuery, isDML, canAutocommit bool) (*sqltypes.Result, error)
ExecuteMultiShard(rss []*srvtopo.ResolvedShard, queries []*querypb.BoundQuery, isDML, canAutocommit bool) (*sqltypes.Result, []error)
ExecuteStandalone(query string, bindvars map[string]*querypb.BindVariable, rs *srvtopo.ResolvedShard) (*sqltypes.Result, error)
StreamExecuteMulti(query string, rss []*srvtopo.ResolvedShard, bindVars []map[string]*querypb.BindVariable, callback func(reply *sqltypes.Result) error) error

Expand Down
61 changes: 39 additions & 22 deletions go/vt/vtgate/engine/route.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import (

"vitess.io/vitess/go/jsonutil"
"vitess.io/vitess/go/sqltypes"
"vitess.io/vitess/go/stats"
"vitess.io/vitess/go/vt/key"
"vitess.io/vitess/go/vt/srvtopo"
"vitess.io/vitess/go/vt/vterrors"
Expand Down Expand Up @@ -71,6 +72,9 @@ type Route struct {

// QueryTimeout contains the optional timeout (in milliseconds) to apply to this query
QueryTimeout int

// ScatterErrorsAsWarnings is true if results should be returned even if some shards have an error
ScatterErrorsAsWarnings bool
}

// OrderbyParams specifies the parameters for ordering.
Expand All @@ -88,25 +92,27 @@ func (route *Route) MarshalJSON() ([]byte, error) {
vindexName = route.Vindex.String()
}
marshalRoute := struct {
Opcode RouteOpcode
Keyspace *vindexes.Keyspace `json:",omitempty"`
Query string `json:",omitempty"`
FieldQuery string `json:",omitempty"`
Vindex string `json:",omitempty"`
Values []sqltypes.PlanValue `json:",omitempty"`
OrderBy []OrderbyParams `json:",omitempty"`
TruncateColumnCount int `json:",omitempty"`
QueryTimeout int `json:",omitempty"`
Opcode RouteOpcode
Keyspace *vindexes.Keyspace `json:",omitempty"`
Query string `json:",omitempty"`
FieldQuery string `json:",omitempty"`
Vindex string `json:",omitempty"`
Values []sqltypes.PlanValue `json:",omitempty"`
OrderBy []OrderbyParams `json:",omitempty"`
TruncateColumnCount int `json:",omitempty"`
QueryTimeout int `json:",omitempty"`
ScatterErrorsAsWarnings bool `json:",omitempty"`
}{
Opcode: route.Opcode,
Keyspace: route.Keyspace,
Query: route.Query,
FieldQuery: route.FieldQuery,
Vindex: vindexName,
Values: route.Values,
OrderBy: route.OrderBy,
TruncateColumnCount: route.TruncateColumnCount,
QueryTimeout: route.QueryTimeout,
Opcode: route.Opcode,
Keyspace: route.Keyspace,
Query: route.Query,
FieldQuery: route.FieldQuery,
Vindex: vindexName,
Values: route.Values,
OrderBy: route.OrderBy,
TruncateColumnCount: route.TruncateColumnCount,
QueryTimeout: route.QueryTimeout,
ScatterErrorsAsWarnings: route.ScatterErrorsAsWarnings,
}
return jsonutil.MarshalNoEscape(marshalRoute)
}
Expand Down Expand Up @@ -152,6 +158,10 @@ var routeName = map[RouteOpcode]string{
SelectDBA: "SelectDBA",
}

var (
partialSuccessScatterQueries = stats.NewCounter("PartialSuccessScatterQueries", "Count of partially successful scatter queries")
)

// MarshalJSON serializes the RouteOpcode as a JSON string.
// It's used for testing and diagnostics.
func (code RouteOpcode) MarshalJSON() ([]byte, error) {
Expand Down Expand Up @@ -209,9 +219,15 @@ func (route *Route) execute(vcursor VCursor, bindVars map[string]*querypb.BindVa
}

queries := getQueries(route.Query, bvs)
result, err := vcursor.ExecuteMultiShard(rss, queries, false /* isDML */, false /* autocommit */)
if err != nil {
return nil, err
result, errs := vcursor.ExecuteMultiShard(rss, queries, false /* isDML */, false /* autocommit */)

if errs != nil {
if route.ScatterErrorsAsWarnings {
partialSuccessScatterQueries.Add(1)
// fall through
} else {
return nil, vterrors.Aggregate(errs)
}
}
if len(route.OrderBy) == 0 {
return result, nil
Expand Down Expand Up @@ -421,12 +437,13 @@ func execAnyShard(vcursor VCursor, query string, bindVars map[string]*querypb.Bi

func execShard(vcursor VCursor, query string, bindVars map[string]*querypb.BindVariable, rs *srvtopo.ResolvedShard, isDML, canAutocommit bool) (*sqltypes.Result, error) {
autocommit := canAutocommit && vcursor.AutocommitApproval()
return vcursor.ExecuteMultiShard([]*srvtopo.ResolvedShard{rs}, []*querypb.BoundQuery{
result, errs := vcursor.ExecuteMultiShard([]*srvtopo.ResolvedShard{rs}, []*querypb.BoundQuery{
{
Sql: query,
BindVariables: bindVars,
},
}, isDML, autocommit)
return result, vterrors.Aggregate(errs)
}

func getQueries(query string, bvs []map[string]*querypb.BindVariable) []*querypb.BoundQuery {
Expand Down
91 changes: 91 additions & 0 deletions go/vt/vtgate/engine/route_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -782,6 +782,7 @@ func TestParamsFail(t *testing.T) {
}

func TestExecFail(t *testing.T) {
// Unsharded error
sel := &Route{
Opcode: SelectUnsharded,
Keyspace: &vindexes.Keyspace{
Expand All @@ -799,4 +800,94 @@ func TestExecFail(t *testing.T) {
vc.Rewind()
_, err = wrapStreamExecute(sel, vc, map[string]*querypb.BindVariable{}, false)
expectError(t, "sel.StreamExecute err", err, "result error")

// Scatter fails if one of N fails without ScatterErrorsAsWarnings
sel = &Route{
Opcode: SelectScatter,
Keyspace: &vindexes.Keyspace{
Name: "ks",
Sharded: true,
},
Query: "dummy_select",
FieldQuery: "dummy_select_field",
}

vc = &loggingVCursor{
shards: []string{"-20", "20-"},
results: []*sqltypes.Result{defaultSelectResult},
multiShardErrs: []error{
errors.New("result error -20"),
},
}
_, err = sel.Execute(vc, map[string]*querypb.BindVariable{}, false)
expectError(t, "sel.Execute err", err, "result error -20")
vc.ExpectLog(t, []string{
`ResolveDestinations ks [] Destinations:DestinationAllShards()`,
`ExecuteMultiShard ks.-20: dummy_select {} ks.20-: dummy_select {} false false`,
})

vc.Rewind()

// Scatter succeeds if all shards fail with ScatterErrorsAsWarnings
sel = &Route{
Opcode: SelectScatter,
Keyspace: &vindexes.Keyspace{
Name: "ks",
Sharded: true,
},
Query: "dummy_select",
FieldQuery: "dummy_select_field",
ScatterErrorsAsWarnings: true,
}

vc = &loggingVCursor{
shards: []string{"-20", "20-"},
results: []*sqltypes.Result{defaultSelectResult},
multiShardErrs: []error{
errors.New("result error -20"),
errors.New("result error 20-"),
},
}
_, err = sel.Execute(vc, map[string]*querypb.BindVariable{}, false)
if err != nil {
t.Errorf("unexpected ScatterErrorsAsWarnings error %v", err)
}
vc.ExpectLog(t, []string{
`ResolveDestinations ks [] Destinations:DestinationAllShards()`,
`ExecuteMultiShard ks.-20: dummy_select {} ks.20-: dummy_select {} false false`,
})

vc.Rewind()

// Scatter succeeds if one of N fails with ScatterErrorsAsWarnings
sel = &Route{
Opcode: SelectScatter,
Keyspace: &vindexes.Keyspace{
Name: "ks",
Sharded: true,
},
Query: "dummy_select",
FieldQuery: "dummy_select_field",
ScatterErrorsAsWarnings: true,
}

vc = &loggingVCursor{
shards: []string{"-20", "20-"},
results: []*sqltypes.Result{defaultSelectResult},
multiShardErrs: []error{
errors.New("result error -20"),
nil,
},
}
result, err := sel.Execute(vc, map[string]*querypb.BindVariable{}, false)
if err != nil {
t.Errorf("unexpected ScatterErrorsAsWarnings error %v", err)
}
vc.ExpectLog(t, []string{
`ResolveDestinations ks [] Destinations:DestinationAllShards()`,
`ExecuteMultiShard ks.-20: dummy_select {} ks.20-: dummy_select {} false false`,
})
expectResult(t, "sel.Execute", result, defaultSelectResult)

vc.Rewind()
}
3 changes: 2 additions & 1 deletion go/vt/vtgate/engine/update.go
Original file line number Diff line number Diff line change
Expand Up @@ -251,5 +251,6 @@ func (upd *Update) execUpdateByDestination(vcursor VCursor, bindVars map[string]
}
}
autocommit := (len(rss) == 1 || upd.MultiShardAutocommit) && vcursor.AutocommitApproval()
return vcursor.ExecuteMultiShard(rss, queries, true /* isDML */, autocommit)
result, errs := vcursor.ExecuteMultiShard(rss, queries, true /* isDML */, autocommit)
return result, vterrors.Aggregate(errs)
}
Loading