From be1c200de9f765dd6412f34c5e0dbb5c617303ba Mon Sep 17 00:00:00 2001 From: wk989898 Date: Mon, 20 Apr 2026 07:54:33 +0000 Subject: [PATCH 1/3] init Signed-off-by: wk989898 --- cmd/storage-consumer/consumer.go | 54 ++++++++++++++++++++++++++++++-- 1 file changed, 52 insertions(+), 2 deletions(-) diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index e3ee2d6eb5..0cce03b1c9 100644 --- a/cmd/storage-consumer/consumer.go +++ b/cmd/storage-consumer/consumer.go @@ -17,6 +17,7 @@ import ( "context" "encoding/json" "fmt" + "regexp" "sort" "strings" "time" @@ -35,6 +36,9 @@ import ( "github.com/pingcap/ticdc/pkg/sink/codec/csv" putil "github.com/pingcap/ticdc/pkg/util" "github.com/pingcap/tidb/br/pkg/storage" + timodel "github.com/pingcap/tidb/pkg/meta/model" + "github.com/pingcap/tidb/pkg/parser" + "github.com/pingcap/tidb/pkg/parser/ast" "go.uber.org/atomic" "go.uber.org/zap" "golang.org/x/sync/errgroup" @@ -46,6 +50,8 @@ const ( fakePartitionNumForSchemaFile = -1 ) +var renameTableQueryRe = regexp.MustCompile(`(?is)^rename\s+table\s+(.+?)\s+to\s+(.+?)$`) + type ( fileIndexRange map[cloudstorage.FileIndexKey]indexRange fileIndexKeyMap map[cloudstorage.FileIndexKey]uint64 @@ -510,6 +516,49 @@ func (c *consumer) mustGetTableDef(key cloudstorage.SchemaPathKey) cloudstorage. return *tableDef } +func getRenameTableOldTableKey(tableDef cloudstorage.TableDefinition) (string, bool) { + if tableDef.Type != byte(timodel.ActionRenameTable) { + return "", false + } + schemaName := tableDef.Schema + tableName := tableDef.Table + stmt, err := parser.New().ParseOneStmt(tableDef.Query, "", "") + if err != nil { + log.Panic("parse statement failed", zap.Any("DDL", tableDef.Query), zap.Error(err)) + } + // The query in job maybe "RENAME TABLE table1 to table2" + oldSchemaName := stmt.(*ast.RenameTableStmt).TableToTables[0].OldTable.Schema.O + if oldSchemaName != "" { + schemaName = oldSchemaName + } + tableName = stmt.(*ast.RenameTableStmt).TableToTables[0].OldTable.Name.O + return commonType.QuoteSchema(schemaName, tableName), true +} + +func getDDLWatermarkKeys(tableDef cloudstorage.TableDefinition) []string { + keys := make(map[string]struct{}, 2) + keys[commonType.QuoteSchema(tableDef.Schema, tableDef.Table)] = struct{}{} + if oldTableKey, ok := getRenameTableOldTableKey(tableDef); ok { + keys[oldTableKey] = struct{}{} + } + + res := make([]string, 0, len(keys)) + for key := range keys { + res = append(res, key) + } + return res +} + +func (c *consumer) updateTableDDLWatermark(tableDef cloudstorage.TableDefinition) []string { + keys := getDDLWatermarkKeys(tableDef) + for _, key := range keys { + if c.tableDDLWatermark[key] < tableDef.TableVersion { + c.tableDDLWatermark[key] = tableDef.TableVersion + } + } + return keys +} + func (c *consumer) handleNewFiles( ctx context.Context, dmlFileMap map[cloudstorage.DmlPathKey]fileIndexRange, @@ -582,12 +631,13 @@ func (c *consumer) handleNewFiles( if err := c.sink.WriteBlockEvent(ddlEvent); err != nil { return errors.Trace(err) } - c.tableDDLWatermark[tableKey] = key.TableVersion + watermarkKeys := c.updateTableDDLWatermark(tableDef) // TODO: need to cleanup tableDefMap in the future. log.Info("execute ddl event successfully", zap.String("query", tableDef.Query), zap.String("schema", key.Schema), zap.String("table", key.Table), - zap.Uint64("ddlWatermark", c.tableDDLWatermark[tableKey])) + zap.Uint64("ddlWatermark", c.tableDDLWatermark[tableKey]), + zap.Strings("watermarkKeys", watermarkKeys)) continue } From 4a2993d0070edd33fb5cf68567d24b0c35cc8547 Mon Sep 17 00:00:00 2001 From: nhsmw Date: Mon, 20 Apr 2026 16:04:49 +0800 Subject: [PATCH 2/3] Update consumer.go --- cmd/storage-consumer/consumer.go | 34 ++++++++++---------------------- 1 file changed, 10 insertions(+), 24 deletions(-) diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index 0cce03b1c9..8b4bd0337e 100644 --- a/cmd/storage-consumer/consumer.go +++ b/cmd/storage-consumer/consumer.go @@ -17,7 +17,6 @@ import ( "context" "encoding/json" "fmt" - "regexp" "sort" "strings" "time" @@ -50,8 +49,6 @@ const ( fakePartitionNumForSchemaFile = -1 ) -var renameTableQueryRe = regexp.MustCompile(`(?is)^rename\s+table\s+(.+?)\s+to\s+(.+?)$`) - type ( fileIndexRange map[cloudstorage.FileIndexKey]indexRange fileIndexKeyMap map[cloudstorage.FileIndexKey]uint64 @@ -535,28 +532,17 @@ func getRenameTableOldTableKey(tableDef cloudstorage.TableDefinition) (string, b return commonType.QuoteSchema(schemaName, tableName), true } -func getDDLWatermarkKeys(tableDef cloudstorage.TableDefinition) []string { - keys := make(map[string]struct{}, 2) - keys[commonType.QuoteSchema(tableDef.Schema, tableDef.Table)] = struct{}{} - if oldTableKey, ok := getRenameTableOldTableKey(tableDef); ok { - keys[oldTableKey] = struct{}{} - } - - res := make([]string, 0, len(keys)) - for key := range keys { - res = append(res, key) +func (c *consumer) updateTableDDLWatermark(tableDef cloudstorage.TableDefinition) string { + key := commonType.QuoteSchema(tableDef.Schema, tableDef.Table) + if c.tableDDLWatermark[key] < tableDef.TableVersion { + c.tableDDLWatermark[key] = tableDef.TableVersion } - return res -} - -func (c *consumer) updateTableDDLWatermark(tableDef cloudstorage.TableDefinition) []string { - keys := getDDLWatermarkKeys(tableDef) - for _, key := range keys { - if c.tableDDLWatermark[key] < tableDef.TableVersion { - c.tableDDLWatermark[key] = tableDef.TableVersion + if oldTableKey, ok := getRenameTableOldTableKey(tableDef); ok { + if c.tableDDLWatermark[oldTableKey] < tableDef.TableVersion { + c.tableDDLWatermark[oldTableKey] = tableDef.TableVersion } } - return keys + return key } func (c *consumer) handleNewFiles( @@ -631,13 +617,13 @@ func (c *consumer) handleNewFiles( if err := c.sink.WriteBlockEvent(ddlEvent); err != nil { return errors.Trace(err) } - watermarkKeys := c.updateTableDDLWatermark(tableDef) + watermarkKey := c.updateTableDDLWatermark(tableDef) // TODO: need to cleanup tableDefMap in the future. log.Info("execute ddl event successfully", zap.String("query", tableDef.Query), zap.String("schema", key.Schema), zap.String("table", key.Table), zap.Uint64("ddlWatermark", c.tableDDLWatermark[tableKey]), - zap.Strings("watermarkKeys", watermarkKeys)) + zap.String("watermarkKey", watermarkKey)) continue } From 01417ce67b660c9a75f2c5a68012dafab96aeca6 Mon Sep 17 00:00:00 2001 From: nhsmw Date: Mon, 20 Apr 2026 16:10:10 +0800 Subject: [PATCH 3/3] Apply suggestions from code review Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- cmd/storage-consumer/consumer.go | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index 8b4bd0337e..ce2fec450b 100644 --- a/cmd/storage-consumer/consumer.go +++ b/cmd/storage-consumer/consumer.go @@ -524,11 +524,15 @@ func getRenameTableOldTableKey(tableDef cloudstorage.TableDefinition) (string, b log.Panic("parse statement failed", zap.Any("DDL", tableDef.Query), zap.Error(err)) } // The query in job maybe "RENAME TABLE table1 to table2" - oldSchemaName := stmt.(*ast.RenameTableStmt).TableToTables[0].OldTable.Schema.O - if oldSchemaName != "" { - schemaName = oldSchemaName + renameStmt, ok := stmt.(*ast.RenameTableStmt) + if !ok || len(renameStmt.TableToTables) == 0 { + log.Panic("invalid rename table statement", zap.Any("DDL", tableDef.Query)) } - tableName = stmt.(*ast.RenameTableStmt).TableToTables[0].OldTable.Name.O + oldTable := renameStmt.TableToTables[0].OldTable + if oldTable.Schema.O != "" { + schemaName = oldTable.Schema.O + } + tableName = oldTable.Name.O return commonType.QuoteSchema(schemaName, tableName), true }