diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index 05ec67cf4b..48ff2e8dee 100644 --- a/cmd/storage-consumer/consumer.go +++ b/cmd/storage-consumer/consumer.go @@ -36,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" @@ -501,6 +504,42 @@ func (c *consumer) mustGetSchemaFile(key cloudstorage.SchemaPathKey) cloudstorag return *schemaFile } +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" + renameStmt, ok := stmt.(*ast.RenameTableStmt) + if !ok || len(renameStmt.TableToTables) == 0 { + log.Panic("invalid rename table statement", zap.Any("DDL", tableDef.Query)) + } + oldTable := renameStmt.TableToTables[0].OldTable + if oldTable.Schema.O != "" { + schemaName = oldTable.Schema.O + } + tableName = oldTable.Name.O + return commonType.QuoteSchema(schemaName, tableName), true +} + +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 + } + if oldTableKey, ok := getRenameTableOldTableKey(tableDef); ok { + if c.tableDDLWatermark[oldTableKey] < tableDef.TableVersion { + c.tableDDLWatermark[oldTableKey] = tableDef.TableVersion + } + } + return key +} + func (c *consumer) handleNewFiles( ctx context.Context, dmlFileMap map[cloudstorage.DMLPathKey]fileIndexRange, @@ -558,12 +597,18 @@ func (c *consumer) handleNewFiles( if err := c.sink.WriteBlockEvent(ddlEvent); err != nil { return errors.Trace(err) } +<<<<<<< HEAD c.tableDDLWatermark[tableKey] = key.TableVersion // TODO: need to cleanup schemaFileMap in the future. +======= + watermarkKey := c.updateTableDDLWatermark(tableDef) + // TODO: need to cleanup tableDefMap in the future. +>>>>>>> 8486a5dbb (consumer: ignore stale event after renaming table in storage consumer (#4871)) log.Info("execute ddl event successfully", zap.String("query", schemaFile.Query), zap.String("schema", key.Schema), zap.String("table", key.Table), - zap.Uint64("ddlWatermark", c.tableDDLWatermark[tableKey])) + zap.Uint64("ddlWatermark", c.tableDDLWatermark[tableKey]), + zap.String("watermarkKey", watermarkKey)) continue }