Skip to content

Commit 890ae70

Browse files
committed
refactor: rename CoprRequestRateLimit to CoprRequestLimiter for consistency
Signed-off-by: Zhigao TONG <tongzhigao@pingcap.com>
1 parent a8e8f3c commit 890ae70

7 files changed

Lines changed: 43 additions & 46 deletions

File tree

pkg/distsql/distsql_test.go

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -129,37 +129,37 @@ func TestSelectAppliesQueryCopStoreLimiter(t *testing.T) {
129129

130130
request := buildRequest(kv.TiKV)
131131
response, err := Select(checkRequest(func(req *kv.Request) {
132-
require.Nil(t, req.CoprRequestRateLimit)
132+
require.Nil(t, req.CoprRequestLimiter)
133133
require.Same(t, dctx.QueryCopStoreLimiter, req.QueryCopStoreLimiter)
134134
}), dctx, request, colTypes)
135135
require.NoError(t, err)
136136
require.NoError(t, response.Close())
137137

138138
request = buildRequest(kv.TiFlash)
139139
response, err = Select(checkRequest(func(req *kv.Request) {
140-
require.Nil(t, req.CoprRequestRateLimit)
140+
require.Nil(t, req.CoprRequestLimiter)
141141
require.Nil(t, req.QueryCopStoreLimiter)
142142
}), dctx, request, colTypes)
143143
require.NoError(t, err)
144144
require.NoError(t, response.Close())
145145

146146
request = buildRequest(kv.TiKV)
147-
explicitRateLimit := kv.NewCoprRequestRateLimit(7)
148-
request.CoprRequestRateLimit = explicitRateLimit
147+
explicitRateLimit := kv.NewCoprRequestLimiter(7)
148+
request.CoprRequestLimiter = explicitRateLimit
149149
response, err = Select(checkRequest(func(req *kv.Request) {
150-
require.Same(t, explicitRateLimit, req.CoprRequestRateLimit)
150+
require.Same(t, explicitRateLimit, req.CoprRequestLimiter)
151151
require.Same(t, dctx.QueryCopStoreLimiter, req.QueryCopStoreLimiter)
152-
require.False(t, req.CoprRequestRateLimit.Acquire(make(chan struct{})))
153-
req.CoprRequestRateLimit.Release()
152+
require.False(t, req.CoprRequestLimiter.Acquire(make(chan struct{})))
153+
req.CoprRequestLimiter.Release()
154154
}), dctx, request, colTypes)
155155
require.NoError(t, err)
156156
require.NoError(t, response.Close())
157157

158158
dctx.QueryCopStoreLimiter = nil
159159
request = buildRequest(kv.TiKV)
160-
request.CoprRequestRateLimit = explicitRateLimit
160+
request.CoprRequestLimiter = explicitRateLimit
161161
response, err = Select(checkRequest(func(req *kv.Request) {
162-
require.Same(t, explicitRateLimit, req.CoprRequestRateLimit)
162+
require.Same(t, explicitRateLimit, req.CoprRequestLimiter)
163163
require.Nil(t, req.QueryCopStoreLimiter)
164164
}), dctx, request, colTypes)
165165
require.NoError(t, err)

pkg/distsql/request_builder.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -390,9 +390,9 @@ func (builder *RequestBuilder) SetConcurrency(concurrency int) *RequestBuilder {
390390
return builder
391391
}
392392

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

pkg/executor/distsql.go

Lines changed: 7 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -935,7 +935,7 @@ func (e *IndexLookUpExecutor) startIndexWorker(ctx context.Context, initBatchSiz
935935
return
936936
}
937937

938-
sharedCoprRequestRateLimit := getMergeSortSharedCoprRequestRateLimit(needMerge, e.dctx.DistSQLConcurrency)
938+
sharedCoprRequestLimiter := getMergeSortSharedCoprRequestLimiter(needMerge, e.dctx.DistSQLConcurrency)
939939
mergeSortIndexScanConcurrency := getMergeSortIndexScanConcurrency(needMerge, len(kvRanges), e.dctx.DistSQLConcurrency)
940940
results := make([]distsql.SelectResult, 0, len(kvRanges))
941941
for idx := range kvRanges {
@@ -960,7 +960,7 @@ func (e *IndexLookUpExecutor) startIndexWorker(ctx context.Context, initBatchSiz
960960
len(kvRanges),
961961
worker.batchSize,
962962
mergeSortIndexScanConcurrency,
963-
sharedCoprRequestRateLimit,
963+
sharedCoprRequestLimiter,
964964
)
965965
if err != nil {
966966
for _, r := range results {
@@ -1021,17 +1021,14 @@ func getSelectResultInFlightCost(result distsql.SelectResult) int {
10211021
return inFlightCost
10221022
}
10231023

1024-
func getMergeSortSharedCoprRequestRateLimit(needMerge bool, distSQLConcurrency int) kv.CoprRequestLimiter {
1024+
func getMergeSortSharedCoprRequestLimiter(needMerge bool, distSQLConcurrency int) kv.CoprRequestLimiter {
10251025
if !needMerge {
10261026
return nil
10271027
}
10281028
// Use a shared limiter to bound aggregate in-flight cop requests across
10291029
// all partitions in merge-sort mode.
1030-
capacity := distSQLConcurrency
1031-
if capacity < 1 {
1032-
capacity = 1
1033-
}
1034-
return kv.NewCoprRequestRateLimit(2 * capacity)
1030+
capacity := max(distSQLConcurrency, 1)
1031+
return kv.NewCoprRequestLimiter(2 * capacity)
10351032
}
10361033

10371034
func getMergeSortIndexScanConcurrency(needMerge bool, kvRangesCount int, distSQLConcurrency int) int {
@@ -1066,7 +1063,7 @@ func (e *IndexLookUpExecutor) buildIndexSelectResultForRange(
10661063
totalRanges int,
10671064
batchSize int,
10681065
indexScanConcurrency int,
1069-
sharedCoprRequestRateLimit kv.CoprRequestLimiter,
1066+
sharedCoprRequestLimiter kv.CoprRequestLimiter,
10701067
) (distsql.SelectResult, error) {
10711068
if tblScanIdxForRewritePartitionID >= 0 {
10721069
// We should set the TblScan's TableID to the partition physical ID to make sure
@@ -1092,7 +1089,7 @@ func (e *IndexLookUpExecutor) buildIndexSelectResultForRange(
10921089
SetClosestReplicaReadAdjuster(newClosestReadAdjuster(e.dctx, &builder.Request, e.idxNetDataSize/float64(totalRanges))).
10931090
SetMemTracker(tracker).
10941091
SetConnIDAndConnAlias(e.dctx.ConnectionID, e.dctx.SessionAlias).
1095-
SetCoprRequestRateLimit(sharedCoprRequestRateLimit)
1092+
SetCoprRequestLimiter(sharedCoprRequestLimiter)
10961093

10971094
if e.indexLookUpPushDown {
10981095
// Paging and Cop-cache is not supported in index lookup push down.

pkg/kv/kv.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -446,8 +446,8 @@ type coprRequestWaiter struct {
446446
admitted bool
447447
}
448448

449-
// NewCoprRequestRateLimit creates a cop request limiter with capacity n.
450-
func NewCoprRequestRateLimit(n int) CoprRequestLimiter {
449+
// NewCoprRequestLimiter creates a cop request limiter with capacity n.
450+
func NewCoprRequestLimiter(n int) CoprRequestLimiter {
451451
if n <= 0 {
452452
return nil
453453
}
@@ -594,7 +594,7 @@ func (l *QueryCopStoreLimiter) GetStoreLimiter(storeID uint64) CoprRequestLimite
594594
if limiter, ok := l.stores.Load(storeID); ok {
595595
return limiter.(CoprRequestLimiter)
596596
}
597-
newLimiter := NewCoprRequestRateLimit(l.limit)
597+
newLimiter := NewCoprRequestLimiter(l.limit)
598598
limiter, _ := l.stores.LoadOrStore(storeID, newLimiter)
599599
return limiter.(CoprRequestLimiter)
600600
}
@@ -789,10 +789,10 @@ type Request struct {
789789
// ResponseIterator.Next is called. If concurrency is greater than 1, the request will be
790790
// sent to multiple storage units concurrently.
791791
Concurrency int
792-
// CoprRequestRateLimit, if not nil, is used as the shared in-flight request
792+
// CoprRequestLimiter, if not nil, is used as the shared in-flight request
793793
// limiter for all cop iterators created from this request. The token lifecycle
794794
// is tied to request send/response receive instead of result consumption.
795-
CoprRequestRateLimit CoprRequestLimiter
795+
CoprRequestLimiter CoprRequestLimiter
796796
// QueryCopStoreLimiter, if not nil, limits in-flight cop request attempts
797797
// per TiKV store within this request's statement.
798798
QueryCopStoreLimiter *QueryCopStoreLimiter

pkg/kv/kv_test.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ func genRandHex(length int) []byte {
3939
}
4040

4141
func TestCoprRequestLimiterWaitsUntilRelease(t *testing.T) {
42-
limiter := NewCoprRequestRateLimit(1)
42+
limiter := NewCoprRequestLimiter(1)
4343
done := make(chan struct{})
4444
require.False(t, limiter.Acquire(done))
4545

@@ -70,7 +70,7 @@ func TestCoprRequestLimiterWaitsUntilRelease(t *testing.T) {
7070
}
7171

7272
func TestCoprRequestLimiterAcquireCanBeCanceled(t *testing.T) {
73-
limiter := NewCoprRequestRateLimit(1)
73+
limiter := NewCoprRequestLimiter(1)
7474
require.False(t, limiter.Acquire(make(chan struct{})))
7575

7676
done := make(chan struct{})
@@ -94,15 +94,15 @@ func TestCoprRequestLimiterAcquireCanBeCanceled(t *testing.T) {
9494
}
9595

9696
func TestCoprRequestLimiterRedundantReleasePanics(t *testing.T) {
97-
limiter := NewCoprRequestRateLimit(1)
97+
limiter := NewCoprRequestLimiter(1)
9898
require.Panics(t, func() {
9999
limiter.Release()
100100
})
101101
}
102102

103103
func TestCoprRequestLimiterConcurrentAcquireRelease(t *testing.T) {
104104
const capacity = int64(3)
105-
limiter := NewCoprRequestRateLimit(int(capacity))
105+
limiter := NewCoprRequestLimiter(int(capacity))
106106
done := make(chan struct{})
107107
var active atomic.Int64
108108
var maxActive atomic.Int64

pkg/store/copr/copr_test/coprocessor_test.go

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -139,7 +139,7 @@ func TestBuildCopIteratorWithRowCountHint(t *testing.T) {
139139
require.Equal(t, rateLimit.GetCapacity(), 4)
140140
}
141141

142-
func TestBuildCopIteratorWithSharedRequestRateLimit(t *testing.T) {
142+
func TestBuildCopIteratorWithSharedRequestLimiter(t *testing.T) {
143143
store, err := mockstore.NewMockStore()
144144
require.NoError(t, err)
145145
defer require.NoError(t, store.Close())
@@ -161,18 +161,18 @@ func TestBuildCopIteratorWithSharedRequestRateLimit(t *testing.T) {
161161

162162
for _, tc := range testCases {
163163
t.Run(tc.name, func(t *testing.T) {
164-
shared := kv.NewCoprRequestRateLimit(7)
164+
shared := kv.NewCoprRequestLimiter(7)
165165
req := &kv.Request{
166-
Tp: kv.ReqTypeDAG,
167-
KeyRanges: kv.NewNonPartitionedKeyRanges(ranges),
168-
Concurrency: 15,
169-
KeepOrder: tc.keepOrder,
170-
CoprRequestRateLimit: shared,
166+
Tp: kv.ReqTypeDAG,
167+
KeyRanges: kv.NewNonPartitionedKeyRanges(ranges),
168+
Concurrency: 15,
169+
KeepOrder: tc.keepOrder,
170+
CoprRequestLimiter: shared,
171171
}
172172
it, errRes := copClient.BuildCopIterator(ctx, req, vars, opt)
173173
require.Nil(t, errRes)
174-
require.Same(t, shared, it.GetRequestRateLimit())
175-
require.Equal(t, 7, it.GetRequestRateLimit().Capacity())
174+
require.Same(t, shared, it.GetRequestLimiter())
175+
require.Equal(t, 7, it.GetRequestLimiter().Capacity())
176176
})
177177
}
178178
}

pkg/store/copr/coprocessor.go

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1348,9 +1348,9 @@ func (it *copIterator) GetSendRate() *util.RateLimit {
13481348
return it.sendRate
13491349
}
13501350

1351-
// GetRequestRateLimit returns the shared request rate-limit object.
1352-
func (it *copIterator) GetRequestRateLimit() kv.CoprRequestLimiter {
1353-
return it.req.CoprRequestRateLimit
1351+
// GetRequestLimiter returns the shared request rate-limit object.
1352+
func (it *copIterator) GetRequestLimiter() kv.CoprRequestLimiter {
1353+
return it.req.CoprRequestLimiter
13541354
}
13551355

13561356
// GetTasks returns the built tasks.
@@ -1814,7 +1814,7 @@ func (worker *copIteratorWorker) handleTaskOnce(bo *Backoffer, task *copTask) (*
18141814
if exit {
18151815
return nil, nil
18161816
}
1817-
releaseRequestRateLimit, exit := acquireCoprRequestLimiter(worker.req.CoprRequestRateLimit, worker.finishCh)
1817+
releaseCoprRequestLimiter, exit := acquireCoprRequestLimiter(worker.req.CoprRequestLimiter, worker.finishCh)
18181818
if exit {
18191819
if releaseQueryCopStoreLimiter != nil {
18201820
releaseQueryCopStoreLimiter()
@@ -1826,8 +1826,8 @@ func (worker *copIteratorWorker) handleTaskOnce(bo *Backoffer, task *copTask) (*
18261826
// remaining panic-safe.
18271827
resp, rpcCtx, storeAddr, err := func() (*tikvrpc.Response, *tikv.RPCContext, string, error) {
18281828
defer func() {
1829-
if releaseRequestRateLimit != nil {
1830-
releaseRequestRateLimit()
1829+
if releaseCoprRequestLimiter != nil {
1830+
releaseCoprRequestLimiter()
18311831
}
18321832
if releaseQueryCopStoreLimiter != nil {
18331833
releaseQueryCopStoreLimiter()

0 commit comments

Comments
 (0)