Skip to content
Open
Show file tree
Hide file tree
Changes from 16 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
24 changes: 15 additions & 9 deletions cmd/tidb-server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -526,10 +526,10 @@ func main() {
keyspaceName := keyspace.GetKeyspaceNameBySettings()
executor.Start()
resourcemanager.InstanceResourceManager.Start()
storage, dom, err := createStoreDDLOwnerMgrAndDomain(keyspaceName)
storage, dom, externalWorkloadManager, err := createStoreDDLOwnerMgrAndDomain(keyspaceName)
terror.MustNil(err)
repository.SetupRepository(dom)
if externalWorkloadManager := extworkload.GetManagerFromStore(storage); externalWorkloadManager != nil {
if externalWorkloadManager != nil {
defer closeExternalWorkloadManager(storage, externalWorkloadManager)
}
svr := createServer(storage, dom)
Expand Down Expand Up @@ -659,7 +659,7 @@ func registerStores() error {
return err
}

func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.Domain, error) {
func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.Domain, extworkload.Manager, error) {
if config.GetGlobalConfig().Store == config.StoreTypeUniStore {
kv.StandAloneTiDB = true
}
Expand All @@ -670,12 +670,12 @@ func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.D
if pdhttpCli != nil {
pdStatus, err := pdhttpCli.GetStatus(context.Background())
if err != nil {
return nil, nil, err
return nil, nil, nil, err
}
if !kerneltype.IsMatch(pdStatus.KernelType) {
log.Error("kernel type mismatch", zap.String("pd", pdStatus.KernelType),
zap.String("tidb", kerneltype.Name()))
return nil, nil, errors.New("kernel type mismatch")
return nil, nil, nil, errors.New("kernel type mismatch")
}
}
}
Expand All @@ -685,14 +685,20 @@ func createStoreDDLOwnerMgrAndDomain(keyspaceName string) (kv.Storage, *domain.D
// Bootstrap a session to load information schema.
err := ddl.StartOwnerManager(context.Background(), storage)
if err != nil {
return nil, nil, err
closeExternalWorkloadManager(storage, externalWorkloadManager)
return nil, nil, nil, err
}
dom, err := session.BootstrapSession(storage)
dom, err := session.BootstrapSessionWithExternalWorkloadManager(storage, externalWorkloadManager)
Comment thread
ystaticy marked this conversation as resolved.
if err != nil {
return nil, nil, err
closeExternalWorkloadManager(storage, externalWorkloadManager)
return nil, nil, nil, err
}
initializeExternalWorkloadGCV2(context.Background(), storage, externalWorkloadManager)
return storage, dom, nil
externalWorkloadManager = extworkload.GetManagerFromStore(storage)
if externalWorkloadManager == nil {
dom.SetExternalWorkloadManager(nil)
}
return storage, dom, externalWorkloadManager, nil
}

// Prometheus push.
Expand Down
2 changes: 2 additions & 0 deletions pkg/ddl/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ go_library(
"//pkg/expression",
"//pkg/expression/exprctx",
"//pkg/expression/exprstatic",
"//pkg/extworkload",
"//pkg/infoschema",
"//pkg/infoschema/context",
"//pkg/ingestor/engineapi",
Expand Down Expand Up @@ -271,6 +272,7 @@ go_test(
"executor_nokit_test.go",
"executor_test.go",
"export_test.go",
"external_workload_ttl_test.go",
"fail_test.go",
"foreign_key_test.go",
"index_change_test.go",
Expand Down
6 changes: 6 additions & 0 deletions pkg/ddl/create_table.go
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,8 @@ func (w *worker) onCreateTable(jobCtx *jobContext, job *model.Job) (ver int64, _
return ver, errors.Trace(err)
}

w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tbInfo)
Comment thread
ystaticy marked this conversation as resolved.
Outdated

Comment thread
coderabbitai[bot] marked this conversation as resolved.
// Finish this job.
job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo)
return ver, errors.Trace(err)
Expand Down Expand Up @@ -293,6 +295,7 @@ func (w *worker) createTableWithForeignKeys(jobCtx *jobContext, job *model.Job,
if err != nil {
return ver, errors.Trace(err)
}
w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tbInfo)

job.FinishTableJob(model.JobStateDone, model.StatePublic, ver, tbInfo)
return ver, nil
Expand Down Expand Up @@ -352,6 +355,9 @@ func (w *worker) onCreateTables(jobCtx *jobContext, job *model.Job) (int64, erro
return ver, errors.Trace(err)
}
}
for i := range tableInfos {
w.tryRegisterTTLTableToExternalWorkload(jobCtx.ctx, tableInfos[i])
}

job.State = model.JobStateDone
job.SchemaState = model.StatePublic
Expand Down
25 changes: 25 additions & 0 deletions pkg/ddl/db_table_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -414,6 +414,31 @@ func TestBatchCreateTable(t *testing.T) {
tk.Session().SetValue(sessionctx.QueryString, "skip")
err = d.BatchCreateTableWithInfo(tk.Session(), ast.NewCIStr("test"), []*model.TableInfo{newinfo}, ddl.WithOnExist(ddl.OnExistError))
require.NoError(t, err)

t.Run("batch create ttl tables registers external workload", func(t *testing.T) {
mgr := &recordingExternalWorkloadManager{role: config.RoleMaster}
tk, store := createTTLExternalWorkloadTestKit(t, mgr)
d := domain.GetDomain(tk.Session()).DDLExecutor()
infos := make([]*model.TableInfo, 0, 2)
for _, name := range []string{"ttl_batch_1", "ttl_batch_2"} {
info, err := testTableInfo(store, name, 2)
require.NoError(t, err)
info.Columns[0].FieldType = *types.NewFieldType(mysql.TypeDatetime)
info.TTLInfo = &model.TTLInfo{
ColumnName: info.Columns[0].Name,
IntervalExprStr: "1",
IntervalTimeUnit: int(ast.TimeUnitDay),
Enable: true,
JobInterval: model.DefaultTTLJobInterval,
}
infos = append(infos, info)
}

tk.Session().SetValue(sessionctx.QueryString, "skip")
err := d.BatchCreateTableWithInfo(tk.Session(), ast.NewCIStr("test"), infos, ddl.WithOnExist(ddl.OnExistError), ddl.WithIDAllocated(true))
require.NoError(t, err)
require.Equal(t, []int64{infos[0].ID, infos[1].ID}, mgr.registeredTTLTables())
})
}

// port from mysql
Expand Down
3 changes: 3 additions & 0 deletions pkg/ddl/ddl.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ import (
"github.com/pingcap/tidb/pkg/dxf/framework/proto"
"github.com/pingcap/tidb/pkg/dxf/framework/scheduler"
"github.com/pingcap/tidb/pkg/dxf/framework/taskexecutor"
"github.com/pingcap/tidb/pkg/extworkload"
"github.com/pingcap/tidb/pkg/infoschema"
"github.com/pingcap/tidb/pkg/kv"
"github.com/pingcap/tidb/pkg/meta"
Expand Down Expand Up @@ -348,6 +349,7 @@ type ddlCtx struct {
etcdCli *clientv3.Client
autoidCli *autoid.ClientDiscover
schemaLoader SchemaLoader
extWorkload extworkload.Manager

// reorgCtx is used for reorganization.
reorgCtx reorgContexts
Expand Down Expand Up @@ -780,6 +782,7 @@ func newDDL(ctx context.Context, options ...Option) (*ddl, *executor) {
etcdCli: opt.EtcdCli,
autoidCli: opt.AutoIDClient,
schemaLoader: opt.SchemaLoader,
extWorkload: opt.ExtWorkloadMgr,
}
ddlCtx.reorgCtx.reorgCtxMap = make(map[int64]*reorgCtx)
ddlCtx.jobCtx.jobCtxMap = make(map[int64]*ReorgContext)
Expand Down
203 changes: 203 additions & 0 deletions pkg/ddl/external_workload_ttl_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,203 @@
// Copyright 2026 PingCAP, Inc.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

package ddl_test

import (
"context"
"sync"
"testing"
"time"

"github.com/pingcap/kvproto/pkg/keyspacepb"
"github.com/pingcap/tidb/pkg/config"
"github.com/pingcap/tidb/pkg/domain"
"github.com/pingcap/tidb/pkg/kv"
"github.com/pingcap/tidb/pkg/session"
"github.com/pingcap/tidb/pkg/sessionctx/vardef"
"github.com/pingcap/tidb/pkg/store/mockstore"
"github.com/pingcap/tidb/pkg/testkit"
"github.com/pingcap/tidb/pkg/testkit/external"
"github.com/stretchr/testify/require"
)

type recordingExternalWorkloadManager struct {
role config.ExternalWorkloadRole

mu sync.Mutex
registeredTables []int64
deletedTables []int64
}

func (*recordingExternalWorkloadManager) Close() error { return nil }

func (m *recordingExternalWorkloadManager) Role() config.ExternalWorkloadRole {
return m.role
}

func (*recordingExternalWorkloadManager) Meta() *keyspacepb.KeyspaceMeta { return nil }

func (*recordingExternalWorkloadManager) InitializeGCV2(context.Context, time.Duration) error {
return nil
}

func (*recordingExternalWorkloadManager) AbortGCV2(context.Context) error { return nil }

func (*recordingExternalWorkloadManager) RegisterGCV2(context.Context, uint64, time.Duration) error {
return nil
}

func (*recordingExternalWorkloadManager) RecycleGCV2(context.Context, uint64) error {
return nil
}

func (*recordingExternalWorkloadManager) UpdateGCLifeTime(context.Context, time.Duration) error {
return nil
}

func (m *recordingExternalWorkloadManager) RegisterTTLTask(_ context.Context, tableID int64, _ bool) error {
m.mu.Lock()
defer m.mu.Unlock()
m.registeredTables = append(m.registeredTables, tableID)
return nil
}

func (m *recordingExternalWorkloadManager) DeleteTTLTableInfo(_ context.Context, tableID int64) error {
m.mu.Lock()
defer m.mu.Unlock()
m.deletedTables = append(m.deletedTables, tableID)
return nil
}

func (*recordingExternalWorkloadManager) RecycleTTLTask(context.Context, uint64) error {
return nil
}

func (*recordingExternalWorkloadManager) UpdateTTLJobEnable(context.Context, bool) error {
return nil
}

func (*recordingExternalWorkloadManager) RegisterAutoAnalyze(context.Context, uint64) error {
return nil
}

func (*recordingExternalWorkloadManager) RecycleAutoAnalyze(context.Context, uint64) error {
return nil
}

func (m *recordingExternalWorkloadManager) registeredTTLTables() []int64 {
m.mu.Lock()
defer m.mu.Unlock()
return append([]int64(nil), m.registeredTables...)
}

func (m *recordingExternalWorkloadManager) deletedTTLTables() []int64 {
m.mu.Lock()
defer m.mu.Unlock()
return append([]int64(nil), m.deletedTables...)
}

func createTTLExternalWorkloadTestKit(t *testing.T, mgr *recordingExternalWorkloadManager) (*testkit.TestKit, kv.Storage) {
store, err := mockstore.NewMockStore()
require.NoError(t, err)

vardef.SetSchemaLease(500 * time.Millisecond)
session.DisableStats4Test()
domain.DisablePlanReplayerBackgroundJob4Test()
domain.DisableDumpHistoricalStats4Test()

dom, err := session.BootstrapSessionWithExternalWorkloadManager(store, mgr)
require.NoError(t, err)
dom.SetStatsUpdating(true)

t.Cleanup(func() {
dom.Close()
require.NoError(t, store.Close())
})

tk := testkit.NewTestKit(t, store)
tk.MustExec("use test")
tk.MustExec("set @@global.tidb_enable_foreign_key=1")
tk.MustExec("set @@foreign_key_checks=1")
return tk, store
}

func TestExternalWorkloadTTLDDLIntegration(t *testing.T) {
t.Run("create table with foreign key registers ttl", func(t *testing.T) {
mgr := &recordingExternalWorkloadManager{role: config.RoleMaster}
tk, _ := createTTLExternalWorkloadTestKit(t, mgr)

tk.MustExec("create table parent(id int primary key)")
tk.MustExec(`create table child(
id int primary key,
parent_id int,
created_at datetime,
index idx_parent(parent_id),
foreign key (parent_id) references parent(id)
) TTL = created_at + interval 1 day`)

childTbl := external.GetTableByName(t, tk, "test", "child")
require.Equal(t, []int64{childTbl.Meta().ID}, mgr.registeredTTLTables())
require.Empty(t, mgr.deletedTTLTables())
})

t.Run("drop table deletes ttl metadata", func(t *testing.T) {
mgr := &recordingExternalWorkloadManager{role: config.RoleMaster}
tk, _ := createTTLExternalWorkloadTestKit(t, mgr)

tk.MustExec(`create table t(
id int primary key,
created_at datetime
) TTL = created_at + interval 1 day`)

tbl := external.GetTableByName(t, tk, "test", "t")
tk.MustExec("drop table t")

require.Equal(t, []int64{tbl.Meta().ID}, mgr.deletedTTLTables())
})

t.Run("truncate table refreshes ttl metadata", func(t *testing.T) {
mgr := &recordingExternalWorkloadManager{role: config.RoleMaster}
tk, _ := createTTLExternalWorkloadTestKit(t, mgr)

tk.MustExec(`create table t(
id int primary key,
created_at datetime
) TTL = created_at + interval 1 day`)

oldTbl := external.GetTableByName(t, tk, "test", "t")
tk.MustExec("truncate table t")
newTbl := external.GetTableByName(t, tk, "test", "t")

require.NotEqual(t, oldTbl.Meta().ID, newTbl.Meta().ID)
require.Equal(t, []int64{oldTbl.Meta().ID}, mgr.deletedTTLTables())
require.Equal(t, []int64{oldTbl.Meta().ID, newTbl.Meta().ID}, mgr.registeredTTLTables())
})

t.Run("ddl syncs ttl metadata from ttl worker role", func(t *testing.T) {
mgr := &recordingExternalWorkloadManager{role: config.RoleTTLTaskWorker}
tk, _ := createTTLExternalWorkloadTestKit(t, mgr)

tk.MustExec(`create table t(
id int primary key,
created_at datetime
) TTL = created_at + interval 1 day`)

tbl := external.GetTableByName(t, tk, "test", "t")
require.Equal(t, []int64{tbl.Meta().ID}, mgr.registeredTTLTables())

tk.MustExec("drop table t")
require.Equal(t, []int64{tbl.Meta().ID}, mgr.deletedTTLTables())
})
}
9 changes: 9 additions & 0 deletions pkg/ddl/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"time"

"github.com/pingcap/tidb/pkg/ddl/notifier"
"github.com/pingcap/tidb/pkg/extworkload"
"github.com/pingcap/tidb/pkg/infoschema"
"github.com/pingcap/tidb/pkg/kv"
"github.com/pingcap/tidb/pkg/meta/autoid"
Expand All @@ -36,6 +37,7 @@ type Options struct {
Lease time.Duration
SchemaLoader SchemaLoader
EventPublishStore notifier.Store
ExtWorkloadMgr extworkload.Manager
}

// WithEtcdClient specifies the `clientv3.Client` of DDL used to request the etcd service
Expand Down Expand Up @@ -86,3 +88,10 @@ func WithEventPublishStore(store notifier.Store) Option {
options.EventPublishStore = store
}
}

// WithExternalWorkloadManager specifies the manager used to coordinate external background workloads.
func WithExternalWorkloadManager(manager extworkload.Manager) Option {
return func(options *Options) {
options.ExtWorkloadMgr = manager
}
}
Loading
Loading