Skip to content

Commit af7312b

Browse files
committed
ddl/server: add async force merge for empty table-id regions (#684)
Signed-off-by: HunDunDM <hundundm@gmail.com>
1 parent d6f385e commit af7312b

22 files changed

Lines changed: 1644 additions & 30 deletions

ddl/BUILD.bazel

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@ go_library(
3535
"multi_schema_change.go",
3636
"options.go",
3737
"partition.go",
38+
"pkdb_force_merge.go",
3839
"placement_policy.go",
3940
"reorg.go",
4041
"rollingback.go",
@@ -186,6 +187,7 @@ go_test(
186187
"multi_schema_change_test.go",
187188
"options_test.go",
188189
"partition_test.go",
190+
"pkdb_force_merge_test.go",
189191
"placement_policy_ddl_test.go",
190192
"placement_policy_test.go",
191193
"placement_sql_test.go",

ddl/delete_range.go

Lines changed: 38 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
"github.com/pingcap/errors"
2626
"github.com/pingcap/kvproto/pkg/kvrpcpb"
2727
"github.com/pingcap/tidb/ddl/util"
28+
"github.com/pingcap/tidb/domain/infosync"
2829
"github.com/pingcap/tidb/kv"
2930
"github.com/pingcap/tidb/parser/model"
3031
"github.com/pingcap/tidb/parser/terror"
@@ -156,51 +157,61 @@ func (dr *delRange) clear() {
156157
func (dr *delRange) startEmulator() {
157158
defer dr.wait.Done()
158159
logutil.BgLogger().Info("[ddl] start delRange emulator")
160+
ctx, cancel := context.WithCancel(context.Background())
161+
defer cancel()
162+
go func() {
163+
select {
164+
case <-dr.quitCh:
165+
cancel()
166+
case <-ctx.Done():
167+
}
168+
}()
159169
for {
160170
select {
161171
case <-dr.emulatorCh:
162172
case <-dr.quitCh:
163173
return
164174
}
165175
if util.IsEmulatorGCEnable() {
166-
err := dr.doDelRangeWork()
176+
err := dr.doDelRangeWork(ctx)
167177
terror.Log(errors.Trace(err))
168178
}
169179
}
170180
}
171181

172-
func (dr *delRange) doDelRangeWork() error {
182+
func (dr *delRange) doDelRangeWork(ctx context.Context) error {
173183
sctx, err := dr.sessPool.get()
174184
if err != nil {
175185
logutil.BgLogger().Error("[ddl] delRange emulator get session failed", zap.Error(err))
176186
return errors.Trace(err)
177187
}
178188
defer dr.sessPool.put(sctx)
179189

180-
ctx := kv.WithInternalSourceType(context.Background(), kv.InternalTxnDDL)
181-
ranges, err := util.LoadDeleteRanges(ctx, sctx, math.MaxInt64)
190+
txnCtx := kv.WithInternalSourceType(ctx, kv.InternalTxnDDL)
191+
ranges, err := util.LoadDeleteRanges(txnCtx, sctx, math.MaxInt64)
182192
if err != nil {
183193
logutil.BgLogger().Error("[ddl] delRange emulator load tasks failed", zap.Error(err))
184194
return errors.Trace(err)
185195
}
186196

197+
gcForceMergeTableCache := make(map[int64]struct{}, len(ranges))
187198
for _, r := range ranges {
188-
if err := dr.doTask(sctx, r); err != nil {
199+
if err := dr.doTask(ctx, sctx, r, gcForceMergeTableCache); err != nil {
189200
logutil.BgLogger().Error("[ddl] delRange emulator do task failed", zap.Error(err))
190201
return errors.Trace(err)
191202
}
192203
}
193204
return nil
194205
}
195206

196-
func (dr *delRange) doTask(sctx sessionctx.Context, r util.DelRangeTask) error {
207+
func (dr *delRange) doTask(ctx context.Context, sctx sessionctx.Context, r util.DelRangeTask, gcForceMergeTableCache map[int64]struct{}) error {
197208
var oldStartKey, newStartKey kv.Key
198209
oldStartKey = r.StartKey
199210
for {
200211
finish := true
201212
dr.keys = dr.keys[:0]
202-
ctx := kv.WithInternalSourceType(context.Background(), kv.InternalTxnDDL)
203-
err := kv.RunInNewTxn(ctx, dr.store, false, func(ctx context.Context, txn kv.Transaction) error {
213+
txnCtx := kv.WithInternalSourceType(ctx, kv.InternalTxnDDL)
214+
err := kv.RunInNewTxn(txnCtx, dr.store, false, func(ctx context.Context, txn kv.Transaction) error {
204215
if topsqlstate.TopSQLEnabled() {
205216
// Only when TiDB run without PD(use unistore as storage for test) will run into here, so just set a mock internal resource tagger.
206217
txn.SetOption(kv.ResourceGroupTagger, util.GetInternalResourceGroupTaggerForTopSQL())
@@ -241,6 +252,9 @@ func (dr *delRange) doTask(sctx sessionctx.Context, r util.DelRangeTask) error {
241252
logutil.BgLogger().Error("[ddl] delRange emulator complete task failed", zap.Error(err))
242253
return errors.Trace(err)
243254
}
255+
if err := dr.reportForceMergeRanges(ctx, sctx, r, gcForceMergeTableCache); err != nil {
256+
logutil.BgLogger().Error("[ddl] delRange emulator report force merge failed", zap.Error(err))
257+
}
244258
startKey, endKey := r.Range()
245259
logutil.BgLogger().Info("[ddl] delRange emulator complete task",
246260
zap.Int64("jobID", r.JobID),
@@ -257,6 +271,22 @@ func (dr *delRange) doTask(sctx sessionctx.Context, r util.DelRangeTask) error {
257271
return nil
258272
}
259273

274+
func (dr *delRange) reportForceMergeRanges(ctx context.Context, sctx sessionctx.Context, r util.DelRangeTask, gcForceMergeTableCache map[int64]struct{}) error {
275+
historyJob, err := GetHistoryJobByID(sctx, r.JobID)
276+
if err != nil {
277+
return err
278+
}
279+
if historyJob == nil {
280+
return errors.Errorf("ddl job %d not found", r.JobID)
281+
}
282+
283+
forceMergeRanges := GetForceMergeRangesForGCDeleteRange(historyJob, r, gcForceMergeTableCache)
284+
if len(forceMergeRanges) == 0 {
285+
return nil
286+
}
287+
return infosync.AddForceMergeRanges(ctx, forceMergeRanges)
288+
}
289+
260290
// insertJobIntoDeleteRangeTable parses the job into delete-range arguments,
261291
// and inserts a new record into gc_delete_range table. The primary key is
262292
// (job ID, element ID), so we ignore key conflict error.

ddl/pkdb_force_merge.go

Lines changed: 256 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,256 @@
1+
// Copyright 2026 PingCAP, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package ddl
16+
17+
import (
18+
"sort"
19+
20+
ddlutil "github.com/pingcap/tidb/ddl/util"
21+
"github.com/pingcap/tidb/domain/infosync"
22+
"github.com/pingcap/tidb/infoschema"
23+
"github.com/pingcap/tidb/meta/autoid"
24+
"github.com/pingcap/tidb/parser/model"
25+
"github.com/pingcap/tidb/sessionctx/variable"
26+
"github.com/pingcap/tidb/tablecodec"
27+
)
28+
29+
// GetForceMergeRangesForGCDeleteRange builds the PD force-merge ranges that
30+
// should be reported for one gc_delete_range task.
31+
func GetForceMergeRangesForGCDeleteRange(
32+
historyJob *model.Job,
33+
dr ddlutil.DelRangeTask,
34+
gcLogicalTableCache map[int64]struct{},
35+
) []infosync.ForceMergeKeyRange {
36+
if !variable.EnableDropTableForceMerge.Load() || historyJob == nil {
37+
return nil
38+
}
39+
if !supportForceMergeForGCDeleteRange(historyJob.Type) {
40+
return nil
41+
}
42+
43+
var tblInfo *model.TableInfo
44+
if historyJob.BinlogInfo != nil {
45+
tblInfo = historyJob.BinlogInfo.TableInfo
46+
}
47+
if shouldSkipForceMerge(historyJob, tblInfo) {
48+
return nil
49+
}
50+
51+
// Reuse the gc_delete_range task range directly for the physical table/partition
52+
// that has just been deleted.
53+
ranges := []infosync.ForceMergeKeyRange{buildForceMergeKeyRange(dr.StartKey, dr.EndKey)}
54+
if shouldReportLogicalTableRange(historyJob, tblInfo, gcLogicalTableCache) {
55+
// DROP/TRUNCATE on a partitioned table only creates gc_delete_range tasks for
56+
// physical partitions. Append the logical table range once as well so PD can
57+
// force merge the partitioned table's global index keyspace.
58+
ranges = append(ranges, buildForceMergeTableKeyRange(historyJob.TableID))
59+
}
60+
return ranges
61+
}
62+
63+
func supportForceMergeForGCDeleteRange(actionType model.ActionType) bool {
64+
switch actionType {
65+
case model.ActionDropSchema,
66+
model.ActionDropTable, model.ActionTruncateTable,
67+
model.ActionDropTablePartition, model.ActionTruncateTablePartition:
68+
return true
69+
}
70+
return false
71+
}
72+
73+
func shouldSkipForceMerge(historyJob *model.Job, tblInfo *model.TableInfo) bool {
74+
if historyJob == nil {
75+
return true
76+
}
77+
78+
tableName := historyJob.TableName
79+
if tblInfo != nil && tableName == "" {
80+
tableName = tblInfo.Name.L
81+
}
82+
return shouldSkipForceMergeTable(historyJob.SchemaName, tableName, tblInfo)
83+
}
84+
85+
func shouldSkipForceMergeTable(schemaName string, tableName string, tblInfo *model.TableInfo) bool {
86+
if tblInfo != nil {
87+
if !tblInfo.IsBaseTable() {
88+
return true
89+
}
90+
if tblInfo.TempTableType != model.TempTableNone {
91+
return true
92+
}
93+
if tblInfo.ID&autoid.SystemSchemaIDFlag != 0 {
94+
return true
95+
}
96+
}
97+
return isSystemTable(schemaName, tableName)
98+
}
99+
100+
func shouldReportLogicalTableRange(
101+
historyJob *model.Job,
102+
tblInfo *model.TableInfo,
103+
gcLogicalTableCache map[int64]struct{},
104+
) bool {
105+
// Partition-level gc_delete_range tasks do not cover the logical table ID range.
106+
// That range must be reported separately for DROP/TRUNCATE of partitioned tables
107+
// because global indexes are stored under the logical table ID prefix.
108+
if historyJob == nil || historyJob.TableID <= 0 || tblInfo == nil || tblInfo.GetPartitionInfo() == nil {
109+
return false
110+
}
111+
if historyJob.Type != model.ActionDropTable && historyJob.Type != model.ActionTruncateTable {
112+
return false
113+
}
114+
if gcLogicalTableCache != nil {
115+
if _, ok := gcLogicalTableCache[historyJob.TableID]; ok {
116+
return false
117+
}
118+
gcLogicalTableCache[historyJob.TableID] = struct{}{}
119+
}
120+
return true
121+
}
122+
123+
func buildForceMergeKeyRange(startKey []byte, endKey []byte) infosync.ForceMergeKeyRange {
124+
return infosync.ForceMergeKeyRange{
125+
StartKey: append([]byte(nil), startKey...),
126+
EndKey: append([]byte(nil), endKey...),
127+
}
128+
}
129+
130+
// GetMergeEmptyRegionsKeyRanges scans the current infoschema in
131+
// [minTableID, maxTableID) and returns the missing table ID keyspace ranges.
132+
func GetMergeEmptyRegionsKeyRanges(
133+
is infoschema.InfoSchema,
134+
minTableID int64,
135+
) (maxTableID int64, ranges []infosync.ForceMergeKeyRange) {
136+
if is == nil {
137+
return 0, nil
138+
}
139+
if minTableID < 1 {
140+
minTableID = 1
141+
}
142+
143+
occupiedTableIDs, maxTableID := collectMergeEmptyRegionsOccupiedTableIDs(is, minTableID)
144+
if maxTableID == 0 || minTableID >= maxTableID {
145+
return maxTableID, nil
146+
}
147+
return maxTableID, buildMergeEmptyRegionsKeyRanges(occupiedTableIDs, minTableID, maxTableID)
148+
}
149+
150+
func collectMergeEmptyRegionsOccupiedTableIDs(is infoschema.InfoSchema, minTableID int64) ([]int64, int64) {
151+
occupiedTableIDs := make([]int64, 0)
152+
maxTableID := int64(0)
153+
for _, dbInfo := range is.AllSchemas() {
154+
schemaName := dbInfo.Name.L
155+
for _, tbl := range is.SchemaTables(dbInfo.Name) {
156+
tblInfo := tbl.Meta()
157+
if !mergeEmptyRegionsTableOccupiesKeyspace(tblInfo) {
158+
continue
159+
}
160+
161+
tableIDs := appendMergeEmptyRegionsTableIDs(nil, tblInfo)
162+
for _, tableID := range tableIDs {
163+
if tableID >= minTableID {
164+
occupiedTableIDs = append(occupiedTableIDs, tableID)
165+
}
166+
}
167+
// Keep skipped live tables as occupied keyspace so the background scan
168+
// never force merges through system or temporary table ranges in the
169+
// current scan window.
170+
if shouldSkipForceMergeTable(schemaName, tblInfo.Name.L, tblInfo) {
171+
continue
172+
}
173+
for _, tableID := range tableIDs {
174+
if tableID > maxTableID {
175+
maxTableID = tableID
176+
}
177+
}
178+
}
179+
}
180+
return occupiedTableIDs, maxTableID
181+
}
182+
183+
func buildMergeEmptyRegionsKeyRanges(
184+
occupiedTableIDs []int64,
185+
minTableID int64,
186+
maxTableID int64,
187+
) []infosync.ForceMergeKeyRange {
188+
if minTableID < 1 || minTableID >= maxTableID {
189+
return nil
190+
}
191+
192+
sort.Slice(occupiedTableIDs, func(i, j int) bool {
193+
return occupiedTableIDs[i] < occupiedTableIDs[j]
194+
})
195+
196+
ranges := make([]infosync.ForceMergeKeyRange, 0)
197+
nextRangeStart := minTableID
198+
for _, occupiedTableID := range occupiedTableIDs {
199+
if occupiedTableID < nextRangeStart {
200+
continue
201+
}
202+
if occupiedTableID >= maxTableID {
203+
break
204+
}
205+
if occupiedTableID > nextRangeStart {
206+
ranges = append(ranges, buildMergeEmptyRegionsTableIDRange(nextRangeStart, occupiedTableID-1))
207+
}
208+
nextRangeStart = occupiedTableID + 1
209+
}
210+
if nextRangeStart < maxTableID {
211+
ranges = append(ranges, buildMergeEmptyRegionsTableIDRange(nextRangeStart, maxTableID-1))
212+
}
213+
return ranges
214+
}
215+
216+
func mergeEmptyRegionsTableOccupiesKeyspace(tblInfo *model.TableInfo) bool {
217+
return tblInfo != nil && tblInfo.IsBaseTable()
218+
}
219+
220+
func appendMergeEmptyRegionsTableIDs(dst []int64, tblInfo *model.TableInfo) []int64 {
221+
if tblInfo == nil {
222+
return dst
223+
}
224+
225+
pi := tblInfo.GetPartitionInfo()
226+
if pi != nil && len(pi.Definitions) > 0 {
227+
for _, def := range pi.Definitions {
228+
if def.ID > 0 {
229+
dst = append(dst, def.ID)
230+
}
231+
}
232+
// Partitioned table global indexes live under the logical table ID prefix.
233+
if hasGlobalIndex(tblInfo) && tblInfo.ID > 0 {
234+
dst = append(dst, tblInfo.ID)
235+
}
236+
return dst
237+
}
238+
if tblInfo.ID > 0 {
239+
dst = append(dst, tblInfo.ID)
240+
}
241+
return dst
242+
}
243+
244+
func buildMergeEmptyRegionsTableIDRange(startTableID int64, endTableID int64) infosync.ForceMergeKeyRange {
245+
return infosync.ForceMergeKeyRange{
246+
StartKey: tablecodec.EncodeTablePrefix(startTableID),
247+
EndKey: tablecodec.EncodeTablePrefix(endTableID + 1),
248+
}
249+
}
250+
251+
func buildForceMergeTableKeyRange(tableID int64) infosync.ForceMergeKeyRange {
252+
return infosync.ForceMergeKeyRange{
253+
StartKey: tablecodec.EncodeTablePrefix(tableID),
254+
EndKey: tablecodec.EncodeTablePrefix(tableID + 1),
255+
}
256+
}

0 commit comments

Comments
 (0)