diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index 05ec67cf4b..399a01e28d 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" @@ -44,6 +47,7 @@ import ( const ( defaultChangefeedName = "storage-consumer" defaultLogInterval = 5 * time.Second + metadataFileName = "metadata" ) type ( @@ -57,6 +61,10 @@ type indexRange struct { end uint64 } +type storageMetadata struct { + CheckpointTs uint64 `json:"checkpoint-ts"` +} + type consumer struct { replicationCfg *config.ReplicaConfig codecCfg *common.Config @@ -76,6 +84,8 @@ type consumer struct { dmlCount atomic.Int64 readSeq atomic.Uint64 + + globalCheckpointTs uint64 } func newConsumer(ctx context.Context) (*consumer, error) { @@ -195,7 +205,30 @@ func diffDMLMaps( return resMap } -// getNewFiles returns newly created dml files in specific ranges +func (c *consumer) getGlobalCheckpointTs(ctx context.Context) error { + exists, err := c.externalStorage.FileExists(ctx, metadataFileName) + if err != nil { + return errors.Trace(err) + } + if !exists { + return nil + } + + data, err := c.externalStorage.ReadFile(ctx, metadataFileName) + if err != nil { + return errors.Trace(err) + } + var metadata storageMetadata + if err := json.Unmarshal(data, &metadata); err != nil { + return errors.Trace(err) + } + if metadata.CheckpointTs > c.globalCheckpointTs { + c.globalCheckpointTs = metadata.CheckpointTs + } + return nil +} + +// getNewFiles returns newly created dml files in specific ranges that are visible under checkpointTs. func (c *consumer) getNewFiles( ctx context.Context, ) (map[cloudstorage.DMLPathKey]fileIndexRange, error) { @@ -397,6 +430,13 @@ func (c *consumer) flushDMLEvents(ctx context.Context, tableID int64) error { } func (c *consumer) parseDMLIndexFile(ctx context.Context, path string, dmlkey cloudstorage.DMLPathKey) { + if c.globalCheckpointTs > 0 && dmlkey.TableVersion > c.globalCheckpointTs { + log.Debug("skip dml index file by checkpoint", + zap.String("path", path), + zap.Uint64("tableVersion", dmlkey.TableVersion), + zap.Uint64("checkpointTs", c.globalCheckpointTs)) + return + } data, err := c.externalStorage.ReadFile(ctx, path) if err != nil { log.Panic("read dml index file failed", @@ -424,6 +464,13 @@ func (c *consumer) parseDMLIndexFile(ctx context.Context, path string, dmlkey cl func (c *consumer) parseSchemaFilePath(ctx context.Context, path string) { var schemaKey cloudstorage.SchemaPathKey schemaKey.Parse(path) + if c.globalCheckpointTs > 0 && schemaKey.TableVersion > c.globalCheckpointTs { + log.Debug("skip schema file by checkpoint", + zap.String("path", path), + zap.Uint64("tableVersion", schemaKey.TableVersion), + zap.Uint64("checkpointTs", c.globalCheckpointTs)) + return + } key := schemaKey.GetKey() if schemaFiles, ok := c.schemaFileMap[key]; ok { if _, ok := schemaFiles[schemaKey.TableVersion]; ok { @@ -501,6 +548,41 @@ func (c *consumer) mustGetSchemaFile(key cloudstorage.SchemaPathKey) cloudstorag return *schemaFile } +func getRenameTableOldTableKey(schemaFile cloudstorage.SchemaFile) (string, bool) { + if schemaFile.Type != byte(timodel.ActionRenameTable) { + return "", false + } + schemaName := schemaFile.Schema + stmt, err := parser.New().ParseOneStmt(schemaFile.Query, "", "") + if err != nil { + log.Panic("parse statement failed", zap.Any("DDL", schemaFile.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", schemaFile.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(schemaFile cloudstorage.SchemaFile) string { + key := commonType.QuoteSchema(schemaFile.Schema, schemaFile.Table) + if c.tableDDLWatermark[key] < schemaFile.TableVersion { + c.tableDDLWatermark[key] = schemaFile.TableVersion + } + if oldTableKey, ok := getRenameTableOldTableKey(schemaFile); ok { + if c.tableDDLWatermark[oldTableKey] < schemaFile.TableVersion { + c.tableDDLWatermark[oldTableKey] = schemaFile.TableVersion + } + } + return key +} + func (c *consumer) handleNewFiles( ctx context.Context, dmlFileMap map[cloudstorage.DMLPathKey]fileIndexRange, @@ -558,12 +640,13 @@ func (c *consumer) handleNewFiles( if err := c.sink.WriteBlockEvent(ddlEvent); err != nil { return errors.Trace(err) } - c.tableDDLWatermark[tableKey] = key.TableVersion + watermarkKey := c.updateTableDDLWatermark(schemaFile) // TODO: need to cleanup schemaFileMap in the future. 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 } @@ -606,7 +689,9 @@ func (c *consumer) handleNewFiles( } } } - c.flushDMLEvents(ctx, tableID) + if err := c.flushDMLEvents(ctx, tableID); err != nil { + return err + } } return nil @@ -641,12 +726,17 @@ func (c *consumer) handle(ctx context.Context) error { } round++ + err := c.getGlobalCheckpointTs(ctx) + if err != nil { + return errors.Trace(err) + } dmlFileMap, err := c.getNewFiles(ctx) if err != nil { return errors.Trace(err) } log.Info("storage consumer scan done", zap.Uint64("round", round), + zap.Uint64("checkpointTs", c.globalCheckpointTs), zap.Int("dmlPathKeyCount", len(dmlFileMap))) err = c.handleNewFiles(ctx, dmlFileMap, round)