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
18 changes: 15 additions & 3 deletions api/v2/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -506,8 +506,16 @@ func (c *ReplicaConfig) toInternalReplicaConfigWithOriginConfig(
}
var debeziumConfig *config.DebeziumConfig
if c.Sink.DebeziumConfig != nil {
// Fall back to the default when OutputOldValue is omitted.
outputOldValue := config.DefaultDebeziumOutputOldValue
if c.Sink.DebeziumConfig.OutputOldValue != nil {
outputOldValue = *c.Sink.DebeziumConfig.OutputOldValue
}
debeziumConfig = &config.DebeziumConfig{
OutputOldValue: c.Sink.DebeziumConfig.OutputOldValue,
OutputOldValue: outputOldValue,
}
if c.Sink.DebeziumConfig.IncludeStartTs != nil {
debeziumConfig.IncludeStartTs = util.AddressOf(*c.Sink.DebeziumConfig.IncludeStartTs)
}
}
var openProtocolConfig *config.OpenProtocolConfig
Expand Down Expand Up @@ -863,7 +871,10 @@ func ToAPIReplicaConfig(c *config.ReplicaConfig) *ReplicaConfig {
var debeziumConfig *DebeziumConfig
if cloned.Sink.Debezium != nil {
debeziumConfig = &DebeziumConfig{
OutputOldValue: cloned.Sink.Debezium.OutputOldValue,
OutputOldValue: util.AddressOf(cloned.Sink.Debezium.OutputOldValue),
}
if cloned.Sink.Debezium.IncludeStartTs != nil {
debeziumConfig.IncludeStartTs = util.AddressOf(*cloned.Sink.Debezium.IncludeStartTs)
}
}
var openProtocolConfig *OpenProtocolConfig
Expand Down Expand Up @@ -1545,7 +1556,8 @@ type OpenProtocolConfig struct {

// DebeziumConfig represents the configurations for debezium protocol encoding
type DebeziumConfig struct {
OutputOldValue bool `json:"output_old_value"`
OutputOldValue *bool `json:"output_old_value,omitempty"`
IncludeStartTs *bool `json:"include_start_ts,omitempty"`
}

type DispatcherCount struct {
Expand Down
21 changes: 21 additions & 0 deletions api/v2/model_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@ func TestReplicaConfigConversion(t *testing.T) {
SpoolDiskQuota: util.AddressOf(int64(1024)),
SpoolBaseDir: util.AddressOf("/tmp/ticdc-spool"),
},
DebeziumConfig: &DebeziumConfig{
IncludeStartTs: util.AddressOf(true),
},
},
Mounter: &MounterConfig{
WorkerNum: util.AddressOf(16),
Expand Down Expand Up @@ -73,6 +76,7 @@ func TestReplicaConfigConversion(t *testing.T) {
require.True(t, util.GetOrZero(internalCfg.Sink.CloudStorageConfig.UseTableIDAsPath))
require.Equal(t, int64(1024), util.GetOrZero(internalCfg.Sink.CloudStorageConfig.SpoolDiskQuota))
require.Equal(t, "/tmp/ticdc-spool", util.GetOrZero(internalCfg.Sink.CloudStorageConfig.SpoolBaseDir))
require.True(t, util.GetOrZero(internalCfg.Sink.Debezium.IncludeStartTs))
require.Equal(t, internalCfg.Mounter.WorkerNum, *apiCfg.Mounter.WorkerNum)
require.True(t, util.GetOrZero(internalCfg.Scheduler.EnableTableAcrossNodes))
require.Equal(t, 1000, util.GetOrZero(internalCfg.Scheduler.RegionThreshold))
Expand All @@ -82,6 +86,21 @@ func TestReplicaConfigConversion(t *testing.T) {
require.Equal(t, int64(128), util.GetOrZero(internalCfg.Consistent.MaxLogSize))
require.Equal(t, int64(2000), util.GetOrZero(internalCfg.Consistent.FlushIntervalInMs))
require.Equal(t, "s3://test", util.GetOrZero(internalCfg.Consistent.Storage))
// output_old_value is omitted in apiCfg and must keep its default (true).
require.True(t, internalCfg.Sink.Debezium.OutputOldValue)

// An explicit output_old_value must be honored.
apiCfgDebezium := &ReplicaConfig{
Sink: &SinkConfig{
DebeziumConfig: &DebeziumConfig{
OutputOldValue: util.AddressOf(false),
IncludeStartTs: util.AddressOf(true),
},
},
}
internalDebezium := apiCfgDebezium.ToInternalReplicaConfig()
require.False(t, internalDebezium.Sink.Debezium.OutputOldValue)
require.True(t, util.GetOrZero(internalDebezium.Sink.Debezium.IncludeStartTs))

// Test case 2: Nil fields (should use defaults or be nil)
apiCfgNil := &ReplicaConfig{}
Expand All @@ -100,6 +119,8 @@ func TestReplicaConfigConversion(t *testing.T) {
require.True(t, *apiCfgBack.Sink.CloudStorageConfig.UseTableIDAsPath)
require.Equal(t, int64(1024), *apiCfgBack.Sink.CloudStorageConfig.SpoolDiskQuota)
require.Equal(t, "/tmp/ticdc-spool", *apiCfgBack.Sink.CloudStorageConfig.SpoolBaseDir)
require.True(t, util.GetOrZero(apiCfgBack.Sink.DebeziumConfig.IncludeStartTs))
require.True(t, util.GetOrZero(apiCfgBack.Sink.DebeziumConfig.OutputOldValue))
require.Equal(t, 16, *apiCfgBack.Mounter.WorkerNum)
require.True(t, *apiCfgBack.Scheduler.EnableTableAcrossNodes)
require.Equal(t, "correctness", *apiCfgBack.Integrity.IntegrityCheckLevel)
Expand Down
2 changes: 1 addition & 1 deletion pkg/config/replica_config.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ var defaultReplicaConfig = &ReplicaConfig{
SendAllBootstrapAtStart: util.AddressOf(DefaultSendAllBootstrapAtStart),
DebeziumDisableSchema: util.AddressOf(false),
OpenProtocol: &OpenProtocolConfig{OutputOldValue: true},
Debezium: &DebeziumConfig{OutputOldValue: true},
Debezium: &DebeziumConfig{OutputOldValue: DefaultDebeziumOutputOldValue},
},
Consistent: &ConsistentConfig{
Level: util.AddressOf("none"),
Expand Down
7 changes: 7 additions & 0 deletions pkg/config/sink.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,10 @@ const (
// to send all tables bootstrap message at changefeed start.
DefaultSendAllBootstrapAtStart = false

// DefaultDebeziumOutputOldValue is the default value of whether
// to output the old value in debezium protocol messages.
DefaultDebeziumOutputOldValue = true

// DefaultMaxReconnectToPulsarBroker is the default max reconnect times to pulsar broker.
// The pulsar client uses an exponential backoff with jitter to reconnect to the broker.
// Based on test, when the max reconnect times is 3,
Expand Down Expand Up @@ -1166,6 +1170,9 @@ type OpenProtocolConfig struct {
// DebeziumConfig represents the configurations for debezium protocol encoding
type DebeziumConfig struct {
OutputOldValue bool `toml:"output-old-value" json:"output-old-value"`
// IncludeStartTs controls whether the transaction start_ts is included in
// the source block of Debezium JSON output.
IncludeStartTs *bool `toml:"include-start-ts" json:"include-start-ts,omitempty"`
}

// validRoutingExpressionRegexp accepts routing expressions made of literal text
Expand Down
26 changes: 25 additions & 1 deletion pkg/sink/codec/common/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,9 @@ type Config struct {
DebeziumDisableSchema bool
// Debezium only. Whether before value should be included in the output.
DebeziumOutputOldValue bool
// Debezium only. Whether the transaction start_ts should be included in
// the source block of the output. JSON protocol only.
DebeziumIncludeStartTs bool
// CSV only. Whether header should be included in the output.
CSVOutputFieldHeader bool
}
Expand Down Expand Up @@ -138,6 +141,7 @@ func NewConfig(protocol config.Protocol) *Config {
DebeziumOutputOldValue: true,
OpenOutputOldValue: true,
DebeziumDisableSchema: false,
DebeziumIncludeStartTs: false,
CSVOutputFieldHeader: false,
}
}
Expand Down Expand Up @@ -177,7 +181,8 @@ type urlConfig struct {
OnlyOutputUpdatedColumns *bool `form:"only-output-updated-columns"`
ContentCompatible *bool `form:"content-compatible"`

DebeziumDisableSchema *bool `form:"debezium-disable-schema"`
DebeziumDisableSchema *bool `form:"debezium-disable-schema"`
DebeziumIncludeStartTs *bool `form:"debezium-include-start-ts"`
// EncodingFormatType is only works for the simple protocol,
// can be `json` and `avro`, default to `json`.
EncodingFormatType *string `form:"encoding-format"`
Expand All @@ -195,6 +200,10 @@ func (c *Config) Apply(sinkURI *url.URL, sinkConfig *config.SinkConfig) error {
if err = binding.Query.Bind(req, urlParameter); err != nil {
return errors.WrapError(errors.ErrSinkInvalidConfig, err)
}
// Keep the raw URI parameters: mergeConfig uses mergo, which cannot
// override a *bool "true" (from the config file) with an explicit
// "false" from the sink URI, so explicit URI values are applied last.
rawURLParameter := urlParameter
if urlParameter, err = mergeConfig(sinkConfig, urlParameter); err != nil {
return err
}
Expand Down Expand Up @@ -300,6 +309,12 @@ func (c *Config) Apply(sinkURI *url.URL, sinkConfig *config.SinkConfig) error {
if urlParameter.DebeziumDisableSchema != nil {
c.DebeziumDisableSchema = *urlParameter.DebeziumDisableSchema
}
if urlParameter.DebeziumIncludeStartTs != nil {
c.DebeziumIncludeStartTs = *urlParameter.DebeziumIncludeStartTs
}
if rawURLParameter.DebeziumIncludeStartTs != nil {
c.DebeziumIncludeStartTs = *rawURLParameter.DebeziumIncludeStartTs
}

return nil
}
Expand Down Expand Up @@ -331,6 +346,9 @@ func mergeConfig(
if sinkConfig.DebeziumDisableSchema != nil {
dest.DebeziumDisableSchema = sinkConfig.DebeziumDisableSchema
}
if sinkConfig.Debezium != nil && sinkConfig.Debezium.IncludeStartTs != nil {
dest.DebeziumIncludeStartTs = sinkConfig.Debezium.IncludeStartTs
}
}
if err := mergo.Merge(dest, urlParameters, mergo.WithOverride); err != nil {
return nil, err
Expand Down Expand Up @@ -360,6 +378,12 @@ func (c *Config) Validate() error {
zap.String("protocol", c.Protocol.String()))
}

if c.DebeziumIncludeStartTs && c.Protocol != config.ProtocolDebezium {
return errors.ErrCodecInvalidConfig.GenWithStack(
`debezium-include-start-ts only takes effect with protocol "debezium"`,
)
}

if c.Protocol == config.ProtocolAvro {
if c.AvroConfluentSchemaRegistry != "" && c.AvroGlueSchemaRegistry != nil {
return errors.ErrCodecInvalidConfig.GenWithStack(
Expand Down
36 changes: 36 additions & 0 deletions pkg/sink/codec/common/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,3 +68,39 @@ func TestValidateMessageLimits(t *testing.T) {
})
}
}

func TestDebeziumIncludeStartTsConfig(t *testing.T) {
// URI parameter
cfg := NewConfig(config.ProtocolDebezium)
sinkURI, err := url.Parse("kafka://127.0.0.1:9092/topic?protocol=debezium&debezium-include-start-ts=true")
require.NoError(t, err)
require.NoError(t, cfg.Apply(sinkURI, config.GetDefaultReplicaConfig().Sink))
require.True(t, cfg.DebeziumIncludeStartTs)
require.NoError(t, cfg.Validate())

// changefeed config file
on := true
cfg2 := NewConfig(config.ProtocolDebezium)
sinkConfig := config.GetDefaultReplicaConfig().Sink
sinkConfig.Debezium.IncludeStartTs = &on
sinkURI2, err := url.Parse("kafka://127.0.0.1:9092/topic?protocol=debezium")
require.NoError(t, err)
require.NoError(t, cfg2.Apply(sinkURI2, sinkConfig))
require.True(t, cfg2.DebeziumIncludeStartTs)

// URI parameter overrides the config file
cfg3 := NewConfig(config.ProtocolDebezium)
sinkConfig3 := config.GetDefaultReplicaConfig().Sink
sinkConfig3.Debezium.IncludeStartTs = &on
sinkURI3, err := url.Parse("kafka://127.0.0.1:9092/topic?protocol=debezium&debezium-include-start-ts=false")
require.NoError(t, err)
require.NoError(t, cfg3.Apply(sinkURI3, sinkConfig3))
require.False(t, cfg3.DebeziumIncludeStartTs)

// only supported by the debezium (JSON) protocol
cfg4 := NewConfig(config.ProtocolCanalJSON)
cfg4.DebeziumIncludeStartTs = true
errCode, ok := errors.RFCCode(cfg4.Validate())
require.True(t, ok)
require.Equal(t, errors.ErrCodecInvalidConfig.RFCCode(), errCode)
}
23 changes: 19 additions & 4 deletions pkg/sink/codec/debezium/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -758,7 +758,10 @@ func (c *dbzCodec) writeBinaryField(writer *util.JSONWriter, fieldName string, v
writer.WriteBase64StringField(fieldName, value)
}

func (c *dbzCodec) writeSourceSchema(writer *util.JSONWriter) {
// includeStartTs indicates whether start_ts should be declared in the source
// schema. DML callers pass the configured value, while DDL and checkpoint
// callers pass false because their payloads do not carry the field.
func (c *dbzCodec) writeSourceSchema(writer *util.JSONWriter, includeStartTs bool) {
writer.WriteObjectElement(func() {
writer.WriteStringField("type", "struct")
writer.WriteArrayField("fields", func() {
Expand Down Expand Up @@ -843,6 +846,13 @@ func (c *dbzCodec) writeSourceSchema(writer *util.JSONWriter) {
writer.WriteBoolField("optional", true)
writer.WriteStringField("field", "query")
})
if includeStartTs {
writer.WriteObjectElement(func() {
writer.WriteStringField("type", "int64")
writer.WriteBoolField("optional", false)
writer.WriteStringField("field", "start_ts")
})
}
})
writer.WriteBoolField("optional", false)
writer.WriteStringField("name", "io.debezium.connector.mysql.Source")
Expand Down Expand Up @@ -933,6 +943,11 @@ func (c *dbzCodec) EncodeValue(

// The followings are TiDB extended fields
jWriter.WriteUint64Field("commit_ts", e.CommitTs)
// start_ts: the start TSO of the transaction that made this change,
// exposed for downstream consumers that need transaction correlation.
if c.config.DebeziumIncludeStartTs {
jWriter.WriteUint64Field("start_ts", e.StartTs)
}
jWriter.WriteStringField("cluster_id", c.clusterID)
})

Expand Down Expand Up @@ -1026,7 +1041,7 @@ func (c *dbzCodec) EncodeValue(
jWriter.WriteRaw(fieldsJSON)
})
})
c.writeSourceSchema(jWriter)
c.writeSourceSchema(jWriter, c.config.DebeziumIncludeStartTs)
jWriter.WriteObjectElement(func() {
jWriter.WriteStringField("type", "string")
jWriter.WriteBoolField("optional", false)
Expand Down Expand Up @@ -1312,7 +1327,7 @@ func (c *dbzCodec) EncodeDDLEvent(
jWriter.WriteIntField("version", 1)
jWriter.WriteStringField("name", "io.debezium.connector.mysql.SchemaChangeValue")
jWriter.WriteArrayField("fields", func() {
c.writeSourceSchema(jWriter)
c.writeSourceSchema(jWriter, false)
jWriter.WriteObjectElement(func() {
jWriter.WriteStringField("field", "ts_ms")
jWriter.WriteBoolField("optional", false)
Expand Down Expand Up @@ -1551,7 +1566,7 @@ func (c *dbzCodec) EncodeCheckpointEvent(
fmt.Sprintf("%s.%s.Envelope", common.SanitizeName(c.clusterID), "watermark"))
jWriter.WriteIntField("version", 1)
jWriter.WriteArrayField("fields", func() {
c.writeSourceSchema(jWriter)
c.writeSourceSchema(jWriter, false)
jWriter.WriteObjectElement(func() {
jWriter.WriteStringField("type", "string")
jWriter.WriteBoolField("optional", false)
Expand Down
41 changes: 41 additions & 0 deletions pkg/sink/codec/debezium/codec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1544,3 +1544,44 @@ func BenchmarkEncodeLargeBinary(b *testing.B) {
codec.EncodeValue(e, buf)
}
}

func TestStartTsNotInDDLAndCheckpointEvents(t *testing.T) {
// Even with debezium-include-start-ts enabled, DDL and checkpoint
// (watermark) messages must not declare start_ts in their schemas:
// their payloads never carry the field (no per-row transaction), and a
// declared-but-absent non-optional field breaks schema-validating consumers.
codec := &dbzCodec{
config: common.NewConfig(config.ProtocolDebezium),
clusterID: "test_cluster",
nowFunc: func() time.Time { return time.Unix(1701326309, 0) },
}
codec.config.DebeziumIncludeStartTs = true
codec.config.DebeziumDisableSchema = false

helper := commonEvent.NewEventTestHelper(t)
defer helper.Close()
helper.Tk().MustExec("use test")
helper.DDL2Job(`create table test.table1(id int(10) primary key)`)
job := helper.DDL2Job(`RENAME TABLE test.table1 to test.table2`)
tableInfo := helper.GetTableInfo(job)

e := &commonEvent.DDLEvent{
FinishedTs: 1,
TableInfo: tableInfo,
SchemaName: "test",
TableName: "table2",
ExtraSchemaName: "test",
ExtraTableName: "table1",
Type: byte(timodel.ActionRenameTable),
Query: job.Query,
}
keyBuf := bytes.NewBuffer(nil)
buf := bytes.NewBuffer(nil)
require.NoError(t, codec.EncodeDDLEvent(e, keyBuf, buf))
require.NotContains(t, buf.String(), "start_ts")

keyBuf.Reset()
buf.Reset()
require.NoError(t, codec.EncodeCheckpointEvent(3, keyBuf, buf))
require.NotContains(t, buf.String(), "start_ts")
}
Loading
Loading