Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 11 additions & 11 deletions api/v2/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -1191,17 +1191,17 @@ type Table struct {
// SinkConfig represents sink config for a changefeed
// This is a duplicate of config.SinkConfig
type SinkConfig struct {
Protocol *string `json:"protocol,omitempty" toml:"protocol,omitempty"`
SchemaRegistry *string `json:"schema_registry,omitempty" toml:"schema-registry,omitempty"`
CSVConfig *CSVConfig `json:"csv,omitempty" toml:"csv,omitempty"`
DispatchRules []*DispatchRule `json:"dispatchers,omitempty" toml:"dispatchers,omitempty"`
ColumnSelectors []*ColumnSelector `json:"column_selectors,omitempty" toml:"column-selectors,omitempty"`
TxnAtomicity *string `json:"transaction_atomicity,omitempty" toml:"transaction-atomicity,omitempty"`
EncoderConcurrency *int `json:"encoder_concurrency,omitempty" toml:"encoder-concurrency,omitempty"`
Terminator *string `json:"terminator,omitempty" toml:"terminator,omitempty"`
DateSeparator *string `json:"date_separator,omitempty" toml:"date-separator,omitempty"`
EnablePartitionSeparator *bool `json:"enable_partition_separator,omitempty" toml:"enable-partition-separator,omitempty"`
FileIndexWidth *int `json:"file_index_width,omitempty" toml:"file-index-digit,omitempty"`
Protocol *string `json:"protocol,omitempty" toml:"protocol,omitempty"`
SchemaRegistry *string `json:"schema_registry,omitempty" toml:"schema-registry,omitempty"`
CSVConfig *CSVConfig `json:"csv,omitempty" toml:"csv,omitempty"`
DispatchRules []*DispatchRule `json:"dispatchers,omitempty" toml:"dispatchers,omitempty"`
ColumnSelectors []*ColumnSelector `json:"column_selectors,omitempty" toml:"column-selectors,omitempty"`
TxnAtomicity *string `json:"transaction_atomicity,omitempty" toml:"transaction-atomicity,omitempty"`
EncoderConcurrency *int `json:"encoder_concurrency,omitempty" toml:"encoder-concurrency,omitempty"`
Terminator *string `json:"terminator,omitempty" toml:"terminator,omitempty"`
DateSeparator *config.DateSeparator `json:"date_separator,omitempty" toml:"date-separator,omitempty"`
EnablePartitionSeparator *bool `json:"enable_partition_separator,omitempty" toml:"enable-partition-separator,omitempty"`
FileIndexWidth *int `json:"file_index_width,omitempty" toml:"file-index-digit,omitempty"`
// deprecated: it's become useless since v9.0.0
EnableKafkaSinkV2 *bool `json:"enable_kafka_sink_v2,omitempty" toml:"enable-kafka-sink-v2,omitempty"`
OnlyOutputUpdatedColumns *bool `json:"only_output_updated_columns,omitempty" toml:"only-output-updated-columns,omitempty"`
Expand Down
13 changes: 13 additions & 0 deletions api/v2/model_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,26 @@
package v2

import (
"encoding/json"
"testing"

"github.com/pingcap/ticdc/pkg/config"
"github.com/pingcap/ticdc/pkg/util"
"github.com/stretchr/testify/require"
)

func TestSinkConfigDateSeparator(t *testing.T) {
t.Parallel()

var sinkConfig SinkConfig
require.NoError(t, json.Unmarshal([]byte(`{"date_separator":"DAY"}`), &sinkConfig))
require.Equal(t, config.DateSeparatorDay, util.GetOrZero(sinkConfig.DateSeparator))

err := json.Unmarshal([]byte(`{"date_separator":"week"}`), &SinkConfig{})
require.Error(t, err)
require.ErrorContains(t, err, "CDC:ErrStorageSinkInvalidConfig")
}

// TestReplicaConfigConversion verifies API/internal replica config conversion,
// including round-tripping the optional event collector batch overrides.
func TestReplicaConfigConversion(t *testing.T) {
Expand Down
6 changes: 4 additions & 2 deletions cmd/storage-consumer/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ type storageMetadata struct {

type consumer struct {
replicationCfg *config.ReplicaConfig
dateSeparator config.DateSeparator
codecCfg *common.Config
columnSelectors *columnselector.ColumnSelectors
externalStorage storeapi.Storage
Expand Down Expand Up @@ -113,6 +114,7 @@ func newConsumer(ctx context.Context) (*consumer, error) {
log.Error("failed to validate replica config", zap.Error(err))
return nil, err
}
dateSeparator := putil.GetOrZero(replicaConfig.Sink.DateSeparator)

switch putil.GetOrZero(replicaConfig.Sink.Protocol) {
case config.ProtocolCsv.String():
Expand Down Expand Up @@ -167,6 +169,7 @@ func newConsumer(ctx context.Context) (*consumer, error) {

return &consumer{
replicationCfg: replicaConfig,
dateSeparator: dateSeparator,
codecCfg: codecConfig,
columnSelectors: columnSelectors,
externalStorage: storage,
Expand Down Expand Up @@ -256,15 +259,14 @@ func (c *consumer) getNewFiles(
origDMLIdxMap[k] = m
}

dateSeparator := putil.GetOrZero(c.replicationCfg.Sink.DateSeparator)
err := c.externalStorage.WalkDir(ctx, opt, func(path string, _ int64) error {
if cloudstorage.IsSchemaFile(path) {
c.parseSchemaFilePath(ctx, path)
return nil
}
if strings.HasSuffix(path, ".index") {
var dmlkey cloudstorage.DMLPathKey
if err := dmlkey.ParseIndexFilePath(dateSeparator, path); err != nil {
if err := dmlkey.ParseIndexFilePath(c.dateSeparator, path); err != nil {
log.Debug("ignore handling unsupported dml index file", zap.String("path", path))
return nil
}
Expand Down
4 changes: 2 additions & 2 deletions downstreamadapter/sink/cloudstorage/dml_writers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ func verifyCloudStorageWriteEventsWithoutDateSeparator(
err = replicaConfig.ValidateAndAdjust(sinkURI)
require.NoError(t, err)

replicaConfig.Sink.DateSeparator = putil.AddressOf(config.DateSeparatorNone.String())
replicaConfig.Sink.DateSeparator = putil.AddressOf(config.DateSeparatorNone)
replicaConfig.Sink.FileIndexWidth = putil.AddressOf(6)

ctx, cancel := context.WithCancel(context.Background())
Expand Down Expand Up @@ -309,7 +309,7 @@ func verifyCloudStorageWriteEventsWithDateSeparator(
err = replicaConfig.ValidateAndAdjust(sinkURI)
require.NoError(t, err)

replicaConfig.Sink.DateSeparator = putil.AddressOf(config.DateSeparatorDay.String())
replicaConfig.Sink.DateSeparator = putil.AddressOf(config.DateSeparatorDay)
replicaConfig.Sink.FileIndexWidth = putil.AddressOf(6)

mockClock := pclock.NewMock()
Expand Down
6 changes: 3 additions & 3 deletions downstreamadapter/sink/cloudstorage/sink.go
Original file line number Diff line number Diff line change
Expand Up @@ -403,11 +403,11 @@ func (s *sink) initCron(
}

func (s *sink) bgCleanup(ctx context.Context) {
if s.cfg.DateSeparator != config.DateSeparatorDay.String() || s.cfg.FileExpirationDays <= 0 {
if s.cfg.DateSeparator != config.DateSeparatorDay || s.cfg.FileExpirationDays <= 0 {
log.Info("skip cleanup expired files for storage sink",
zap.String("keyspace", s.changefeedID.Keyspace()),
zap.String("changefeedID", s.changefeedID.Name()),
zap.String("dateSeparator", s.cfg.DateSeparator),
zap.Stringer("dateSeparator", s.cfg.DateSeparator),
zap.Int("expiredFileTTL", s.cfg.FileExpirationDays))
return
}
Expand All @@ -417,7 +417,7 @@ func (s *sink) bgCleanup(ctx context.Context) {
log.Info("start schedule cleanup expired files for storage sink",
zap.String("keyspace", s.changefeedID.Keyspace()),
zap.String("changefeedID", s.changefeedID.Name()),
zap.String("dateSeparator", s.cfg.DateSeparator),
zap.Stringer("dateSeparator", s.cfg.DateSeparator),
zap.Int("expiredFileTTL", s.cfg.FileExpirationDays))

// wait for the context done
Expand Down
4 changes: 2 additions & 2 deletions downstreamadapter/sink/cloudstorage/sink_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,7 @@ func TestCloudStorageSinkWithColumnSelector(t *testing.T) {
}
err = replicaConfig.ValidateAndAdjust(sinkURI)
require.NoError(t, err)
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone.String())
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
Expand Down Expand Up @@ -785,7 +785,7 @@ func TestCleanupExpiredFiles(t *testing.T) {
cloudStorageSink := &sink{
changefeedID: common.NewChangefeedID4Test("test", "test"),
cfg: &cloudstorage.Config{
DateSeparator: config.DateSeparatorDay.String(),
DateSeparator: config.DateSeparatorDay,
FileExpirationDays: 1,
FileCleanupCronSpec: util.GetOrZero(replicaConfig.Sink.CloudStorageConfig.FileCleanupCronSpec),
},
Expand Down
8 changes: 4 additions & 4 deletions downstreamadapter/sink/cloudstorage/writer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ func testWriter(ctx context.Context, t *testing.T, dir string) *writer {
require.NoError(t, err)
cfg := cloudstorage.NewConfig()
replicaConfig := config.GetDefaultReplicaConfig()
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone.String())
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone)
err = cfg.Apply(context.TODO(), sinkURI, replicaConfig.Sink, true)
cfg.FileIndexWidth = 6
require.NoError(t, err)
Expand Down Expand Up @@ -467,7 +467,7 @@ func TestWriterStoresPendingMessagesInSpoolBeforeFlush(t *testing.T) {

cfg := cloudstorage.NewConfig()
replicaConfig := config.GetDefaultReplicaConfig()
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone.String())
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone)
replicaConfig.Sink.CloudStorageConfig = &config.CloudStorageConfig{
// Keep the quota larger than this encoded batch so the controller still
// spills it to local spool files instead of taking the oversized in-memory fast path.
Expand Down Expand Up @@ -641,7 +641,7 @@ func TestWriterIndexWriteError(t *testing.T) {
require.NoError(t, err)
cfg := cloudstorage.NewConfig()
replicaConfig := config.GetDefaultReplicaConfig()
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone.String())
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone)
err = cfg.Apply(context.TODO(), sinkURI, replicaConfig.Sink, true)
require.NoError(t, err)
cfg.FileIndexWidth = 6
Expand Down Expand Up @@ -706,7 +706,7 @@ func TestWriterDataFileCloseError(t *testing.T) {
require.NoError(t, err)
cfg := cloudstorage.NewConfig()
replicaConfig := config.GetDefaultReplicaConfig()
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone.String())
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorNone)
err = cfg.Apply(context.TODO(), sinkURI, replicaConfig.Sink, true)
require.NoError(t, err)
cfg.FileIndexWidth = 6
Expand Down
2 changes: 1 addition & 1 deletion pkg/cloudstorage/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ type Config struct {
FlushInterval time.Duration
FileSize int
FileIndexWidth int
DateSeparator string
DateSeparator config.DateSeparator
FileExpirationDays int
FileCleanupCronSpec string
EnablePartitionSeparator bool
Expand Down
4 changes: 3 additions & 1 deletion pkg/cloudstorage/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/pingcap/ticdc/pkg/config"
"github.com/pingcap/ticdc/pkg/errors"
"github.com/pingcap/ticdc/pkg/util"
"github.com/stretchr/testify/require"
)

Expand All @@ -30,7 +31,7 @@ func TestConfigApply(t *testing.T) {
expected.FlushInterval = 10 * time.Second
expected.FileSize = 16 * 1024 * 1024
expected.FileIndexWidth = config.DefaultFileIndexWidth
expected.DateSeparator = config.DateSeparatorDay.String()
expected.DateSeparator = config.DateSeparatorDay
expected.EnablePartitionSeparator = true
expected.FlushConcurrency = 1
expected.SpoolDiskQuota = 10 * 1024 * 1024 * 1024
Expand All @@ -40,6 +41,7 @@ func TestConfigApply(t *testing.T) {
require.NoError(t, err)

replicaConfig := config.GetDefaultReplicaConfig()
replicaConfig.Sink.DateSeparator = util.AddressOf(config.DateSeparatorDay)
err = replicaConfig.ValidateAndAdjust(sinkURI)
require.NoError(t, err)
cfg := NewConfig()
Expand Down
8 changes: 4 additions & 4 deletions pkg/cloudstorage/generator.go
Original file line number Diff line number Diff line change
Expand Up @@ -323,11 +323,11 @@ func (f *FilePathGenerator) GenerateDateStr() string {
currTime := f.pdClock.CurrentTime()
// Note: `dateStr` is formatted using local TZ.
switch f.config.DateSeparator {
case config.DateSeparatorYear.String():
case config.DateSeparatorYear:
dateStr = currTime.Format("2006")
case config.DateSeparatorMonth.String():
case config.DateSeparatorMonth:
dateStr = currTime.Format("2006-01")
case config.DateSeparatorDay.String():
case config.DateSeparatorDay:
dateStr = currTime.Format("2006-01-02")
default:
}
Expand Down Expand Up @@ -519,7 +519,7 @@ func RemoveExpiredFiles(
cfg *Config,
checkpointTs uint64,
) error {
if cfg.DateSeparator != config.DateSeparatorDay.String() {
if cfg.DateSeparator != config.DateSeparatorDay {
return nil
}
if dateSeparatorDayRegexp == nil {
Expand Down
16 changes: 9 additions & 7 deletions pkg/cloudstorage/path_key.go
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,7 @@ func (d DMLPathKey) generateDMLDataDirPath() string {
}

func (d *DMLPathKey) parseDMLDataDir(
dateSeparator string, parts []string, filePath string,
dateSeparator config.DateSeparator, parts []string, filePath string,
) error {
var (
key DMLPathKey
Expand All @@ -228,14 +228,14 @@ func (d *DMLPathKey) parseDMLDataDir(
dateRE string
)
switch dateSeparator {
case config.DateSeparatorNone.String():
case config.DateSeparatorYear.String():
case config.DateSeparatorNone:
case config.DateSeparatorYear:
hasDate = true
dateRE = config.DateSeparatorYear.GetPattern()
case config.DateSeparatorMonth.String():
case config.DateSeparatorMonth:
hasDate = true
dateRE = config.DateSeparatorMonth.GetPattern()
case config.DateSeparatorDay.String():
case config.DateSeparatorDay:
hasDate = true
dateRE = config.DateSeparatorDay.GetPattern()
default:
Expand Down Expand Up @@ -304,7 +304,7 @@ func (d *DMLPathKey) parseDMLDataDir(
// <data-dir>/meta/CDC_<dispatcherID>.index. Only the data directory portion is
// parsed here; the file index itself is stored in the index file content and is
// read by the caller.
func (d *DMLPathKey) ParseIndexFilePath(dateSeparator, path string) error {
func (d *DMLPathKey) ParseIndexFilePath(dateSeparator config.DateSeparator, path string) error {
parts := strings.Split(path, "/")
if len(parts) < 4 || parts[len(parts)-2] != "meta" || !isDMLIndexFileName(parts[len(parts)-1]) {
return invalidDMLPathError(path)
Expand All @@ -316,7 +316,9 @@ func (d *DMLPathKey) ParseIndexFilePath(dateSeparator, path string) error {
// index encoded in the file name.
// Input is <data-dir>/CDC<idx><extension> or
// <data-dir>/CDC_<dispatcherID>_<idx><extension>. Invalid paths panic.
func (d *DMLPathKey) ParseDMLFilePath(dateSeparator, filePath, extension string) FileIndex {
func (d *DMLPathKey) ParseDMLFilePath(
dateSeparator config.DateSeparator, filePath, extension string,
) FileIndex {
parts := strings.Split(filePath, "/")
fileIndex, err := ParseFileIndexFromFileName(parts[len(parts)-1], extension)
if err != nil {
Expand Down
16 changes: 8 additions & 8 deletions pkg/cloudstorage/path_key_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,15 +63,15 @@ func TestGenerateDMLFilePath(t *testing.T) {
index uint64
fileIndexWidth int
extension string
dateSeparator string
dateSeparator config.DateSeparator
path string
dmlkey DMLPathKey
}{
{
index: 10,
fileIndexWidth: 20,
extension: ".csv",
dateSeparator: config.DateSeparatorDay.String(),
dateSeparator: config.DateSeparatorDay,
path: fmt.Sprintf("schema1/table1/123456/2023-05-09/CDC_%s_00000000000000000010.csv", dispatcherID.String()),
dmlkey: DMLPathKey{
SchemaPathKey: SchemaPathKey{
Expand All @@ -87,7 +87,7 @@ func TestGenerateDMLFilePath(t *testing.T) {
index: 10,
fileIndexWidth: 20,
extension: ".csv",
dateSeparator: config.DateSeparatorNone.String(),
dateSeparator: config.DateSeparatorNone,
path: fmt.Sprintf("12345/123456/CDC_%s_00000000000000000010.csv", dispatcherID.String()),
dmlkey: DMLPathKey{
SchemaPathKey: SchemaPathKey{
Expand All @@ -102,7 +102,7 @@ func TestGenerateDMLFilePath(t *testing.T) {
index: 10,
fileIndexWidth: 20,
extension: ".csv",
dateSeparator: config.DateSeparatorDay.String(),
dateSeparator: config.DateSeparatorDay,
path: fmt.Sprintf("schema1/table1/123456/55/2023-05-09/CDC_%s_00000000000000000010.csv", dispatcherID.String()),
dmlkey: DMLPathKey{
SchemaPathKey: SchemaPathKey{
Expand Down Expand Up @@ -168,22 +168,22 @@ func TestParseIndexFilePathRejectsUnsupportedPath(t *testing.T) {

testCases := []struct {
name string
dateSeparator string
dateSeparator config.DateSeparator
path string
}{
{
name: "legacy schema table date index",
dateSeparator: config.DateSeparatorDay.String(),
dateSeparator: config.DateSeparatorDay,
path: "test/binary_columns_dummy/2026-06-23/meta/CDC.index",
},
{
name: "invalid index file name",
dateSeparator: config.DateSeparatorNone.String(),
dateSeparator: config.DateSeparatorNone,
path: "schema1/table1/123456/meta/notCDC.index",
},
{
name: "date does not match separator",
dateSeparator: config.DateSeparatorMonth.String(),
dateSeparator: config.DateSeparatorMonth,
path: "schema1/table1/123456/2023-05-09/meta/CDC.index",
},
}
Expand Down
Loading
Loading