Skip to content

Commit 3905eea

Browse files
committed
feat: introduce QueryCopStoreLimiter to manage TiKV cop request concurrency
- Added QueryCopStoreLimiter to limit the number of in-flight cop requests per TiKV store within a single query. - Updated DistSQLContext to include QueryCopStoreLimiter. - Modified Select function to apply the QueryCopStoreLimiter when sending requests to TiKV. - Enhanced tests to validate the behavior of the new limiter, ensuring it correctly manages concurrent requests. - Adjusted related variables and configurations to support the new functionality. Signed-off-by: Zhigao TONG <tongzhigao@pingcap.com>
1 parent c473608 commit 3905eea

20 files changed

Lines changed: 686 additions & 26 deletions

File tree

pkg/distsql/BUILD.bazel

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@ go_test(
6868
embed = [":distsql"],
6969
flaky = True,
7070
race = "on",
71-
shard_count = 33,
71+
shard_count = 34,
7272
deps = [
7373
"//pkg/distsql/context",
7474
"//pkg/errctx",

pkg/distsql/context/context.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ type DistSQLContext struct {
6363
TiFlashQuerySpillRatio float64
6464
TiFlashHashJoinVersion string
6565

66+
QueryCopStoreLimiter *kv.QueryCopStoreLimiter
6667
DistSQLConcurrency int
6768
ReplicaReadType kv.ReplicaReadType
6869
WeakConsistency bool

pkg/distsql/context/context_test.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,7 @@ func TestContextDetach(t *testing.T) {
111111
"$.RunawayChecker",
112112
"$.RUConsumptionReporter",
113113
"$.ExecDetails",
114+
"$.QueryCopStoreLimiter",
114115
}))
115116

116117
staticObj := obj.Detach()

pkg/distsql/distsql.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,10 @@ func Select(ctx context.Context, dctx *distsqlctx.DistSQLContext, kvReq *kv.Requ
6060
r, ctx := tracing.StartRegionEx(ctx, "distsql.Select")
6161
defer r.End()
6262

63+
if kvReq.StoreType == kv.TiKV && dctx.QueryCopStoreLimiter != nil {
64+
kvReq.QueryCopStoreLimiter = dctx.QueryCopStoreLimiter
65+
}
66+
6367
// For testing purpose.
6468
if hook := ctx.Value("CheckSelectRequestHook"); hook != nil {
6569
hook.(func(*kv.Request))(kvReq)

pkg/distsql/distsql_test.go

Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,69 @@ func TestSelectWithRuntimeStats(t *testing.T) {
103103
require.NoError(t, response.Close())
104104
}
105105

106+
func TestSelectAppliesQueryCopStoreLimiter(t *testing.T) {
107+
sctx := newMockSessionContext()
108+
sctx.GetSessionVars().QueryCopStoreLimit = 3
109+
dctx := sctx.GetDistSQLCtx()
110+
require.NotNil(t, dctx.QueryCopStoreLimiter)
111+
require.Equal(t, 3, dctx.QueryCopStoreLimiter.Capacity())
112+
113+
colTypes := []*types.FieldType{types.NewFieldType(mysql.TypeLonglong)}
114+
buildRequest := func(storeType kv.StoreType) *kv.Request {
115+
request, err := (&RequestBuilder{}).SetKeyRanges(nil).
116+
SetDAGRequest(&tipb.DAGRequest{}).
117+
SetStoreType(storeType).
118+
SetFromSessionVars(DefaultDistSQLContext).
119+
SetMemTracker(memory.NewTracker(-1, -1)).
120+
Build()
121+
require.NoError(t, err)
122+
return request
123+
}
124+
checkRequest := func(check func(*kv.Request)) context.Context {
125+
return context.WithValue(context.TODO(), "CheckSelectRequestHook", func(req *kv.Request) {
126+
check(req)
127+
})
128+
}
129+
130+
request := buildRequest(kv.TiKV)
131+
response, err := Select(checkRequest(func(req *kv.Request) {
132+
require.Nil(t, req.CoprRequestRateLimit)
133+
require.Same(t, dctx.QueryCopStoreLimiter, req.QueryCopStoreLimiter)
134+
}), dctx, request, colTypes)
135+
require.NoError(t, err)
136+
require.NoError(t, response.Close())
137+
138+
request = buildRequest(kv.TiFlash)
139+
response, err = Select(checkRequest(func(req *kv.Request) {
140+
require.Nil(t, req.CoprRequestRateLimit)
141+
require.Nil(t, req.QueryCopStoreLimiter)
142+
}), dctx, request, colTypes)
143+
require.NoError(t, err)
144+
require.NoError(t, response.Close())
145+
146+
request = buildRequest(kv.TiKV)
147+
explicitRateLimit := kv.NewCoprRequestRateLimit(7)
148+
request.CoprRequestRateLimit = explicitRateLimit
149+
response, err = Select(checkRequest(func(req *kv.Request) {
150+
require.Same(t, explicitRateLimit, req.CoprRequestRateLimit)
151+
require.Same(t, dctx.QueryCopStoreLimiter, req.QueryCopStoreLimiter)
152+
require.False(t, req.CoprRequestRateLimit.Acquire(make(chan struct{})))
153+
req.CoprRequestRateLimit.Release()
154+
}), dctx, request, colTypes)
155+
require.NoError(t, err)
156+
require.NoError(t, response.Close())
157+
158+
dctx.QueryCopStoreLimiter = nil
159+
request = buildRequest(kv.TiKV)
160+
request.CoprRequestRateLimit = explicitRateLimit
161+
response, err = Select(checkRequest(func(req *kv.Request) {
162+
require.Same(t, explicitRateLimit, req.CoprRequestRateLimit)
163+
require.Nil(t, req.QueryCopStoreLimiter)
164+
}), dctx, request, colTypes)
165+
require.NoError(t, err)
166+
require.NoError(t, response.Close())
167+
}
168+
106169
func TestSelectResultRuntimeStats(t *testing.T) {
107170
stmtStats := execdetails.NewRuntimeStatsColl(nil)
108171
basic := stmtStats.GetBasicRuntimeStats(1, true)

pkg/distsql/request_builder.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -391,7 +391,7 @@ func (builder *RequestBuilder) SetConcurrency(concurrency int) *RequestBuilder {
391391
}
392392

393393
// SetCoprRequestRateLimit sets a shared in-flight cop request limiter for this request.
394-
func (builder *RequestBuilder) SetCoprRequestRateLimit(rateLimit *util.RateLimit) *RequestBuilder {
394+
func (builder *RequestBuilder) SetCoprRequestRateLimit(rateLimit kv.CoprRequestLimiter) *RequestBuilder {
395395
builder.Request.CoprRequestRateLimit = rateLimit
396396
return builder
397397
}

pkg/executor/distsql.go

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,6 @@ import (
6565
rangerctx "github.com/pingcap/tidb/pkg/util/ranger/context"
6666
"github.com/pingcap/tidb/pkg/util/size"
6767
"github.com/pingcap/tipb/go-tipb"
68-
tikvutil "github.com/tikv/client-go/v2/util"
6968
"go.uber.org/zap"
7069
)
7170

@@ -1022,7 +1021,7 @@ func getSelectResultInFlightCost(result distsql.SelectResult) int {
10221021
return inFlightCost
10231022
}
10241023

1025-
func getMergeSortSharedCoprRequestRateLimit(needMerge bool, distSQLConcurrency int) *tikvutil.RateLimit {
1024+
func getMergeSortSharedCoprRequestRateLimit(needMerge bool, distSQLConcurrency int) kv.CoprRequestLimiter {
10261025
if !needMerge {
10271026
return nil
10281027
}
@@ -1032,7 +1031,7 @@ func getMergeSortSharedCoprRequestRateLimit(needMerge bool, distSQLConcurrency i
10321031
if capacity < 1 {
10331032
capacity = 1
10341033
}
1035-
return tikvutil.NewRateLimit(2 * capacity)
1034+
return kv.NewCoprRequestRateLimit(2 * capacity)
10361035
}
10371036

10381037
func getMergeSortIndexScanConcurrency(needMerge bool, kvRangesCount int, distSQLConcurrency int) int {
@@ -1067,7 +1066,7 @@ func (e *IndexLookUpExecutor) buildIndexSelectResultForRange(
10671066
totalRanges int,
10681067
batchSize int,
10691068
indexScanConcurrency int,
1070-
sharedCoprRequestRateLimit *tikvutil.RateLimit,
1069+
sharedCoprRequestRateLimit kv.CoprRequestLimiter,
10711070
) (distsql.SelectResult, error) {
10721071
if tblScanIdxForRewritePartitionID >= 0 {
10731072
// We should set the TblScan's TableID to the partition physical ID to make sure

pkg/executor/set_test.go

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1447,6 +1447,7 @@ func TestSetConcurrency(t *testing.T) {
14471447
tk.MustQuery("select @@tidb_streamagg_concurrency;").Check(testkit.Rows(strconv.Itoa(vardef.DefTiDBStreamAggConcurrency)))
14481448
tk.MustQuery("select @@tidb_projection_concurrency;").Check(testkit.Rows(strconv.Itoa(vardef.ConcurrencyUnset)))
14491449
tk.MustQuery("select @@tidb_distsql_scan_concurrency;").Check(testkit.Rows(strconv.Itoa(vardef.DefDistSQLScanConcurrency)))
1450+
tk.MustQuery("select @@tidb_query_cop_store_limit;").Check(testkit.Rows(strconv.Itoa(vardef.DefTiDBQueryCopStoreLimit)))
14501451

14511452
tk.MustQuery("select @@tidb_index_serial_scan_concurrency;").Check(testkit.Rows(strconv.Itoa(vardef.DefIndexSerialScanConcurrency)))
14521453

@@ -1461,6 +1462,7 @@ func TestSetConcurrency(t *testing.T) {
14611462
require.Equal(t, vardef.DefTiDBStreamAggConcurrency, vars.StreamAggConcurrency())
14621463
require.Equal(t, vardef.DefExecutorConcurrency, vars.ProjectionConcurrency())
14631464
require.Equal(t, vardef.DefDistSQLScanConcurrency, vars.DistSQLScanConcurrency())
1465+
require.Equal(t, vardef.DefTiDBQueryCopStoreLimit, vars.QueryCopStoreLimit)
14641466

14651467
// test setting deprecated variables
14661468
warnTpl := "Warning 1287 '%s' is deprecated and will be removed in a future release. Please use tidb_executor_concurrency instead"
@@ -1499,6 +1501,10 @@ func TestSetConcurrency(t *testing.T) {
14991501
tk.MustQuery(fmt.Sprintf("select @@%s;", vardef.TiDBDistSQLScanConcurrency)).Check(testkit.Rows("1"))
15001502
require.Equal(t, 1, vars.DistSQLScanConcurrency())
15011503

1504+
tk.MustExec(fmt.Sprintf("set @@%s=8;", vardef.TiDBQueryCopStoreLimit))
1505+
tk.MustQuery(fmt.Sprintf("select @@%s;", vardef.TiDBQueryCopStoreLimit)).Check(testkit.Rows("8"))
1506+
require.Equal(t, 8, vars.QueryCopStoreLimit)
1507+
15021508
tk.MustExec("set @@tidb_index_serial_scan_concurrency=4")
15031509
tk.MustQuery("show warnings").Check(testkit.Rows("Warning 1287 The 'tidb_index_serial_scan_concurrency' variable is deprecated. Sequential scans follow 'tidb_executor_concurrency', and index statistics collection uses 'tidb_analyze_distsql_scan_concurrency'."))
15041510
tk.MustQuery("select @@tidb_index_serial_scan_concurrency;").Check(testkit.Rows("4"))

pkg/kv/BUILD.bazel

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -85,7 +85,7 @@ go_test(
8585
],
8686
embed = [":kv"],
8787
flaky = True,
88-
shard_count = 27,
88+
shard_count = 32,
8989
deps = [
9090
"//pkg/config/kerneltype",
9191
"//pkg/keyspace",

0 commit comments

Comments
 (0)