diff --git a/api/v2/model.go b/api/v2/model.go index 0575a73e23..69841c9909 100644 --- a/api/v2/model.go +++ b/api/v2/model.go @@ -514,6 +514,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, @@ -845,6 +846,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, @@ -1531,6 +1533,7 @@ type MySQLConfig struct { 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"` diff --git a/api/v2/model_test.go b/api/v2/model_test.go index 88fb821cdb..31867ab2d2 100644 --- a/api/v2/model_test.go +++ b/api/v2/model_test.go @@ -207,3 +207,23 @@ func TestReplicaConfigConversionRedoBatchField(t *testing.T) { require.NotNil(t, apiCfgBack.Consistent.EventCollectorBatchCount) require.Equal(t, 4096, *apiCfgBack.Consistent.EventCollectorBatchCount) } + +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 2856d5fce4..54d32fec67 100644 --- a/downstreamadapter/sink/mysql/sink.go +++ b/downstreamadapter/sink/mysql/sink.go @@ -50,11 +50,12 @@ type Sink struct { // enableActiveActive is false. progressTableWriter *mysql.ProgressTableWriter - // 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 @@ -81,12 +82,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 } @@ -97,7 +101,7 @@ func New( sinkURI *url.URL, keyspaceID uint32, ) (*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 } @@ -113,9 +117,10 @@ func New( metrics.ChangefeedDownstreamIsTiDBGauge.DeleteLabelValues(keyspace, name) } - return newMySQLSinkWithControlDB(ctx, changefeedID, cfg, dmlDB, controlDB, config.BDRMode, config.EnableActiveActive, config.ActiveActiveProgressInterval, keyspaceID), nil + return newMySQLSinkWithControlAsyncDB(ctx, changefeedID, cfg, dmlDB, controlDB, controlAsyncDB, config.BDRMode, config.EnableActiveActive, config.ActiveActiveProgressInterval, keyspaceID), nil } +// NewMySQLSink used for test func NewMySQLSink( ctx context.Context, changefeedID common.ChangeFeedID, @@ -126,7 +131,11 @@ func NewMySQLSink( progressInterval time.Duration, keyspaceID uint32, ) *Sink { - return newMySQLSinkWithDBs(ctx, changefeedID, cfg, db, db, bdrMode, enableActiveActive, progressInterval, keyspaceID) + var controlAsyncDB *sql.DB + if cfg.IsTiDB { + controlAsyncDB = db + } + return newMySQLSinkWithDBs(ctx, changefeedID, cfg, db, db, controlAsyncDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) } // newMySQLSinkWithControlDB creates a MySQL sink with separate pools for DML and @@ -144,7 +153,26 @@ func newMySQLSinkWithControlDB( progressInterval time.Duration, keyspaceID uint32, ) *Sink { - return newMySQLSinkWithDBs(ctx, changefeedID, cfg, dmlDB, controlDB, bdrMode, enableActiveActive, progressInterval, keyspaceID) + 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( @@ -153,11 +181,18 @@ func newMySQLSinkWithDBs( 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 @@ -178,11 +213,12 @@ func newMySQLSinkWithDBs( } 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, @@ -201,6 +237,7 @@ func newMySQLSinkWithDBs( 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) } @@ -453,6 +490,9 @@ func (s *Sink) Close() { if s.controlDB != s.dmlDB { s.closeDBPool("control", s.controlDB) } + 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() } diff --git a/downstreamadapter/sink/mysql/sink_test.go b/downstreamadapter/sink/mysql/sink_test.go index d1da54d9e7..b80e25fc7a 100644 --- a/downstreamadapter/sink/mysql/sink_test.go +++ b/downstreamadapter/sink/mysql/sink_test.go @@ -25,6 +25,7 @@ import ( "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/sink/mysql" + timodel "github.com/pingcap/tidb/pkg/meta/model" "github.com/pingcap/tidb/pkg/sessionctx/vardef" "github.com/stretchr/testify/require" ) @@ -79,6 +80,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 dca4d50572..a0f415a34d 100644 --- a/pkg/config/sink.go +++ b/pkg/config/sink.go @@ -721,6 +721,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 ca0ec2bd81..b69213fd3b 100644 --- a/pkg/sink/mysql/config.go +++ b/pkg/sink/mysql/config.go @@ -64,13 +64,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 @@ -112,12 +113,15 @@ type Config struct { ReadTimeout string WriteTimeout string DialTimeout string - SafeMode bool - Timezone string - TLS string - SSLCa string - SSLCert string - SSLKey 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 // retry number for dml DMLMaxRetry uint64 @@ -180,6 +184,7 @@ func New() *Config { ReadTimeout: defaultReadTimeout, WriteTimeout: defaultWriteTimeout, DialTimeout: defaultDialTimeout, + AsyncDDLTimeout: defaultAsyncDDLTimeout, SafeMode: defaultSafeMode, BatchDMLEnable: defaultBatchDMLEnable, MultiStmtEnable: defaultMultiStmtEnable, @@ -217,6 +222,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) @@ -277,6 +283,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 } @@ -321,16 +330,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) @@ -339,10 +348,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( @@ -433,6 +474,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 ee437a8a06..0f8170c2ad 100644 --- a/pkg/sink/mysql/config_test.go +++ b/pkg/sink/mysql/config_test.go @@ -232,6 +232,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) @@ -315,6 +331,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() @@ -444,6 +537,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 4c38fadaab..fc94a5bee4 100644 --- a/pkg/sink/mysql/mysql_writer.go +++ b/pkg/sink/mysql/mysql_writer.go @@ -46,10 +46,13 @@ const ( // Writer is responsible for writing various dml events, ddl events, syncpoint events to mysql downstream. type Writer struct { - id int - ctx context.Context - cancel context.CancelFunc - 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 cfg *Config ChangefeedID common.ChangeFeedID @@ -133,6 +136,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 8424bdbdbd..d0774de330 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 } @@ -231,7 +233,7 @@ func (w *Writer) execDDLWithMaxRetries(event *commonEvent.DDLEvent) error { log.Warn("Wait the asynchronous ddl to synchronize", zap.Uint64("startTs", event.GetStartTs()), zap.Uint64("commitTs", event.GetCommitTs()), zap.String("ddl", event.Query), zap.String("ddlCreateTime", ddlCreateTime), - zap.String("readTimeout", w.cfg.ReadTimeout), zap.Error(err)) + zap.String("readTimeout", w.getDDLReadTimeout(event)), zap.Error(err)) return w.waitDDLDone(w.ctx, event, ddlCreateTime) } log.Warn("Execute DDL with error, retry later", @@ -250,6 +252,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 00c37f9f85..3814ff373c 100644 --- a/pkg/sink/mysql/mysql_writer_test.go +++ b/pkg/sink/mysql/mysql_writer_test.go @@ -553,6 +553,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)