From 4b6c357d42f54fe3bc5e282d9f4e52265e5ce4eb Mon Sep 17 00:00:00 2001 From: nhsmw Date: Fri, 7 Aug 2026 18:34:05 +0800 Subject: [PATCH 1/4] This is an automated cherry-pick of #5836 Signed-off-by: ti-chi-bot --- api/v2/model.go | 21 ++++ api/v2/model_test.go | 20 ++++ cmd/cdc/cli/cli_changefeed_create_test.go | 3 + downstreamadapter/sink/mysql/sink.go | 120 +++++++++++++++++++-- downstreamadapter/sink/mysql/sink_test.go | 96 +++++++++++++++++ pkg/config/sink.go | 1 + pkg/sink/mysql/config.go | 122 +++++++++++++++++++--- pkg/sink/mysql/config_test.go | 94 +++++++++++++++++ pkg/sink/mysql/mysql_writer.go | 15 +++ pkg/sink/mysql/mysql_writer_ddl.go | 29 ++++- pkg/sink/mysql/mysql_writer_test.go | 73 +++++++++++++ 11 files changed, 567 insertions(+), 27 deletions(-) diff --git a/api/v2/model.go b/api/v2/model.go index f0f2316304..4db4c3d33b 100644 --- a/api/v2/model.go +++ b/api/v2/model.go @@ -472,6 +472,7 @@ func (c *ReplicaConfig) toInternalReplicaConfigWithOriginConfig( WriteTimeout: c.Sink.MySQLConfig.WriteTimeout, ReadTimeout: c.Sink.MySQLConfig.ReadTimeout, Timeout: c.Sink.MySQLConfig.Timeout, + AsyncDDLTimeout: c.Sink.MySQLConfig.AsyncDDLTimeout, EnableBatchDML: c.Sink.MySQLConfig.EnableBatchDML, EnableMultiStatement: c.Sink.MySQLConfig.EnableMultiStatement, EnableCachePreparedStatement: c.Sink.MySQLConfig.EnableCachePreparedStatement, @@ -794,6 +795,7 @@ func ToAPIReplicaConfig(c *config.ReplicaConfig) *ReplicaConfig { WriteTimeout: cloned.Sink.MySQLConfig.WriteTimeout, ReadTimeout: cloned.Sink.MySQLConfig.ReadTimeout, Timeout: cloned.Sink.MySQLConfig.Timeout, + AsyncDDLTimeout: cloned.Sink.MySQLConfig.AsyncDDLTimeout, EnableBatchDML: cloned.Sink.MySQLConfig.EnableBatchDML, EnableMultiStatement: cloned.Sink.MySQLConfig.EnableMultiStatement, EnableCachePreparedStatement: cloned.Sink.MySQLConfig.EnableCachePreparedStatement, @@ -1473,6 +1475,7 @@ type KafkaConfig struct { // MySQLConfig represents a MySQL sink configuration type MySQLConfig struct { +<<<<<<< HEAD WorkerCount *int `json:"worker_count,omitempty"` MaxTxnRow *int `json:"max_txn_row,omitempty"` MaxMultiUpdateRowSize *int `json:"max_multi_update_row_size,omitempty"` @@ -1488,6 +1491,24 @@ type MySQLConfig struct { EnableBatchDML *bool `json:"enable_batch_dml,omitempty"` EnableMultiStatement *bool `json:"enable_multi_statement,omitempty"` EnableCachePreparedStatement *bool `json:"enable_cache_prepared_statement,omitempty"` +======= + WorkerCount *int `json:"worker_count,omitempty" toml:"worker-count,omitempty"` + MaxTxnRow *int `json:"max_txn_row,omitempty" toml:"max-txn-row,omitempty"` + MaxMultiUpdateRowSize *int `json:"max_multi_update_row_size,omitempty" toml:"max-multi-update-row-size,omitempty"` + MaxMultiUpdateRowCount *int `json:"max_multi_update_row_count,omitempty" toml:"max-multi-update-row-count,omitempty"` + TiDBTxnMode *string `json:"tidb_txn_mode,omitempty" toml:"tidb-txn-mode,omitempty"` + SSLCa *string `json:"ssl_ca,omitempty" toml:"ssl-ca,omitempty"` + SSLCert *string `json:"ssl_cert,omitempty" toml:"ssl-cert,omitempty"` + SSLKey *string `json:"ssl_key,omitempty" toml:"ssl-key,omitempty"` + TimeZone *string `json:"time_zone,omitempty" toml:"time-zone,omitempty"` + WriteTimeout *string `json:"write_timeout,omitempty" toml:"write-timeout,omitempty"` + ReadTimeout *string `json:"read_timeout,omitempty" toml:"read-timeout,omitempty"` + Timeout *string `json:"timeout,omitempty" toml:"timeout,omitempty"` + AsyncDDLTimeout *string `json:"async_ddl_timeout,omitempty" toml:"async-ddl-timeout,omitempty"` + EnableBatchDML *bool `json:"enable_batch_dml,omitempty" toml:"enable-batch-dml,omitempty"` + EnableMultiStatement *bool `json:"enable_multi_statement,omitempty" toml:"enable-multi-statement,omitempty"` + EnableCachePreparedStatement *bool `json:"enable_cache_prepared_statement,omitempty" toml:"enable-cache-prepared-statement,omitempty"` +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } // CloudStorageConfig represents a cloud storage sink configuration diff --git a/api/v2/model_test.go b/api/v2/model_test.go index fe0afaee00..0169f37689 100644 --- a/api/v2/model_test.go +++ b/api/v2/model_test.go @@ -103,3 +103,23 @@ func TestReplicaConfigConversion(t *testing.T) { require.Equal(t, "correctness", *apiCfgBack.Integrity.IntegrityCheckLevel) require.Equal(t, "eventual", *apiCfgBack.Consistent.Level) } + +func TestReplicaConfigConversionMySQLAsyncDDLTimeout(t *testing.T) { + t.Parallel() + + apiCfg := &ReplicaConfig{ + Sink: &SinkConfig{ + MySQLConfig: &MySQLConfig{ + AsyncDDLTimeout: util.AddressOf("45m"), + }, + }, + } + + internalCfg := apiCfg.ToInternalReplicaConfig() + require.NotNil(t, internalCfg.Sink.MySQLConfig) + require.Equal(t, "45m", util.GetOrZero(internalCfg.Sink.MySQLConfig.AsyncDDLTimeout)) + + apiCfgBack := ToAPIReplicaConfig(internalCfg) + require.NotNil(t, apiCfgBack.Sink.MySQLConfig) + require.Equal(t, "45m", util.GetOrZero(apiCfgBack.Sink.MySQLConfig.AsyncDDLTimeout)) +} diff --git a/cmd/cdc/cli/cli_changefeed_create_test.go b/cmd/cdc/cli/cli_changefeed_create_test.go index c390d2c731..e69c9372ec 100644 --- a/cmd/cdc/cli/cli_changefeed_create_test.go +++ b/cmd/cdc/cli/cli_changefeed_create_test.go @@ -74,6 +74,9 @@ func TestTomlFileToApiModel(t *testing.T) { content := ` [filter] rules = ['*.*', '!test.*'] + + [sink.mysql-config] + async-ddl-timeout = "45m" ` err := os.WriteFile(path, []byte(content), 0o644) require.Nil(t, err) diff --git a/downstreamadapter/sink/mysql/sink.go b/downstreamadapter/sink/mysql/sink.go index a62c864bea..abd8e3cfe0 100644 --- a/downstreamadapter/sink/mysql/sink.go +++ b/downstreamadapter/sink/mysql/sink.go @@ -46,11 +46,12 @@ type Sink struct { dmlWriter []*mysql.Writer ddlWriter *mysql.Writer - // dmlDB and controlDB are the DB pools this sink is responsible for closing. + // dmlDB, controlDB, and controlAsyncDB are the DB pools this sink is responsible for closing. // Compatibility callers built through NewMySQLSink use one shared pool. - dmlDB *sql.DB - controlDB *sql.DB - statistics *metrics.Statistics + dmlDB *sql.DB + controlDB *sql.DB + controlAsyncDB *sql.DB + statistics *metrics.Statistics conflictDetector *causality.ConflictDetector @@ -71,12 +72,15 @@ func Verify( config *config.ChangefeedConfig, ) error { testID := common.NewChangefeedID4Test("test", "mysql_create_sink_test") - _, dmlDB, controlDB, err := mysql.NewMysqlConfigAndDBs(ctx, testID, uri, config) + _, dmlDB, controlDB, controlAsyncDB, err := mysql.NewMysqlConfigAndDBs(ctx, testID, uri, config) if err != nil { return err } _ = dmlDB.Close() _ = controlDB.Close() + if controlAsyncDB != nil { + _ = controlAsyncDB.Close() + } return nil } @@ -86,7 +90,7 @@ func New( config *config.ChangefeedConfig, sinkURI *url.URL, ) (*Sink, error) { - cfg, dmlDB, controlDB, err := mysql.NewMysqlConfigAndDBs(ctx, changefeedID, sinkURI, config) + cfg, dmlDB, controlDB, controlAsyncDB, err := mysql.NewMysqlConfigAndDBs(ctx, changefeedID, sinkURI, config) if err != nil { return nil, err } @@ -102,9 +106,14 @@ func New( metrics.ChangefeedDownstreamIsTiDBGauge.DeleteLabelValues(keyspace, name) } +<<<<<<< HEAD return newMySQLSinkWithControlDB(ctx, changefeedID, cfg, dmlDB, controlDB, config.BDRMode), nil +======= + return newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, config.BDRMode, config.EnableActiveActive, config.ActiveActiveProgressInterval, keyspaceID), nil +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } +// NewMySQLSink used for test func NewMySQLSink( ctx context.Context, changefeedID common.ChangeFeedID, @@ -112,7 +121,15 @@ func NewMySQLSink( db *sql.DB, bdrMode bool, ) *Sink { +<<<<<<< HEAD return newMySQLSinkWithControlDB(ctx, changefeedID, cfg, db, db, bdrMode) +======= + var controlAsyncDB *sql.DB + if cfg.IsTiDB { + controlAsyncDB = db + } + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, db, db, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } // newMySQLSinkWithControlDB creates a MySQL sink with separate pools for DML and @@ -126,13 +143,76 @@ func newMySQLSinkWithControlDB( controlDB *sql.DB, bdrMode bool, ) *Sink { +<<<<<<< HEAD stat := metrics.NewStatistics(changefeedID, "TxnSink") +======= + var controlAsyncDB *sql.DB + if cfg.IsTiDB { + controlAsyncDB = controlDB + } + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) +} + +func newMySQLSinkWithControlAsyncDB( + ctx context.Context, + changefeedID common.ChangeFeedID, + cfg *mysql.Config, + dmlDB *sql.DB, + controlDB *sql.DB, + controlAsyncDB *sql.DB, + bdrMode bool, + enableActiveActive bool, + progressInterval time.Duration, + keyspaceID uint32, +) *Sink { + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) +} + +func newMySQLSinkWithDBs( + ctx context.Context, + changefeedID common.ChangeFeedID, + cfg *mysql.Config, + dmlDB *sql.DB, + controlDB *sql.DB, + controlAsyncDB *sql.DB, + bdrMode bool, + enableActiveActive bool, + progressInterval time.Duration, + keyspaceID uint32, +) *Sink { + if !cfg.IsTiDB { + controlAsyncDB = nil + } else if controlAsyncDB == nil { + controlAsyncDB = controlDB + } + + stat := metrics.NewStatistics(changefeedID, keyspaceID, "TxnSink") + + var activeActiveSyncStatsCollector *mysql.ActiveActiveSyncStatsCollector + if enableActiveActive && cfg.IsTiDB && cfg.ActiveActiveSyncStatsInterval > 0 { + supported, err := mysql.CheckActiveActiveSyncStatsSupported(ctx, dmlDB) + if err != nil { + log.Info("failed to check tidb_cdc_active_active_sync_stats support, disable metric collection", + zap.String("keyspace", changefeedID.Keyspace()), + zap.Stringer("changefeed", changefeedID), + zap.Error(err)) + } else if supported { + activeActiveSyncStatsCollector = mysql.NewActiveActiveSyncStatsCollector(changefeedID) + } else { + log.Info("downstream does not support tidb_cdc_active_active_sync_stats, disable metric collection", + zap.String("keyspace", changefeedID.Keyspace()), + zap.Stringer("changefeed", changefeedID)) + } + } + +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) result := &Sink{ - changefeedID: changefeedID, - dmlDB: dmlDB, - controlDB: controlDB, - dmlWriter: make([]*mysql.Writer, cfg.WorkerCount), - statistics: stat, + changefeedID: changefeedID, + dmlDB: dmlDB, + controlDB: controlDB, + controlAsyncDB: controlAsyncDB, + dmlWriter: make([]*mysql.Writer, cfg.WorkerCount), + statistics: stat, conflictDetector: causality.New(defaultConflictDetectorSlots, causality.TxnCacheOption{ Count: cfg.WorkerCount, @@ -146,7 +226,16 @@ func newMySQLSinkWithControlDB( bdrMode: bdrMode, } for i := 0; i < len(result.dmlWriter); i++ { +<<<<<<< HEAD result.dmlWriter[i] = mysql.NewWriter(ctx, i, dmlDB, cfg, changefeedID, stat) +======= + result.dmlWriter[i] = mysql.NewWriter(ctx, i, dmlDB, cfg, changefeedID, stat, activeActiveSyncStatsCollector) + } + result.ddlWriter = mysql.NewWriter(ctx, len(result.dmlWriter), controlDB, cfg, changefeedID, stat, nil) + result.ddlWriter.SetControlAsyncDB(controlAsyncDB) + if enableActiveActive { + result.progressTableWriter = mysql.NewProgressTableWriter(ctx, controlDB, changefeedID, cfg.MaxTxnRow, progressInterval) +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } result.ddlWriter = mysql.NewWriter(ctx, len(result.dmlWriter), controlDB, cfg, changefeedID, stat) return result @@ -369,6 +458,15 @@ func (s *Sink) Close() { if s.controlDB != s.dmlDB { s.closeDBPool("control", s.controlDB) } +<<<<<<< HEAD +======= + if s.controlAsyncDB != nil && s.controlAsyncDB != s.dmlDB && s.controlAsyncDB != s.controlDB { + s.closeDBPool("control async", s.controlAsyncDB) + } + if s.activeActiveSyncStatsCollector != nil { + s.activeActiveSyncStatsCollector.Close() + } +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) s.statistics.Close() metrics.ChangefeedDownstreamIsTiDBGauge.DeleteLabelValues(s.changefeedID.Keyspace(), s.changefeedID.Name()) diff --git a/downstreamadapter/sink/mysql/sink_test.go b/downstreamadapter/sink/mysql/sink_test.go index 99c596128e..85eb0f002b 100644 --- a/downstreamadapter/sink/mysql/sink_test.go +++ b/downstreamadapter/sink/mysql/sink_test.go @@ -25,7 +25,12 @@ import ( "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/sink/mysql" +<<<<<<< HEAD "github.com/pingcap/tidb/pkg/sessionctx/variable" +======= + timodel "github.com/pingcap/tidb/pkg/meta/model" + "github.com/pingcap/tidb/pkg/sessionctx/vardef" +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) "github.com/stretchr/testify/require" ) @@ -79,6 +84,97 @@ func getMysqlSinkWithSeparateDBs(t *testing.T) (context.Context, *Sink, sqlmock. return ctx, sink, dmlMock, controlMock } +func TestMysqlSinkControlAsyncDBOnlyForTiDB(t *testing.T) { + ctx := context.Background() + changefeedID := common.NewChangefeedID4Test("test", "test") + + t.Run("mysql downstream has no control async db", func(t *testing.T) { + db, mock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual)) + require.NoError(t, err) + + cfg := mysql.New() + cfg.WorkerCount = 1 + cfg.MaxAllowedPacket = int64(vardef.DefMaxAllowedPacket) + cfg.CachePrepStmts = false + cfg.IsTiDB = false + + sink := NewMySQLSink(ctx, changefeedID, cfg, db, false, false, time.Minute, common.DefaultKeyspaceID) + require.Nil(t, sink.controlAsyncDB) + + mock.ExpectClose() + sink.Close() + require.NoError(t, mock.ExpectationsWereMet()) + }) + + t.Run("tidb downstream has control async db", func(t *testing.T) { + db, mock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual)) + require.NoError(t, err) + + cfg := mysql.New() + cfg.WorkerCount = 1 + cfg.MaxAllowedPacket = int64(vardef.DefMaxAllowedPacket) + cfg.CachePrepStmts = false + cfg.IsTiDB = true + + sink := NewMySQLSink(ctx, changefeedID, cfg, db, false, false, time.Minute, common.DefaultKeyspaceID) + require.Same(t, db, sink.controlAsyncDB) + + mock.ExpectClose() + sink.Close() + require.NoError(t, mock.ExpectationsWereMet()) + }) +} + +func TestMysqlSinkUsesControlAsyncDBForTiDBAddIndex(t *testing.T) { + dmlDB, dmlMock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual)) + require.NoError(t, err) + controlDB, controlMock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual)) + require.NoError(t, err) + controlAsyncDB, controlAsyncMock, err := sqlmock.New(sqlmock.QueryMatcherOption(sqlmock.QueryMatcherEqual)) + require.NoError(t, err) + + ctx := context.Background() + changefeedID := common.NewChangefeedID4Test("test", "test") + cfg := mysql.New() + cfg.WorkerCount = 1 + cfg.MaxAllowedPacket = int64(vardef.DefMaxAllowedPacket) + cfg.CachePrepStmts = false + cfg.EnableDDLTs = false + cfg.IsTiDB = true + + sink := newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, false, false, time.Minute, common.DefaultKeyspaceID) + + ddl := &commonEvent.DDLEvent{ + Type: byte(timodel.ActionAddIndex), + Query: "alter table t add index idx_name(name);", + SchemaName: "test", + TableName: "t", + BlockedTables: &commonEvent.InfluencedTables{ + InfluenceType: commonEvent.InfluenceTypeNormal, + TableIDs: []int64{1}, + }, + } + + controlMock.ExpectQuery("BEGIN; SET @ticdc_ts := TIDB_PARSE_TSO(@@tidb_current_ts); ROLLBACK; SELECT @ticdc_ts; SET @ticdc_ts=NULL;"). + WillReturnRows(sqlmock.NewRows([]string{"@ticdc_ts"}).AddRow("2021-05-26 11:33:37.776000")) + controlAsyncMock.ExpectBegin() + controlAsyncMock.ExpectExec("USE `test`;").WillReturnResult(sqlmock.NewResult(1, 1)) + controlAsyncMock.ExpectExec("SET TIMESTAMP = DEFAULT").WillReturnResult(sqlmock.NewResult(1, 1)) + controlAsyncMock.ExpectExec("alter table t add index idx_name(name);").WillReturnResult(sqlmock.NewResult(1, 1)) + controlAsyncMock.ExpectCommit() + + require.NoError(t, sink.WriteBlockEvent(ddl)) + + dmlMock.ExpectClose() + controlMock.ExpectClose() + controlAsyncMock.ExpectClose() + sink.Close() + + require.NoError(t, dmlMock.ExpectationsWereMet()) + require.NoError(t, controlMock.ExpectationsWereMet()) + require.NoError(t, controlAsyncMock.ExpectationsWereMet()) +} + func MysqlSinkForTest() (*Sink, sqlmock.Sqlmock) { ctx, sink, mock := getMysqlSink() go sink.Run(ctx) diff --git a/pkg/config/sink.go b/pkg/config/sink.go index e0edb288ef..991db8d277 100644 --- a/pkg/config/sink.go +++ b/pkg/config/sink.go @@ -720,6 +720,7 @@ type MySQLConfig struct { WriteTimeout *string `toml:"write-timeout" json:"write-timeout,omitempty"` ReadTimeout *string `toml:"read-timeout" json:"read-timeout,omitempty"` Timeout *string `toml:"timeout" json:"timeout,omitempty"` + AsyncDDLTimeout *string `toml:"async-ddl-timeout" json:"async-ddl-timeout,omitempty"` EnableBatchDML *bool `toml:"enable-batch-dml" json:"enable-batch-dml,omitempty"` EnableMultiStatement *bool `toml:"enable-multi-statement" json:"enable-multi-statement,omitempty"` EnableCachePreparedStatement *bool `toml:"enable-cache-prepared-statement" json:"enable-cache-prepared-statement,omitempty"` diff --git a/pkg/sink/mysql/config.go b/pkg/sink/mysql/config.go index 0f88eae1e4..cfc07815c2 100644 --- a/pkg/sink/mysql/config.go +++ b/pkg/sink/mysql/config.go @@ -65,13 +65,14 @@ const ( // The upper limit of max multi update row size(8KB). maxMaxMultiUpdateRowSize = 8192 - defaultTiDBTxnMode = txnModeOptimistic - defaultReadTimeout = "2m" - defaultWriteTimeout = "2m" - defaultDialTimeout = "2m" - defaultSafeMode = false - defaultTxnIsolationRC = "READ-COMMITTED" - defaultCharacterSet = "utf8mb4" + defaultTiDBTxnMode = txnModeOptimistic + defaultReadTimeout = "2m" + defaultWriteTimeout = "2m" + defaultDialTimeout = "2m" + defaultAsyncDDLTimeout = "10s" + defaultSafeMode = false + defaultTxnIsolationRC = "READ-COMMITTED" + defaultCharacterSet = "utf8mb4" // BackoffBaseDelay indicates the base delay time for retrying. BackoffBaseDelay = 100 * time.Millisecond @@ -107,6 +108,7 @@ type Config struct { MaxMultiUpdateRowCount int MaxMultiUpdateRowSize int TidbTxnMode string +<<<<<<< HEAD ReadTimeout string WriteTimeout string DialTimeout string @@ -116,6 +118,23 @@ type Config struct { SSLCa string SSLCert string SSLKey string +======= + // tidbTxnModeSpecified indicates whether TidbTxnMode is explicitly set by user via sink URI or changefeed config. + // It is used to avoid overriding user configuration when applying downstream-specific defaults. + tidbTxnModeSpecified bool + ReadTimeout string + WriteTimeout string + DialTimeout string + // AsyncDDLTimeout controls the read timeout for the async DDL DB pool. + // If it is not explicitly set, it defaults to defaultAsyncDDLTimeout. + AsyncDDLTimeout string + SafeMode bool + Timezone string + TLS string + SSLCa string + SSLCert string + SSLKey string +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) // retry number for dml DMLMaxRetry uint64 @@ -163,6 +182,7 @@ type Config struct { // New returns the default mysql backend config. func New() *Config { return &Config{ +<<<<<<< HEAD WorkerCount: DefaultTiDBWorkerCount, workerCountSpecified: false, MaxTxnRow: DefaultMaxTxnRow, @@ -182,6 +202,29 @@ func New() *Config { EnableDDLTs: defaultEnableDDLTs, SlowQuery: slowQuery, whereClause: sqlmodel.DefaultWhereClause, +======= + WorkerCount: DefaultTiDBWorkerCount, + workerCountSpecified: false, + MaxTxnRow: DefaultMaxTxnRow, + MaxMultiUpdateRowCount: defaultMaxMultiUpdateRowCount, + MaxMultiUpdateRowSize: defaultMaxMultiUpdateRowSize, + TidbTxnMode: defaultTiDBTxnMode, + ReadTimeout: defaultReadTimeout, + WriteTimeout: defaultWriteTimeout, + DialTimeout: defaultDialTimeout, + AsyncDDLTimeout: defaultAsyncDDLTimeout, + SafeMode: defaultSafeMode, + BatchDMLEnable: defaultBatchDMLEnable, + MultiStmtEnable: defaultMultiStmtEnable, + CachePrepStmts: defaultCachePrepStmts, + SourceID: config.DefaultTiDBSourceID, + DMLMaxRetry: 8, + HasVectorType: defaultHasVectorType, + EnableDDLTs: defaultEnableDDLTs, + SlowQuery: slowQuery, + ActiveActiveSyncStatsInterval: time.Minute, + whereClause: sqlmodel.DefaultWhereClause, +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } } @@ -205,6 +248,7 @@ func (c *Config) mergeConfig(cfg *config.ChangefeedConfig) { merge(&c.WriteTimeout, mConfig.WriteTimeout) merge(&c.ReadTimeout, mConfig.ReadTimeout) merge(&c.DialTimeout, mConfig.Timeout) + merge(&c.AsyncDDLTimeout, mConfig.AsyncDDLTimeout) merge(&c.BatchDMLEnable, mConfig.EnableBatchDML) merge(&c.MultiStmtEnable, mConfig.EnableMultiStatement) merge(&c.CachePrepStmts, mConfig.EnableCachePreparedStatement) @@ -265,6 +309,9 @@ func (c *Config) Apply( if err = getDuration(query, "timeout", &c.DialTimeout); err != nil { return err } + if err = getDuration(query, "async-ddl-timeout", &c.AsyncDDLTimeout); err != nil { + return err + } if err = getBatchDMLEnable(query, &c.BatchDMLEnable); err != nil { return err } @@ -309,16 +356,16 @@ func NewMysqlConfigAndDB( } // NewMysqlConfigAndDBs creates the effective MySQL sink config and independent -// database pools for DML and control-plane work. The DML pool follows the worker -// based sizing, while the control pool remains small and independent so DDL, -// DDL-ts, syncpoint, and progress metadata operations cannot be starved by -// long-lived DML sessions. +// database pools for DML, control-plane work, and TiDB asynchronous DDL +// execution. The DML pool follows the worker based sizing, while the control +// pool remains small and independent so DDL, DDL-ts, syncpoint, and progress +// metadata operations cannot be starved by long-lived DML sessions. func NewMysqlConfigAndDBs( ctx context.Context, changefeedID common.ChangeFeedID, sinkURI *url.URL, config *config.ChangefeedConfig, -) (*Config, *sql.DB, *sql.DB, error) { +) (*Config, *sql.DB, *sql.DB, *sql.DB, error) { cfg, dmlDB, dsnStr, err := newMysqlConfigAndDB(ctx, changefeedID, sinkURI, config) if err != nil { - return nil, nil, nil, err + return nil, nil, nil, nil, err } controlDB, err := CreateMysqlDBConn(dsnStr) @@ -327,10 +374,42 @@ func NewMysqlConfigAndDBs( log.Warn("close mysql dml db after control db creation failed", zap.String("changefeed", changefeedID.String()), zap.Error(closeErr)) } - return nil, nil, nil, err + return nil, nil, nil, nil, err } configureControlDBConn(controlDB) - return cfg, dmlDB, controlDB, nil + + if !cfg.IsTiDB { + return cfg, dmlDB, controlDB, nil, nil + } + + controlAsyncDSNStr, err := setDSNReadTimeout(dsnStr, cfg.AsyncDDLTimeout) + if err != nil { + closeDMLAndControlDBAfterFailure(changefeedID, dmlDB, controlDB, "async ddl db dsn creation failed") + return nil, nil, nil, nil, err + } + controlAsyncDB, err := CreateMysqlDBConn(controlAsyncDSNStr) + if err != nil { + closeDMLAndControlDBAfterFailure(changefeedID, dmlDB, controlDB, "async ddl db creation failed") + return nil, nil, nil, nil, err + } + configureControlDBConn(controlAsyncDB) + return cfg, dmlDB, controlDB, controlAsyncDB, nil +} + +func closeDMLAndControlDBAfterFailure( + changefeedID common.ChangeFeedID, + dmlDB *sql.DB, + controlDB *sql.DB, + failureContext string, +) { + if closeErr := dmlDB.Close(); closeErr != nil { + log.Warn("close mysql dml db after "+failureContext, + zap.String("changefeed", changefeedID.String()), zap.Error(closeErr)) + } + if closeErr := controlDB.Close(); closeErr != nil { + log.Warn("close mysql control db after "+failureContext, + zap.String("changefeed", changefeedID.String()), zap.Error(closeErr)) + } } func newMysqlConfigAndDB( @@ -417,6 +496,19 @@ func newMysqlConfigAndDB( return cfg, db, dsnStr, nil } +func setDSNReadTimeout(dsnStr string, readTimeout string) (string, error) { + dsn, err := dmysql.ParseDSN(dsnStr) + if err != nil { + return "", errors.WrapError(errors.ErrMySQLInvalidConfig, err) + } + readTimeoutDuration, err := time.ParseDuration(readTimeout) + if err != nil { + return "", errors.WrapError(errors.ErrMySQLInvalidConfig, err) + } + dsn.ReadTimeout = readTimeoutDuration + return dsn.FormatDSN(), nil +} + func configureDMLDBConn(db *sql.DB, cfg *Config) { // Keep one spare DML connection so db.Prepare on a statement-cache miss // cannot wait behind all writer-owned transaction sessions. Control-plane diff --git a/pkg/sink/mysql/config_test.go b/pkg/sink/mysql/config_test.go index 331c965f6f..62b366db38 100644 --- a/pkg/sink/mysql/config_test.go +++ b/pkg/sink/mysql/config_test.go @@ -183,6 +183,22 @@ func TestGenerateDSNByConfig(t *testing.T) { testIsolationConfig() } +func TestSetDSNReadTimeout(t *testing.T) { + t.Parallel() + + dsnStr, err := setDSNReadTimeout( + "root:123456@tcp(127.0.0.1:4000)/?readTimeout=2m&writeTimeout=5m&timeout=3m", + "10m", + ) + require.NoError(t, err) + + dsn, err := dmysql.ParseDSN(dsnStr) + require.NoError(t, err) + require.Equal(t, 10*time.Minute, dsn.ReadTimeout) + require.Equal(t, 5*time.Minute, dsn.WriteTimeout) + require.Equal(t, 3*time.Minute, dsn.Timeout) +} + func TestConfigureControlDBConn(t *testing.T) { db, _, err := sqlmock.New() require.NoError(t, err) @@ -265,6 +281,83 @@ func TestApplySinkURIParamsToConfig(t *testing.T) { require.Equal(t, expected, cfg) } +func TestApplyAsyncDDLTimeout(t *testing.T) { + t.Parallel() + + newChangefeedConfig := func(mysqlConfig *config.MySQLConfig) *config.ChangefeedConfig { + return &config.ChangefeedConfig{ + TimeZone: "UTC", + SinkConfig: &config.SinkConfig{ + TiDBSourceID: 1, + MySQLConfig: mysqlConfig, + }, + } + } + + cases := []struct { + name string + uri string + mysqlConfig *config.MySQLConfig + expectedReadTimeout string + expectedAsyncDDLTimeout string + }{ + { + name: "default async ddl timeout", + uri: "mysql://127.0.0.1:3306/?read-timeout=4m", + expectedReadTimeout: "4m", + expectedAsyncDDLTimeout: defaultAsyncDDLTimeout, + }, + { + name: "sink uri async ddl timeout", + uri: "mysql://127.0.0.1:3306/?read-timeout=4m&async-ddl-timeout=30m", + expectedReadTimeout: "4m", + expectedAsyncDDLTimeout: "30m", + }, + { + name: "config async ddl timeout", + uri: "mysql://127.0.0.1:3306/?read-timeout=4m", + mysqlConfig: &config.MySQLConfig{ + AsyncDDLTimeout: util.AddressOf("20m"), + }, + expectedReadTimeout: "4m", + expectedAsyncDDLTimeout: "20m", + }, + { + name: "sink uri async ddl timeout overrides config", + uri: "mysql://127.0.0.1:3306/?read-timeout=4m&async-ddl-timeout=30m", + mysqlConfig: &config.MySQLConfig{ + AsyncDDLTimeout: util.AddressOf("20m"), + }, + expectedReadTimeout: "4m", + expectedAsyncDDLTimeout: "30m", + }, + { + name: "config read timeout does not change async ddl timeout", + uri: "mysql://127.0.0.1:3306/", + mysqlConfig: &config.MySQLConfig{ + ReadTimeout: util.AddressOf("6m"), + }, + expectedReadTimeout: "6m", + expectedAsyncDDLTimeout: defaultAsyncDDLTimeout, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + uri, err := url.Parse(tc.uri) + require.NoError(t, err) + cfg := New() + err = cfg.Apply(uri, common.NewChangefeedID4Test("default", "changefeed-01"), newChangefeedConfig(tc.mysqlConfig)) + require.NoError(t, err) + + require.Equal(t, tc.expectedReadTimeout, cfg.ReadTimeout) + require.Equal(t, tc.expectedAsyncDDLTimeout, cfg.AsyncDDLTimeout) + }) + } +} + func TestDefaultWorkerCountByDownstream(t *testing.T) { t.Parallel() @@ -394,6 +487,7 @@ func TestParseSinkURIBadQueryString(t *testing.T) { "mysql://127.0.0.1:3306/?write-timeout=badduration", "mysql://127.0.0.1:3306/?read-timeout=badduration", "mysql://127.0.0.1:3306/?timeout=badduration", + "mysql://127.0.0.1:3306/?async-ddl-timeout=badduration", } var uri *url.URL var err error diff --git a/pkg/sink/mysql/mysql_writer.go b/pkg/sink/mysql/mysql_writer.go index b8f98c3749..63cbcb4796 100644 --- a/pkg/sink/mysql/mysql_writer.go +++ b/pkg/sink/mysql/mysql_writer.go @@ -45,9 +45,19 @@ const ( // Writer is responsible for writing various dml events, ddl events, syncpoint events to mysql downstream. type Writer struct { +<<<<<<< HEAD id int ctx context.Context db *sql.DB +======= + id int + ctx context.Context + cancel context.CancelFunc + db *sql.DB + // asyncDB is used only by the TiDB ADD INDEX execution path, whose + // read timeout is intentionally independent from the regular DB. + asyncDB *sql.DB +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) cfg *Config ChangefeedID common.ChangeFeedID @@ -109,6 +119,11 @@ func (w *Writer) SetTableSchemaStore(tableSchemaStore *commonEvent.TableSchemaSt w.tableSchemaStore = tableSchemaStore } +// SetControlAsyncDB sets the DB pool used to execute TiDB ADD INDEX DDLs. +func (w *Writer) SetControlAsyncDB(db *sql.DB) { + w.asyncDB = db +} + func (w *Writer) FlushDDLEvent(event *commonEvent.DDLEvent) error { if w.cfg.IsTiDB { // first we check whether there is some async ddl executed now. diff --git a/pkg/sink/mysql/mysql_writer_ddl.go b/pkg/sink/mysql/mysql_writer_ddl.go index 7c49cb3ef6..7d714ee89c 100644 --- a/pkg/sink/mysql/mysql_writer_ddl.go +++ b/pkg/sink/mysql/mysql_writer_ddl.go @@ -15,6 +15,7 @@ package mysql import ( "context" + "database/sql" "fmt" "strconv" "strings" @@ -100,7 +101,8 @@ func (w *Writer) execDDL(event *commonEvent.DDLEvent) error { } }) - tx, err := w.db.BeginTx(ctx, nil) + db := w.getDDLExecDB(event) + tx, err := db.BeginTx(ctx, nil) if err != nil { return err } @@ -230,8 +232,13 @@ func (w *Writer) execDDLWithMaxRetries(event *commonEvent.DDLEvent) error { if w.cfg.IsTiDB && ddlCreateTime != "" && errors.Cause(err) == mysql.ErrInvalidConn { log.Warn("Wait the asynchronous ddl to synchronize", zap.Uint64("startTs", event.GetStartTs()), zap.Uint64("commitTs", event.GetCommitTs()), +<<<<<<< HEAD zap.String("ddlCreateTime", ddlCreateTime), zap.String("ddl", event.Query), zap.String("readTimeout", w.cfg.ReadTimeout), zap.Error(err)) +======= + zap.String("ddl", event.Query), zap.String("ddlCreateTime", ddlCreateTime), + zap.String("readTimeout", w.getDDLReadTimeout(event)), zap.Error(err)) +>>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) return w.waitDDLDone(w.ctx, event, ddlCreateTime) } log.Warn("Execute DDL with error, retry later", @@ -250,6 +257,26 @@ func (w *Writer) execDDLWithMaxRetries(event *commonEvent.DDLEvent) error { retry.WithIsRetryableErr(errors.IsRetryableDDLError)) } +func (w *Writer) getDDLExecDB(event *commonEvent.DDLEvent) *sql.DB { + if w.useAsyncDB(event) { + return w.asyncDB + } + return w.db +} + +func (w *Writer) getDDLReadTimeout(event *commonEvent.DDLEvent) string { + if w.useAsyncDB(event) { + return w.cfg.AsyncDDLTimeout + } + return w.cfg.ReadTimeout +} + +func (w *Writer) useAsyncDB(event *commonEvent.DDLEvent) bool { + return w.asyncDB != nil && + w.cfg.IsTiDB && + event.GetDDLType() == timodel.ActionAddIndex +} + // waitDDLDone wait current ddl func (w *Writer) waitDDLDone(ctx context.Context, ddl *commonEvent.DDLEvent, ddlCreateTime string) error { ticker := time.NewTicker(5 * time.Second) diff --git a/pkg/sink/mysql/mysql_writer_test.go b/pkg/sink/mysql/mysql_writer_test.go index eb43a439de..585c53cf83 100644 --- a/pkg/sink/mysql/mysql_writer_test.go +++ b/pkg/sink/mysql/mysql_writer_test.go @@ -512,6 +512,79 @@ func TestWaitAsyncDDLDone_CreateTableLikeShouldQueryDownstreamAddIndexJob(t *tes require.NoError(t, mock.ExpectationsWereMet()) } +func TestExecDDLUsesControlAsyncDBOnlyForTiDBAddIndex(t *testing.T) { + writer, controlDB, controlMock := newTestMysqlWriterForTiDB(t) + defer controlDB.Close() + + controlAsyncDB, controlAsyncMock := newTestMockDB(t) + defer controlAsyncDB.Close() + writer.SetControlAsyncDB(controlAsyncDB) + writer.cfg.ReadTimeout = "2m" + writer.cfg.AsyncDDLTimeout = "30m" + + addIndexEvent := &commonEvent.DDLEvent{ + Type: byte(timodel.ActionAddIndex), + Query: "alter table t add index idx_name(name);", + SchemaName: "test", + TableName: "t", + } + controlAsyncMock.ExpectBegin() + controlAsyncMock.ExpectExec("USE `test`;").WillReturnResult(sqlmock.NewResult(1, 1)) + controlAsyncMock.ExpectExec("SET TIMESTAMP = DEFAULT").WillReturnResult(sqlmock.NewResult(1, 1)) + controlAsyncMock.ExpectExec("alter table t add index idx_name(name);").WillReturnResult(sqlmock.NewResult(1, 1)) + controlAsyncMock.ExpectCommit() + + require.NoError(t, writer.execDDL(addIndexEvent)) + require.Equal(t, "30m", writer.getDDLReadTimeout(addIndexEvent)) + require.NoError(t, controlAsyncMock.ExpectationsWereMet()) + require.NoError(t, controlMock.ExpectationsWereMet()) + + addColumnEvent := &commonEvent.DDLEvent{ + Type: byte(timodel.ActionAddColumn), + Query: "alter table t add column age int;", + SchemaName: "test", + TableName: "t", + } + controlMock.ExpectBegin() + controlMock.ExpectExec("USE `test`;").WillReturnResult(sqlmock.NewResult(1, 1)) + controlMock.ExpectExec("SET TIMESTAMP = DEFAULT").WillReturnResult(sqlmock.NewResult(1, 1)) + controlMock.ExpectExec("alter table t add column age int;").WillReturnResult(sqlmock.NewResult(1, 1)) + controlMock.ExpectCommit() + + require.NoError(t, writer.execDDL(addColumnEvent)) + require.Equal(t, "2m", writer.getDDLReadTimeout(addColumnEvent)) + require.NoError(t, controlMock.ExpectationsWereMet()) + require.NoError(t, controlAsyncMock.ExpectationsWereMet()) +} + +func TestExecDDLUsesControlDBForMySQLAddIndex(t *testing.T) { + writer, controlDB, controlMock := newTestMysqlWriter(t) + defer controlDB.Close() + + controlAsyncDB, controlAsyncMock := newTestMockDB(t) + defer controlAsyncDB.Close() + writer.SetControlAsyncDB(controlAsyncDB) + writer.cfg.ReadTimeout = "2m" + writer.cfg.AsyncDDLTimeout = "30m" + + addIndexEvent := &commonEvent.DDLEvent{ + Type: byte(timodel.ActionAddIndex), + Query: "alter table t add index idx_name(name);", + SchemaName: "test", + TableName: "t", + } + controlMock.ExpectBegin() + controlMock.ExpectExec("USE `test`;").WillReturnResult(sqlmock.NewResult(1, 1)) + controlMock.ExpectExec("SET TIMESTAMP = DEFAULT").WillReturnResult(sqlmock.NewResult(1, 1)) + controlMock.ExpectExec("alter table t add index idx_name(name);").WillReturnResult(sqlmock.NewResult(1, 1)) + controlMock.ExpectCommit() + + require.NoError(t, writer.execDDL(addIndexEvent)) + require.Equal(t, "2m", writer.getDDLReadTimeout(addIndexEvent)) + require.NoError(t, controlMock.ExpectationsWereMet()) + require.NoError(t, controlAsyncMock.ExpectationsWereMet()) +} + // Test the async ddl can be write successfully func TestMysqlWriter_AsyncDDL(t *testing.T) { writer, db, mock := newTestMysqlWriterForTiDB(t) From 3f7f8139d529049bdf9182200d625e01fdbd404a Mon Sep 17 00:00:00 2001 From: nhsmw Date: Thu, 13 Aug 2026 14:14:15 +0800 Subject: [PATCH 2/4] Increase defaultAsyncDDLTimeout to 2 minutes --- pkg/sink/mysql/config.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/sink/mysql/config.go b/pkg/sink/mysql/config.go index cfc07815c2..179778c6f2 100644 --- a/pkg/sink/mysql/config.go +++ b/pkg/sink/mysql/config.go @@ -69,7 +69,7 @@ const ( defaultReadTimeout = "2m" defaultWriteTimeout = "2m" defaultDialTimeout = "2m" - defaultAsyncDDLTimeout = "10s" + defaultAsyncDDLTimeout = "2m" defaultSafeMode = false defaultTxnIsolationRC = "READ-COMMITTED" defaultCharacterSet = "utf8mb4" From 8764d65ec022a3e6af10b4a4840b9135dbba2003 Mon Sep 17 00:00:00 2001 From: wk989898 Date: Fri, 14 Aug 2026 10:24:35 +0000 Subject: [PATCH 3/4] update Signed-off-by: wk989898 --- api/v2/model.go | 18 ------- downstreamadapter/sink/mysql/sink.go | 61 +++-------------------- downstreamadapter/sink/mysql/sink_test.go | 20 +++----- pkg/sink/mysql/config.go | 39 +-------------- pkg/sink/mysql/mysql_writer.go | 6 --- pkg/sink/mysql/mysql_writer_ddl.go | 5 -- 6 files changed, 16 insertions(+), 133 deletions(-) diff --git a/api/v2/model.go b/api/v2/model.go index 4db4c3d33b..000b9660e7 100644 --- a/api/v2/model.go +++ b/api/v2/model.go @@ -1475,23 +1475,6 @@ type KafkaConfig struct { // MySQLConfig represents a MySQL sink configuration type MySQLConfig struct { -<<<<<<< HEAD - WorkerCount *int `json:"worker_count,omitempty"` - MaxTxnRow *int `json:"max_txn_row,omitempty"` - MaxMultiUpdateRowSize *int `json:"max_multi_update_row_size,omitempty"` - MaxMultiUpdateRowCount *int `json:"max_multi_update_row_count,omitempty"` - TiDBTxnMode *string `json:"tidb_txn_mode,omitempty"` - SSLCa *string `json:"ssl_ca,omitempty"` - SSLCert *string `json:"ssl_cert,omitempty"` - SSLKey *string `json:"ssl_key,omitempty"` - TimeZone *string `json:"time_zone,omitempty"` - WriteTimeout *string `json:"write_timeout,omitempty"` - ReadTimeout *string `json:"read_timeout,omitempty"` - Timeout *string `json:"timeout,omitempty"` - EnableBatchDML *bool `json:"enable_batch_dml,omitempty"` - EnableMultiStatement *bool `json:"enable_multi_statement,omitempty"` - EnableCachePreparedStatement *bool `json:"enable_cache_prepared_statement,omitempty"` -======= WorkerCount *int `json:"worker_count,omitempty" toml:"worker-count,omitempty"` MaxTxnRow *int `json:"max_txn_row,omitempty" toml:"max-txn-row,omitempty"` MaxMultiUpdateRowSize *int `json:"max_multi_update_row_size,omitempty" toml:"max-multi-update-row-size,omitempty"` @@ -1508,7 +1491,6 @@ type MySQLConfig struct { EnableBatchDML *bool `json:"enable_batch_dml,omitempty" toml:"enable-batch-dml,omitempty"` EnableMultiStatement *bool `json:"enable_multi_statement,omitempty" toml:"enable-multi-statement,omitempty"` EnableCachePreparedStatement *bool `json:"enable_cache_prepared_statement,omitempty" toml:"enable-cache-prepared-statement,omitempty"` ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } // CloudStorageConfig represents a cloud storage sink configuration diff --git a/downstreamadapter/sink/mysql/sink.go b/downstreamadapter/sink/mysql/sink.go index abd8e3cfe0..cb24956c33 100644 --- a/downstreamadapter/sink/mysql/sink.go +++ b/downstreamadapter/sink/mysql/sink.go @@ -106,11 +106,7 @@ func New( metrics.ChangefeedDownstreamIsTiDBGauge.DeleteLabelValues(keyspace, name) } -<<<<<<< HEAD - return newMySQLSinkWithControlDB(ctx, changefeedID, cfg, dmlDB, controlDB, config.BDRMode), nil -======= - return newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, config.BDRMode, config.EnableActiveActive, config.ActiveActiveProgressInterval, keyspaceID), nil ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) + return newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, config.BDRMode), nil } // NewMySQLSink used for test @@ -121,15 +117,11 @@ func NewMySQLSink( db *sql.DB, bdrMode bool, ) *Sink { -<<<<<<< HEAD - return newMySQLSinkWithControlDB(ctx, changefeedID, cfg, db, db, bdrMode) -======= var controlAsyncDB *sql.DB if cfg.IsTiDB { controlAsyncDB = db } - return newMySQLSinkWithDBs(ctx, changefeedID, cfg, db, db, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, db, db, controlAsyncDB, bdrMode) } // newMySQLSinkWithControlDB creates a MySQL sink with separate pools for DML and @@ -143,14 +135,11 @@ func newMySQLSinkWithControlDB( controlDB *sql.DB, bdrMode bool, ) *Sink { -<<<<<<< HEAD - stat := metrics.NewStatistics(changefeedID, "TxnSink") -======= var controlAsyncDB *sql.DB if cfg.IsTiDB { controlAsyncDB = controlDB } - return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, bdrMode) } func newMySQLSinkWithControlAsyncDB( @@ -161,11 +150,8 @@ func newMySQLSinkWithControlAsyncDB( controlDB *sql.DB, controlAsyncDB *sql.DB, bdrMode bool, - enableActiveActive bool, - progressInterval time.Duration, - keyspaceID uint32, ) *Sink { - return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, bdrMode) } func newMySQLSinkWithDBs( @@ -176,9 +162,6 @@ func newMySQLSinkWithDBs( controlDB *sql.DB, controlAsyncDB *sql.DB, bdrMode bool, - enableActiveActive bool, - progressInterval time.Duration, - keyspaceID uint32, ) *Sink { if !cfg.IsTiDB { controlAsyncDB = nil @@ -186,26 +169,8 @@ func newMySQLSinkWithDBs( controlAsyncDB = controlDB } - stat := metrics.NewStatistics(changefeedID, keyspaceID, "TxnSink") - - var activeActiveSyncStatsCollector *mysql.ActiveActiveSyncStatsCollector - if enableActiveActive && cfg.IsTiDB && cfg.ActiveActiveSyncStatsInterval > 0 { - supported, err := mysql.CheckActiveActiveSyncStatsSupported(ctx, dmlDB) - if err != nil { - log.Info("failed to check tidb_cdc_active_active_sync_stats support, disable metric collection", - zap.String("keyspace", changefeedID.Keyspace()), - zap.Stringer("changefeed", changefeedID), - zap.Error(err)) - } else if supported { - activeActiveSyncStatsCollector = mysql.NewActiveActiveSyncStatsCollector(changefeedID) - } else { - log.Info("downstream does not support tidb_cdc_active_active_sync_stats, disable metric collection", - zap.String("keyspace", changefeedID.Keyspace()), - zap.Stringer("changefeed", changefeedID)) - } - } + stat := metrics.NewStatistics(changefeedID, "TxnSink") ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) result := &Sink{ changefeedID: changefeedID, dmlDB: dmlDB, @@ -226,18 +191,10 @@ func newMySQLSinkWithDBs( bdrMode: bdrMode, } for i := 0; i < len(result.dmlWriter); i++ { -<<<<<<< HEAD result.dmlWriter[i] = mysql.NewWriter(ctx, i, dmlDB, cfg, changefeedID, stat) -======= - result.dmlWriter[i] = mysql.NewWriter(ctx, i, dmlDB, cfg, changefeedID, stat, activeActiveSyncStatsCollector) - } - result.ddlWriter = mysql.NewWriter(ctx, len(result.dmlWriter), controlDB, cfg, changefeedID, stat, nil) - result.ddlWriter.SetControlAsyncDB(controlAsyncDB) - if enableActiveActive { - result.progressTableWriter = mysql.NewProgressTableWriter(ctx, controlDB, changefeedID, cfg.MaxTxnRow, progressInterval) ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } result.ddlWriter = mysql.NewWriter(ctx, len(result.dmlWriter), controlDB, cfg, changefeedID, stat) + result.ddlWriter.SetControlAsyncDB(controlAsyncDB) return result } @@ -458,15 +415,9 @@ func (s *Sink) Close() { if s.controlDB != s.dmlDB { s.closeDBPool("control", s.controlDB) } -<<<<<<< HEAD -======= if s.controlAsyncDB != nil && s.controlAsyncDB != s.dmlDB && s.controlAsyncDB != s.controlDB { s.closeDBPool("control async", s.controlAsyncDB) } - if s.activeActiveSyncStatsCollector != nil { - s.activeActiveSyncStatsCollector.Close() - } ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) s.statistics.Close() metrics.ChangefeedDownstreamIsTiDBGauge.DeleteLabelValues(s.changefeedID.Keyspace(), s.changefeedID.Name()) diff --git a/downstreamadapter/sink/mysql/sink_test.go b/downstreamadapter/sink/mysql/sink_test.go index 85eb0f002b..036e1c6d4b 100644 --- a/downstreamadapter/sink/mysql/sink_test.go +++ b/downstreamadapter/sink/mysql/sink_test.go @@ -25,12 +25,8 @@ import ( "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/sink/mysql" -<<<<<<< HEAD + "github.com/pingcap/tidb/pkg/meta/model" "github.com/pingcap/tidb/pkg/sessionctx/variable" -======= - timodel "github.com/pingcap/tidb/pkg/meta/model" - "github.com/pingcap/tidb/pkg/sessionctx/vardef" ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) "github.com/stretchr/testify/require" ) @@ -94,11 +90,11 @@ func TestMysqlSinkControlAsyncDBOnlyForTiDB(t *testing.T) { cfg := mysql.New() cfg.WorkerCount = 1 - cfg.MaxAllowedPacket = int64(vardef.DefMaxAllowedPacket) + cfg.MaxAllowedPacket = int64(variable.DefMaxAllowedPacket) cfg.CachePrepStmts = false cfg.IsTiDB = false - sink := NewMySQLSink(ctx, changefeedID, cfg, db, false, false, time.Minute, common.DefaultKeyspaceID) + sink := NewMySQLSink(ctx, changefeedID, cfg, db, false) require.Nil(t, sink.controlAsyncDB) mock.ExpectClose() @@ -112,11 +108,11 @@ func TestMysqlSinkControlAsyncDBOnlyForTiDB(t *testing.T) { cfg := mysql.New() cfg.WorkerCount = 1 - cfg.MaxAllowedPacket = int64(vardef.DefMaxAllowedPacket) + cfg.MaxAllowedPacket = int64(variable.DefMaxAllowedPacket) cfg.CachePrepStmts = false cfg.IsTiDB = true - sink := NewMySQLSink(ctx, changefeedID, cfg, db, false, false, time.Minute, common.DefaultKeyspaceID) + sink := NewMySQLSink(ctx, changefeedID, cfg, db, false) require.Same(t, db, sink.controlAsyncDB) mock.ExpectClose() @@ -137,15 +133,15 @@ func TestMysqlSinkUsesControlAsyncDBForTiDBAddIndex(t *testing.T) { changefeedID := common.NewChangefeedID4Test("test", "test") cfg := mysql.New() cfg.WorkerCount = 1 - cfg.MaxAllowedPacket = int64(vardef.DefMaxAllowedPacket) + cfg.MaxAllowedPacket = int64(variable.DefMaxAllowedPacket) cfg.CachePrepStmts = false cfg.EnableDDLTs = false cfg.IsTiDB = true - sink := newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, false, false, time.Minute, common.DefaultKeyspaceID) + sink := newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, false) ddl := &commonEvent.DDLEvent{ - Type: byte(timodel.ActionAddIndex), + Type: byte(model.ActionAddIndex), Query: "alter table t add index idx_name(name);", SchemaName: "test", TableName: "t", diff --git a/pkg/sink/mysql/config.go b/pkg/sink/mysql/config.go index 179778c6f2..34cd7280d1 100644 --- a/pkg/sink/mysql/config.go +++ b/pkg/sink/mysql/config.go @@ -25,11 +25,11 @@ import ( dmysql "github.com/go-sql-driver/mysql" lru "github.com/hashicorp/golang-lru" - "github.com/pingcap/errors" "github.com/pingcap/failpoint" "github.com/pingcap/log" "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/config" + "github.com/pingcap/ticdc/pkg/errors" cerror "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/security" "github.com/pingcap/ticdc/pkg/sink/sqlmodel" @@ -108,17 +108,6 @@ type Config struct { MaxMultiUpdateRowCount int MaxMultiUpdateRowSize int TidbTxnMode string -<<<<<<< HEAD - ReadTimeout string - WriteTimeout string - DialTimeout string - SafeMode bool - Timezone string - TLS string - SSLCa string - SSLCert string - SSLKey string -======= // tidbTxnModeSpecified indicates whether TidbTxnMode is explicitly set by user via sink URI or changefeed config. // It is used to avoid overriding user configuration when applying downstream-specific defaults. tidbTxnModeSpecified bool @@ -134,7 +123,6 @@ type Config struct { SSLCa string SSLCert string SSLKey string ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) // retry number for dml DMLMaxRetry uint64 @@ -182,7 +170,6 @@ type Config struct { // New returns the default mysql backend config. func New() *Config { return &Config{ -<<<<<<< HEAD WorkerCount: DefaultTiDBWorkerCount, workerCountSpecified: false, MaxTxnRow: DefaultMaxTxnRow, @@ -192,6 +179,7 @@ func New() *Config { ReadTimeout: defaultReadTimeout, WriteTimeout: defaultWriteTimeout, DialTimeout: defaultDialTimeout, + AsyncDDLTimeout: defaultAsyncDDLTimeout, SafeMode: defaultSafeMode, BatchDMLEnable: defaultBatchDMLEnable, MultiStmtEnable: defaultMultiStmtEnable, @@ -202,29 +190,6 @@ func New() *Config { EnableDDLTs: defaultEnableDDLTs, SlowQuery: slowQuery, whereClause: sqlmodel.DefaultWhereClause, -======= - WorkerCount: DefaultTiDBWorkerCount, - workerCountSpecified: false, - MaxTxnRow: DefaultMaxTxnRow, - MaxMultiUpdateRowCount: defaultMaxMultiUpdateRowCount, - MaxMultiUpdateRowSize: defaultMaxMultiUpdateRowSize, - TidbTxnMode: defaultTiDBTxnMode, - ReadTimeout: defaultReadTimeout, - WriteTimeout: defaultWriteTimeout, - DialTimeout: defaultDialTimeout, - AsyncDDLTimeout: defaultAsyncDDLTimeout, - SafeMode: defaultSafeMode, - BatchDMLEnable: defaultBatchDMLEnable, - MultiStmtEnable: defaultMultiStmtEnable, - CachePrepStmts: defaultCachePrepStmts, - SourceID: config.DefaultTiDBSourceID, - DMLMaxRetry: 8, - HasVectorType: defaultHasVectorType, - EnableDDLTs: defaultEnableDDLTs, - SlowQuery: slowQuery, - ActiveActiveSyncStatsInterval: time.Minute, - whereClause: sqlmodel.DefaultWhereClause, ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) } } diff --git a/pkg/sink/mysql/mysql_writer.go b/pkg/sink/mysql/mysql_writer.go index 63cbcb4796..923bed9886 100644 --- a/pkg/sink/mysql/mysql_writer.go +++ b/pkg/sink/mysql/mysql_writer.go @@ -45,11 +45,6 @@ const ( // Writer is responsible for writing various dml events, ddl events, syncpoint events to mysql downstream. type Writer struct { -<<<<<<< HEAD - id int - ctx context.Context - db *sql.DB -======= id int ctx context.Context cancel context.CancelFunc @@ -57,7 +52,6 @@ type Writer struct { // asyncDB is used only by the TiDB ADD INDEX execution path, whose // read timeout is intentionally independent from the regular DB. asyncDB *sql.DB ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) cfg *Config ChangefeedID common.ChangeFeedID diff --git a/pkg/sink/mysql/mysql_writer_ddl.go b/pkg/sink/mysql/mysql_writer_ddl.go index 7d714ee89c..16560b43a5 100644 --- a/pkg/sink/mysql/mysql_writer_ddl.go +++ b/pkg/sink/mysql/mysql_writer_ddl.go @@ -232,13 +232,8 @@ func (w *Writer) execDDLWithMaxRetries(event *commonEvent.DDLEvent) error { if w.cfg.IsTiDB && ddlCreateTime != "" && errors.Cause(err) == mysql.ErrInvalidConn { log.Warn("Wait the asynchronous ddl to synchronize", zap.Uint64("startTs", event.GetStartTs()), zap.Uint64("commitTs", event.GetCommitTs()), -<<<<<<< HEAD - zap.String("ddlCreateTime", ddlCreateTime), zap.String("ddl", event.Query), - zap.String("readTimeout", w.cfg.ReadTimeout), zap.Error(err)) -======= zap.String("ddl", event.Query), zap.String("ddlCreateTime", ddlCreateTime), zap.String("readTimeout", w.getDDLReadTimeout(event)), zap.Error(err)) ->>>>>>> 430b0a8cc (sink: add async ddl timeout for add index (#5836)) return w.waitDDLDone(w.ctx, event, ddlCreateTime) } log.Warn("Execute DDL with error, retry later", From 71e44fc6a11b9748bf23e262279e8b6e5f572a16 Mon Sep 17 00:00:00 2001 From: wk989898 Date: Fri, 14 Aug 2026 10:53:33 +0000 Subject: [PATCH 4/4] update comment Signed-off-by: wk989898 --- pkg/common/table_info_shared_schema_guard_test.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pkg/common/table_info_shared_schema_guard_test.go b/pkg/common/table_info_shared_schema_guard_test.go index 705bd35df4..ffe29b1ddf 100644 --- a/pkg/common/table_info_shared_schema_guard_test.go +++ b/pkg/common/table_info_shared_schema_guard_test.go @@ -66,7 +66,7 @@ type goSourceFile struct { // whose shared-schema compatibility has been reviewed and recorded here. // New upstream fields should be reviewed and then added here intentionally. func TestLatestTiDBTableInfoSharedSchemaGuard(t *testing.T) { - // This guard intentionally checks the latest TiDB master to detect + // This guard intentionally checks the latest TiDB release-8.5 to detect // upstream struct-field changes before TiCDC upgrades its pinned TiDB version. cases := []structGuardCase{ { @@ -205,12 +205,12 @@ func buildRequiredTypesByPackage(cases []structGuardCase) map[packageKey][]strin return requiredTypesByPackage } -// queryModuleInfo resolves a module at @master and returns location/version +// queryModuleInfo resolves a module at @release-8.5 and returns location/version // information for source-level contract checks. func queryModuleInfo(t *testing.T, modulePath string) *moduleInfo { t.Helper() - cmd := exec.Command("go", "mod", "download", "-json", modulePath+"@master") + cmd := exec.Command("go", "mod", "download", "-json", modulePath+"@release-8.5") output, cmdErr := cmd.CombinedOutput() mod, decodeErr := decodeModuleInfo(output) require.NoError(t, decodeErr, "unmarshal module metadata for %s failed", modulePath)