Skip to content

Commit 0a5fce0

Browse files
committed
logpuller: keep bootstrap requests flow-controlled
1 parent 07e9447 commit 0a5fce0

6 files changed

Lines changed: 237 additions & 60 deletions

File tree

logservice/logpuller/region_failure_handler.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -138,12 +138,12 @@ func (r *regionFailureHandler) handleError(ctx context.Context, errInfo regionEr
138138
}
139139
if innerErr.GetCongested() != nil {
140140
metricKvCongestedCounter.Inc()
141-
r.client.scheduleRegionRequest(ctx, errInfo.regionInfo, retryPriority)
141+
r.client.scheduleRegionRequest(ctx, errInfo.regionInfo, TaskLowPrior)
142142
return nil
143143
}
144144
if innerErr.GetServerIsBusy() != nil {
145145
metricKvIsBusyCounter.Inc()
146-
r.client.scheduleRegionRequest(ctx, errInfo.regionInfo, retryPriority)
146+
r.client.scheduleRegionRequest(ctx, errInfo.regionInfo, TaskLowPrior)
147147
return nil
148148
}
149149
if duplicated := innerErr.GetDuplicateRequest(); duplicated != nil {

logservice/logpuller/region_req_cache.go

Lines changed: 21 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -92,27 +92,32 @@ func newRequestCache(maxPendingCount int) *requestCache {
9292
}
9393

9494
// add adds a new region request to the cache
95-
// It blocks if pendingCount >= maxPendingCount until there's space or ctx is cancelled
95+
// Normal data requests are limited to maxPendingCount. Forced data requests can
96+
// use one additional slot, while stop/control requests keep their existing bypass.
9697
func (c *requestCache) add(ctx context.Context, region regionInfo, force bool) (bool, error) {
9798
start := time.Now()
9899
ticker := time.NewTicker(addReqRetryInterval)
99100
defer ticker.Stop()
100101
addReqRetryLimit := addReqRetryLimit
101102

102103
for {
103-
current := c.pendingCount.Load()
104-
if current < c.maxPendingCount || force {
104+
limit := c.maxPendingCount
105+
if force {
106+
limit++
107+
}
108+
if c.tryAcquireSlot(limit, region.isStopped()) {
105109
// Try to add the request
106110
req := newRegionReq(region)
107111
select {
108112
case <-ctx.Done():
113+
c.markDone()
109114
return false, ctx.Err()
110115
case c.pendingQueue <- req:
111-
c.pendingCount.Inc()
112116
cost := time.Since(start)
113117
metrics.SubscriptionClientAddRegionRequestDuration.Observe(cost.Seconds())
114118
return true, nil
115119
case <-ticker.C:
120+
c.markDone()
116121
addReqRetryLimit--
117122
if addReqRetryLimit <= 0 {
118123
return false, nil
@@ -137,6 +142,18 @@ func (c *requestCache) add(ctx context.Context, region regionInfo, force bool) (
137142
}
138143
}
139144

145+
func (c *requestCache) tryAcquireSlot(limit int64, bypassLimit bool) bool {
146+
for {
147+
current := c.pendingCount.Load()
148+
if !bypassLimit && current >= limit {
149+
return false
150+
}
151+
if c.pendingCount.CompareAndSwap(current, current+1) {
152+
return true
153+
}
154+
}
155+
}
156+
140157
// pop gets the next pending request.
141158
// Note: it doesn't change pendingCount. The slot acquired in add() should be released later
142159
// (e.g. resolve/markStopped/markDone).

logservice/logpuller/region_req_cache_test.go

Lines changed: 84 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import (
1919
"time"
2020

2121
"github.com/pingcap/ticdc/heartbeatpb"
22+
"github.com/pingcap/ticdc/logservice/logpuller/regionlock"
2223
"github.com/stretchr/testify/require"
2324
"github.com/tikv/client-go/v2/tikv"
2425
)
@@ -38,7 +39,9 @@ func createTestRegionInfo(subID SubscriptionID, regionID uint64) regionInfo {
3839
span: span,
3940
}
4041

41-
return newRegionInfo(verID, span, nil, subscribedSpan, false)
42+
region := newRegionInfo(verID, span, nil, subscribedSpan, false)
43+
region.lockedRangeState = &regionlock.LockedRangeState{}
44+
return region
4245
}
4346

4447
func TestRequestCacheAdd_NormalCase(t *testing.T) {
@@ -77,14 +80,7 @@ func TestRequestCacheAdd_ForceFlag(t *testing.T) {
7780
require.False(t, ok)
7881
require.NoError(t, err)
7982

80-
// With force=true, it should still fail because the channel is full
81-
// The force flag only bypasses the pendingCount check, not the channel capacity
82-
region3 := createTestRegionInfo(1, 3)
83-
ok, err = cache.add(ctx, region3, true)
84-
require.False(t, ok)
85-
require.NoError(t, err)
86-
87-
// consume the pending queue ann add with force
83+
// Move the normal request to sentRequests so the pending queue has room.
8884
req, err := cache.pop(ctx)
8985
require.NoError(t, err)
9086
require.NotNil(t, req)
@@ -93,15 +89,91 @@ func TestRequestCacheAdd_ForceFlag(t *testing.T) {
9389
cache.markSent(req)
9490
require.Equal(t, 1, cache.getPendingCount())
9591

92+
// A forced data request can use one extra slot.
93+
region3 := createTestRegionInfo(1, 3)
9694
ok, err = cache.add(ctx, region3, true)
9795
require.True(t, ok)
9896
require.NoError(t, err)
99-
// It is 2 since region1 is unresolved
10097
require.Equal(t, 2, cache.getPendingCount())
10198

102-
// resolve region1
103-
cache.resolve(region1.subscribedSpan.subID, region1.verID.GetID())
99+
// No additional forced data request can exceed the N+1 ceiling.
100+
req, err = cache.pop(ctx)
101+
require.NoError(t, err)
102+
cache.markSent(req)
103+
region4 := createTestRegionInfo(1, 4)
104+
ok, err = cache.add(ctx, region4, true)
105+
require.False(t, ok)
106+
require.NoError(t, err)
107+
require.Equal(t, 2, cache.getPendingCount())
108+
109+
// Stop/control requests keep their existing bypass and remain accounted.
110+
stopRegion := createTestRegionInfo(2, 5)
111+
stopRegion.lockedRangeState = nil
112+
ok, err = cache.add(ctx, stopRegion, true)
113+
require.True(t, ok)
114+
require.NoError(t, err)
115+
require.Equal(t, 3, cache.getPendingCount())
116+
117+
stopReq, err := cache.pop(ctx)
118+
require.NoError(t, err)
119+
cache.markSent(stopReq)
120+
cache.markStopped(stopReq.regionInfo.subscribedSpan.subID, stopReq.regionInfo.verID.GetID())
121+
require.Equal(t, 2, cache.getPendingCount())
122+
123+
require.True(t, cache.resolve(region1.subscribedSpan.subID, region1.verID.GetID()))
104124
require.Equal(t, 1, cache.getPendingCount())
125+
ok, err = cache.add(ctx, region4, true)
126+
require.True(t, ok)
127+
require.NoError(t, err)
128+
require.Equal(t, 2, cache.getPendingCount())
129+
}
130+
131+
func TestRequestCacheAddRollsBackReservedSlot(t *testing.T) {
132+
cache := newRequestCache(1)
133+
cache.pendingQueue <- newRegionReq(createTestRegionInfo(1, 1))
134+
135+
ctx, cancel := context.WithCancel(context.Background())
136+
cancel()
137+
ok, err := cache.add(ctx, createTestRegionInfo(1, 2), false)
138+
require.False(t, ok)
139+
require.ErrorIs(t, err, context.Canceled)
140+
require.Equal(t, 0, cache.getPendingCount())
141+
}
142+
143+
func TestRequestCacheConcurrentForcedAddsStayWithinCeiling(t *testing.T) {
144+
const normalLimit = 10
145+
cache := newRequestCache(normalLimit)
146+
ctx := context.Background()
147+
148+
for i := range normalLimit {
149+
ok, err := cache.add(ctx, createTestRegionInfo(1, uint64(i+1)), false)
150+
require.True(t, ok)
151+
require.NoError(t, err)
152+
}
153+
for range normalLimit {
154+
req, err := cache.pop(ctx)
155+
require.NoError(t, err)
156+
cache.markSent(req)
157+
}
158+
159+
const addCount = 20
160+
results := make(chan bool, addCount)
161+
for i := range addCount {
162+
go func(regionID uint64) {
163+
ok, err := cache.add(ctx, createTestRegionInfo(1, regionID), true)
164+
require.NoError(t, err)
165+
results <- ok
166+
}(uint64(normalLimit + i + 1))
167+
}
168+
169+
successes := 0
170+
for range addCount {
171+
if <-results {
172+
successes++
173+
}
174+
}
175+
require.Equal(t, 1, successes)
176+
require.Equal(t, normalLimit+1, cache.getPendingCount())
105177
}
106178

107179
func TestRequestCacheAdd_ContextCancellation(t *testing.T) {

logservice/logpuller/scan_priority_test.go

Lines changed: 76 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -141,11 +141,86 @@ func TestScanPriorityUsesRestoredRegionProgress(t *testing.T) {
141141
retryRegion := newRegionInfo(tikv.NewRegionVerID(1, 1, 2), rawSpan, nil, span, false)
142142
client.scheduleRegionRequest(context.Background(), retryRegion, TaskLowPrior)
143143
retryTask := popRegionPriorityTask(t, client.regionTaskQueue)
144-
require.Equal(t, TaskHighPrior, retryTask.taskType)
144+
require.Equal(t, TaskLowPrior, retryTask.taskType)
145145
require.Equal(t, cdcpb.ScanPriority_SCAN_PRIORITY_HIGH, retryTask.GetRegionInfo().scanPriority)
146146
require.False(t, span.priorityPolicy.everCaughtUp.Load())
147147
}
148148

149+
func TestScheduleRegionRequestSeparatesRemoteAndLocalPriority(t *testing.T) {
150+
currentTime := time.Date(2026, time.June, 27, 12, 0, 0, 0, time.UTC)
151+
currentTs := oracle.GoTimeToTS(currentTime)
152+
153+
for _, tc := range []struct {
154+
name string
155+
priorRemote cdcpb.ScanPriority
156+
inheritedLocal TaskType
157+
startTs uint64
158+
expectedRemote cdcpb.ScanPriority
159+
expectedLocal TaskType
160+
}{
161+
{
162+
name: "busy retry preserves remote high",
163+
priorRemote: cdcpb.ScanPriority_SCAN_PRIORITY_HIGH,
164+
inheritedLocal: TaskLowPrior,
165+
startTs: oracle.GoTimeToTS(currentTime.Add(-time.Hour)),
166+
expectedRemote: cdcpb.ScanPriority_SCAN_PRIORITY_HIGH,
167+
expectedLocal: TaskLowPrior,
168+
},
169+
{
170+
name: "repair is high locally and remotely",
171+
priorRemote: cdcpb.ScanPriority_SCAN_PRIORITY_LOW,
172+
inheritedLocal: TaskHighPrior,
173+
startTs: oracle.GoTimeToTS(currentTime.Add(-time.Hour)),
174+
expectedRemote: cdcpb.ScanPriority_SCAN_PRIORITY_HIGH,
175+
expectedLocal: TaskHighPrior,
176+
},
177+
{
178+
name: "recent bootstrap is only high remotely",
179+
priorRemote: cdcpb.ScanPriority_SCAN_PRIORITY_LOW,
180+
inheritedLocal: TaskLowPrior,
181+
startTs: oracle.GoTimeToTS(currentTime.Add(-time.Minute)),
182+
expectedRemote: cdcpb.ScanPriority_SCAN_PRIORITY_HIGH,
183+
expectedLocal: TaskLowPrior,
184+
},
185+
{
186+
name: "old bootstrap stays low",
187+
priorRemote: cdcpb.ScanPriority_SCAN_PRIORITY_LOW,
188+
inheritedLocal: TaskLowPrior,
189+
startTs: oracle.GoTimeToTS(currentTime.Add(-time.Hour)),
190+
expectedRemote: cdcpb.ScanPriority_SCAN_PRIORITY_LOW,
191+
expectedLocal: TaskLowPrior,
192+
},
193+
} {
194+
t.Run(tc.name, func(t *testing.T) {
195+
pdClock := pdutil.NewClock4Test()
196+
pdClock.(*pdutil.Clock4Test).SetTS(currentTs)
197+
client := &subscriptionClient{
198+
pdClock: pdClock,
199+
regionTaskQueue: priorityqueue.New[PriorityTask](),
200+
}
201+
rawSpan := heartbeatpb.TableSpan{
202+
TableID: 1,
203+
StartKey: []byte("a"),
204+
EndKey: []byte("z"),
205+
}
206+
span := &subscribedSpan{
207+
subID: SubscriptionID(1),
208+
span: rawSpan,
209+
startTs: tc.startTs,
210+
rangeLock: regionlock.NewRangeLock(1, rawSpan.StartKey, rawSpan.EndKey, tc.startTs),
211+
priorityPolicy: newScanPriorityPolicy(pdClock, 30*time.Minute),
212+
}
213+
region := newRegionInfo(tikv.NewRegionVerID(1, 1, 1), rawSpan, nil, span, false)
214+
region.scanPriority = tc.priorRemote
215+
216+
client.scheduleRegionRequest(context.Background(), region, tc.inheritedLocal)
217+
task := popRegionPriorityTask(t, client.regionTaskQueue)
218+
require.Equal(t, tc.expectedLocal, task.taskType)
219+
require.Equal(t, tc.expectedRemote, task.GetRegionInfo().scanPriority)
220+
})
221+
}
222+
}
223+
149224
func popRegionPriorityTask(
150225
t *testing.T,
151226
queue *priorityqueue.PriorityQueue[PriorityTask],

logservice/logpuller/subscription_client.go

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import (
1919
"sync/atomic"
2020
"time"
2121

22+
"github.com/pingcap/kvproto/pkg/cdcpb"
2223
"github.com/pingcap/kvproto/pkg/metapb"
2324
"github.com/pingcap/log"
2425
"github.com/pingcap/ticdc/heartbeatpb"
@@ -652,13 +653,17 @@ func (s *subscriptionClient) scheduleRegionRequest(
652653
case regionlock.LockRangeStatusSuccess:
653654
region.lockedRangeState = lockRangeResult.LockedRangeState
654655
currentTs := s.pdClock.CurrentTS()
655-
priority := region.subscribedSpan.priorityPolicy.resolve(
656-
inheritedPriority,
656+
remoteBase := inheritedPriority
657+
if region.scanPriority == cdcpb.ScanPriority_SCAN_PRIORITY_HIGH {
658+
remoteBase = TaskHighPrior
659+
}
660+
remotePriority := region.subscribedSpan.priorityPolicy.resolve(
661+
remoteBase,
657662
region.resolvedTs(),
658663
oracle.GetTimeFromTS(currentTs),
659664
)
660-
region.scanPriority = priority.scanPriority()
661-
s.regionTaskQueue.Push(NewRegionPriorityTask(priority, region, currentTs))
665+
region.scanPriority = remotePriority.scanPriority()
666+
s.regionTaskQueue.Push(NewRegionPriorityTask(inheritedPriority, region, currentTs))
662667
if log.GetLevel() <= zapcore.DebugLevel {
663668
log.Debug("cdc region scan task enqueued",
664669
zap.Uint64("subscriptionID", uint64(region.subscribedSpan.subID)),
@@ -667,7 +672,7 @@ func (s *subscriptionClient) scheduleRegionRequest(
667672
zap.Uint64("regionID", region.verID.GetID()),
668673
zap.Uint64("regionEpochVersion", region.verID.GetVer()),
669674
zap.Uint64("regionEpochConfVer", region.verID.GetConfVer()),
670-
zap.String("priority", priority.String()),
675+
zap.String("priority", inheritedPriority.String()),
671676
zap.String("scanPriority", region.scanPriority.String()),
672677
zap.String("span", common.FormatTableSpan(&region.span)))
673678
}

0 commit comments

Comments
 (0)