From 347ca6ea4e636ec2fa2bc8d7add285a28bdecd09 Mon Sep 17 00:00:00 2001 From: nhsmw Date: Fri, 31 Jul 2026 18:14:37 +0800 Subject: [PATCH 1/2] This is an automated cherry-pick of #5824 Signed-off-by: ti-chi-bot --- cmd/kafka-consumer/writer.go | 45 +++++ cmd/kafka-consumer/writer_test.go | 133 +++++++++++++++ cmd/pulsar-consumer/writer.go | 39 +++++ cmd/pulsar-consumer/writer_test.go | 264 +++++++++++++++++++++++++++++ cmd/storage-consumer/consumer.go | 10 ++ cmd/util/event_group.go | 80 +++++++++ cmd/util/event_group_test.go | 142 +++++++++++++++- 7 files changed, 705 insertions(+), 8 deletions(-) diff --git a/cmd/kafka-consumer/writer.go b/cmd/kafka-consumer/writer.go index 865c61249f..2cf2c68ddb 100644 --- a/cmd/kafka-consumer/writer.go +++ b/cmd/kafka-consumer/writer.go @@ -612,11 +612,26 @@ func (w *writer) appendRow2Group(dml *event.DMLEvent, progress *partitionProgres table = dml.TableInfo.GetTableName() commitTs = dml.GetCommitTs() ) + globalWatermark := w.globalWatermark() + if commitTs < globalWatermark { + log.Warn("DML event fallback row, since less than the global watermark, ignore it", + zap.Int64("tableID", tableID), zap.Int32("partition", progress.partition), + zap.Uint64("commitTs", commitTs), zap.Any("offset", offset), + zap.Uint64("globalWatermark", globalWatermark), + zap.Uint64("partitionWatermark", progress.watermark), + zap.Any("watermarkOffset", progress.watermarkOffset), + zap.String("schema", schema), zap.String("table", table), + zap.Stringer("eventType", message.RowType), + zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) + return + } + group := progress.eventsGroup[tableID] if group == nil { group = util.NewEventsGroup(progress.partition, tableID) progress.eventsGroup[tableID] = group } +<<<<<<< HEAD // IMPORTANT: Kafka offsets are append-only, but CommitTs can go backwards after // a TiCDC restart/retry (at-least-once replay). We must not drop such events // solely based on a "seen" watermark (e.g. HighWatermark). The only safe @@ -650,6 +665,36 @@ func (w *writer) appendRow2Group(dml *event.DMLEvent, progress *partitionProgres zap.Uint64("appliedWatermark", group.AppliedWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", dml.RowTypes[0])) +======= + message = w.messageWithPartitionCheck(message, progress.partition, offset) + group.AppendMessage(message) + if commitTs < progress.watermark { + log.Warn("DML event fallback row, since less than the partition watermark, append it and sort before flush", + zap.Int64("tableID", tableID), zap.Int32("partition", group.Partition), + zap.Uint64("commitTs", commitTs), zap.Any("offset", offset), + zap.Uint64("watermark", progress.watermark), zap.Any("watermarkOffset", progress.watermarkOffset), + zap.Uint64("globalWatermark", globalWatermark), + zap.String("schema", schema), zap.String("table", table), + zap.Stringer("eventType", message.RowType), + zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) + return + } + if commitTs >= group.HighWatermark { + log.Debug("DML event append to the group", + zap.Int32("partition", group.Partition), zap.Any("offset", offset), + zap.Uint64("commitTs", commitTs), zap.Uint64("HighWatermark", group.HighWatermark), + zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), + zap.Stringer("eventType", message.RowType)) + return + } + log.Warn("DML event commit ts fallback, append it and sort before flush", + zap.Int32("partition", progress.partition), zap.Any("offset", offset), + zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), + zap.Any("partitionWatermark", progress.watermark), zap.Any("watermarkOffset", progress.watermarkOffset), + zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), + zap.Stringer("eventType", message.RowType), + zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } func openDB(ctx context.Context, dsn string) (*sql.DB, error) { diff --git a/cmd/kafka-consumer/writer_test.go b/cmd/kafka-consumer/writer_test.go index 70fed688de..8476189cac 100644 --- a/cmd/kafka-consumer/writer_test.go +++ b/cmd/kafka-consumer/writer_test.go @@ -284,6 +284,7 @@ func TestWriterWrite_handlesOutOfOrderDDLsByCommitTs(t *testing.T) { require.Equal(t, "CREATE TABLE `common_1`.`a` (`a` BIGINT PRIMARY KEY,`b` INT)", w.ddlList[0].Query) } +<<<<<<< HEAD func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) { // Scenario: // 1) TiCDC writes DML messages to Kafka in commitTs order. @@ -293,6 +294,99 @@ func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) // The kafka-consumer must not drop these "fallback commitTs" events unless they have // already been flushed to downstream (AppliedWatermark), otherwise the replay cannot // heal the missing window. +======= +func TestWriterWrite_sortsOutOfOrderDMLByWatermark(t *testing.T) { + ctx := context.Background() + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + flushedCommitTs := make([]uint64, 0) + s.EXPECT().AddDMLEvent(gomock.Any()).Do(func(event *commonEvent.DMLEvent) { + flushedCommitTs = append(flushedCommitTs, event.GetCommitTs()) + event.PostFlush() + }).Times(2) + + replicaCfg := config.GetDefaultReplicaConfig() + eventRouter, err := eventrouter.NewEventRouter(replicaCfg.Sink, "test-topic", false, false) + require.NoError(t, err) + + p := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + watermark: 0, + } + w := &writer{ + progresses: []*partitionProgress{p}, + mysqlSink: s, + eventRouter: eventRouter, + protocol: config.ProtocolOpen, + } + + w.appendMessage2Group(newDMLMessageForWriterTest(20), p, kafka.Offset(1)) + w.appendMessage2Group(newDMLMessageForWriterTest(10), p, kafka.Offset(2)) + w.appendMessage2Group(newDMLMessageForWriterTest(20), p, kafka.Offset(3)) + + p.watermark = 20 + require.True(t, w.Write(ctx, codeccommon.MessageTypeResolved)) + require.Equal(t, []uint64{10, 20}, flushedCommitTs) +} + +func TestWriteMessageIgnoresFallbackDMLBelowGlobalWatermark(t *testing.T) { + ctx := context.Background() + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + s.EXPECT().AddDMLEvent(gomock.Any()).Times(0) + + progress := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + watermark: 20, + decoder: &singleDMLDecoder{message: newDMLMessageForWriterTest(10)}, + } + w := &writer{ + progresses: []*partitionProgress{progress}, + mysqlSink: s, + protocol: config.ProtocolOpen, + maxBatchSize: 64, + maxMessageBytes: 1, + } + + needCommit := w.WriteMessage(ctx, &kafka.Message{ + TopicPartition: kafka.TopicPartition{Partition: 0, Offset: kafka.Offset(10)}, + }) + + require.False(t, needCommit) + require.Nil(t, progress.eventsGroup[1]) +} + +func TestAppendMessageKeepsFallbackDMLAboveGlobalWatermark(t *testing.T) { + replicaCfg := config.GetDefaultReplicaConfig() + eventRouter, err := eventrouter.NewEventRouter(replicaCfg.Sink, "test-topic", false, false) + require.NoError(t, err) + + progress := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + watermark: 20, + } + w := &writer{ + progresses: []*partitionProgress{ + progress, + {partition: 1, watermark: 5}, + }, + eventRouter: eventRouter, + protocol: config.ProtocolOpen, + } + + w.appendMessage2Group(newDMLMessageForWriterTest(10), progress, kafka.Offset(10)) + + require.NotNil(t, progress.eventsGroup[1]) + resolved := progress.eventsGroup[1].ResolveInto(20, nil) + require.Len(t, resolved, 1) + require.Equal(t, uint64(10), resolved[0].GetCommitTs()) +} + +func TestOnDDLMarksRoutedCreateTableLikePartitionTableForAvro(t *testing.T) { +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) replicaCfg := config.GetDefaultReplicaConfig() eventRouter, err := eventrouter.NewEventRouter(replicaCfg.Sink, "test-topic", false, false) require.NoError(t, err) @@ -340,3 +434,42 @@ func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) resolved = group.ResolveInto(150, resolvedEvents) require.Empty(t, resolved) } + +func newDMLMessageForWriterTest(commitTs uint64) *codeccommon.DMLMessage { + return codeccommon.NewDMLMessage(1, "test", "t", commitTs, common.RowTypeUpdate, func() *commonEvent.DMLEvent { + return &commonEvent.DMLEvent{ + PhysicalTableID: 1, + CommitTs: commitTs, + RowTypes: []common.RowType{common.RowTypeUpdate}, + Rows: chunk.NewChunkWithCapacity(nil, 0), + TableInfo: &common.TableInfo{ + TableName: common.TableName{Schema: "test", Table: "t", TableID: 1}, + }, + } + }) +} + +type singleDMLDecoder struct { + message *codeccommon.DMLMessage + consumed bool +} + +func (d *singleDMLDecoder) AddKeyValue(_, _ []byte) { +} + +func (d *singleDMLDecoder) HasNext() (codeccommon.MessageType, bool) { + return codeccommon.MessageTypeRow, !d.consumed +} + +func (d *singleDMLDecoder) NextResolvedEvent() uint64 { + return 0 +} + +func (d *singleDMLDecoder) NextDMLMessage() *codeccommon.DMLMessage { + d.consumed = true + return d.message +} + +func (d *singleDMLDecoder) NextDDLEvent() *commonEvent.DDLEvent { + return nil +} diff --git a/cmd/pulsar-consumer/writer.go b/cmd/pulsar-consumer/writer.go index fbbf94c3aa..e2083ee138 100644 --- a/cmd/pulsar-consumer/writer.go +++ b/cmd/pulsar-consumer/writer.go @@ -505,11 +505,25 @@ func (w *writer) appendRow2Group(dml *commonEvent.DMLEvent, progress *partitionP table = dml.TableInfo.GetTableName() commitTs = dml.GetCommitTs() ) + globalWatermark := w.globalWatermark() + if commitTs < globalWatermark { + log.Warn("DML event fallback row, since less than the global watermark, ignore it", + zap.Int64("tableID", tableID), zap.Int32("partition", progress.partition), + zap.Uint64("commitTs", commitTs), + zap.Uint64("globalWatermark", globalWatermark), + zap.Uint64("partitionWatermark", progress.watermark), + zap.String("schema", schema), zap.String("table", table), + zap.Stringer("eventType", message.RowType), + zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) + return + } + group := progress.eventsGroup[tableID] if group == nil { group = util.NewEventsGroup(progress.partition, tableID) progress.eventsGroup[tableID] = group } +<<<<<<< HEAD if commitTs <= group.AppliedWatermark { log.Warn("DML event replayed after applied, ignore it", zap.Int64("tableID", tableID), zap.Int32("partition", group.Partition), @@ -523,6 +537,21 @@ func (w *writer) appendRow2Group(dml *commonEvent.DMLEvent, progress *partitionP if forceInsert { log.Warn("DML event commit ts fallback, append with forceInsert", zap.Int32("partition", group.Partition), +======= + group.AppendMessage(message) + if commitTs < progress.watermark { + log.Warn("DML event fallback row, since less than the partition watermark, append it and sort before flush", + zap.Int64("tableID", tableID), zap.Int32("partition", group.Partition), + zap.Uint64("commitTs", commitTs), zap.Uint64("watermark", progress.watermark), + zap.Uint64("globalWatermark", globalWatermark), + zap.String("schema", schema), zap.String("table", table), + zap.Stringer("eventType", message.RowType), + zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) + return + } + if commitTs >= group.HighWatermark { + log.Debug("DML event append to the group", +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), zap.Uint64("appliedWatermark", group.AppliedWatermark), zap.Uint64("partitionWatermark", progress.watermark), @@ -532,6 +561,7 @@ func (w *writer) appendRow2Group(dml *commonEvent.DMLEvent, progress *partitionP group.Append(dml, true) return } +<<<<<<< HEAD group.Append(dml, false) log.Info("DML event append to the group", zap.Int32("partition", group.Partition), @@ -539,4 +569,13 @@ func (w *writer) appendRow2Group(dml *commonEvent.DMLEvent, progress *partitionP zap.Uint64("appliedWatermark", group.AppliedWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", dml.RowTypes[0])) +======= + log.Warn("DML event commit ts fallback, append it and sort before flush", + zap.Int32("partition", progress.partition), + zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), + zap.Any("partitionWatermark", progress.watermark), + zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), + zap.Stringer("eventType", message.RowType), + zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } diff --git a/cmd/pulsar-consumer/writer_test.go b/cmd/pulsar-consumer/writer_test.go index 8fa315d686..54bead2d6c 100644 --- a/cmd/pulsar-consumer/writer_test.go +++ b/cmd/pulsar-consumer/writer_test.go @@ -282,6 +282,7 @@ func TestWriterWrite_handlesOutOfOrderDDLsByCommitTs(t *testing.T) { require.Equal(t, "CREATE TABLE `common_1`.`a` (`a` BIGINT PRIMARY KEY,`b` INT)", w.ddlList[0].Query) } +<<<<<<< HEAD func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) { // Scenario: // 1) TiCDC writes DML messages to Pulsar in commitTs order. @@ -291,6 +292,95 @@ func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) // The pulsar-consumer must not drop these "fallback commitTs" events unless they // have already been flushed to downstream (AppliedWatermark), otherwise replayed // messages cannot heal missing windows. +======= +func TestWriterWrite_sortsOutOfOrderDMLByWatermark(t *testing.T) { + ctx := context.Background() + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + flushedCommitTs := make([]uint64, 0) + s.EXPECT().AddDMLEvent(gomock.Any()).Do(func(event *commonEvent.DMLEvent) { + flushedCommitTs = append(flushedCommitTs, event.GetCommitTs()) + event.PostFlush() + }).Times(2) + + p := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + watermark: 0, + } + w := &writer{ + progresses: []*partitionProgress{p}, + mysqlSink: s, + protocol: config.ProtocolCanalJSON, + } + + w.appendMessage2Group(newDMLMessageForWriterTest(20), p) + w.appendMessage2Group(newDMLMessageForWriterTest(10), p) + w.appendMessage2Group(newDMLMessageForWriterTest(20), p) + + p.watermark = 20 + require.True(t, w.Write(ctx, codeccommon.MessageTypeResolved)) + require.Equal(t, []uint64{10, 20}, flushedCommitTs) +} + +func TestWriteMessageIgnoresFallbackDMLBelowGlobalWatermark(t *testing.T) { + ctx := context.Background() + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + s.EXPECT().AddDMLEvent(gomock.Any()).Times(0) + + decoder := &deferredDMLDecoder{ + row: &commonEvent.DMLEvent{ + PhysicalTableID: 1, + CommitTs: 10, + RowTypes: []common.RowType{common.RowTypeInsert}, + TableInfo: &common.TableInfo{ + TableName: common.TableName{Schema: "test", Table: "t", TableID: 1}, + }, + }, + } + progress := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + watermark: 20, + decoder: decoder, + } + w := &writer{ + progresses: []*partitionProgress{progress}, + mysqlSink: s, + protocol: config.ProtocolCanalJSON, + } + + needCommit := w.WriteMessage(ctx, fakePulsarMessage{key: "k", payload: []byte(`{"fake":"row"}`)}) + + require.False(t, needCommit) + require.Nil(t, progress.eventsGroup[1]) +} + +func TestAppendMessageKeepsFallbackDMLAboveGlobalWatermark(t *testing.T) { + progress := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + watermark: 20, + } + w := &writer{ + progresses: []*partitionProgress{ + progress, + {partition: 1, watermark: 5}, + }, + protocol: config.ProtocolCanalJSON, + } + + w.appendMessage2Group(newDMLMessageForWriterTest(10), progress) + + require.NotNil(t, progress.eventsGroup[1]) + resolved := progress.eventsGroup[1].ResolveInto(20, nil) + require.Len(t, resolved, 1) + require.Equal(t, uint64(10), resolved[0].GetCommitTs()) +} + +func TestOnDDLMarksRoutedCreateTableLikePartitionTable(t *testing.T) { +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) w := &writer{ progresses: []*partitionProgress{ { @@ -328,6 +418,7 @@ func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) // Expect: commitTs=100 is still kept and can be resolved. resolved := group.ResolveInto(150, nil) require.Len(t, resolved, 1) +<<<<<<< HEAD require.Equal(t, uint64(100), resolved[0].CommitTs) // Step 3: once downstream has flushed beyond commitTs=100, replay is safe to ignore. @@ -335,4 +426,177 @@ func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) w.appendRow2Group(newDMLEvent(1, 100), progress) resolved = group.ResolveInto(150, nil) require.Empty(t, resolved) +======= + require.Equal(t, uint64(100), resolved[0].GetCommitTs()) +} + +func TestWriteMessageDefersDMLAssemblyUntilFlush(t *testing.T) { + ctx := context.Background() + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + s.EXPECT().AddDMLEvent(gomock.Any()).Do(func(event *commonEvent.DMLEvent) { + event.PostFlush() + }).Times(1) + + decoder := &deferredDMLDecoder{ + row: &commonEvent.DMLEvent{ + PhysicalTableID: 1, + CommitTs: 100, + RowTypes: []common.RowType{common.RowTypeInsert}, + TableInfo: &common.TableInfo{ + TableName: common.TableName{Schema: "test", Table: "t", TableID: 1}, + }, + }, + } + progress := &partitionProgress{ + partition: 0, + eventsGroup: make(map[int64]*util.EventsGroup), + decoder: decoder, + } + w := &writer{ + progresses: []*partitionProgress{progress}, + mysqlSink: s, + protocol: config.ProtocolCanalJSON, + } + + needCommit := w.WriteMessage(ctx, fakePulsarMessage{key: "k", payload: []byte(`{"fake":"row"}`)}) + require.False(t, needCommit) + require.Equal(t, 1, decoder.addKeyValueCount) + require.Equal(t, 1, decoder.hasNextCount) + require.Equal(t, 1, decoder.nextDMLMessageCount) + require.Zero(t, decoder.toDMLEventCount) + require.Len(t, progress.eventsGroup[1].ResolveInto(99, nil), 0) + + progress.watermark = 100 + require.True(t, w.Write(ctx, codeccommon.MessageTypeResolved)) + require.Equal(t, 1, decoder.addKeyValueCount) + require.Equal(t, 1, decoder.hasNextCount) + require.Equal(t, 1, decoder.nextDMLMessageCount) + require.Equal(t, 1, decoder.toDMLEventCount) + require.Empty(t, progress.eventsGroup[1].ResolveInto(100, nil)) + require.Equal(t, []byte(`{"fake":"row"}`), decoder.lastValue) +} + +type deferredDMLDecoder struct { + row *commonEvent.DMLEvent + + addKeyValueCount int + hasNextCount int + nextDMLMessageCount int + toDMLEventCount int + lastValue []byte +} + +func (d *deferredDMLDecoder) AddKeyValue(_, value []byte) { + d.addKeyValueCount++ + d.lastValue = append(d.lastValue[:0], value...) +} + +func (d *deferredDMLDecoder) HasNext() (codeccommon.MessageType, bool) { + d.hasNextCount++ + return codeccommon.MessageTypeRow, true +} + +func (d *deferredDMLDecoder) NextResolvedEvent() uint64 { + return 0 +} + +func (d *deferredDMLDecoder) NextDMLMessage() *codeccommon.DMLMessage { + d.nextDMLMessageCount++ + return codeccommon.NewDMLMessage(1, "test", "t", d.row.CommitTs, common.RowTypeInsert, func() *commonEvent.DMLEvent { + d.toDMLEventCount++ + return d.row + }) +} + +func (d *deferredDMLDecoder) NextDDLEvent() *commonEvent.DDLEvent { + return nil +} + +func newDMLMessageForWriterTest(commitTs uint64) *codeccommon.DMLMessage { + return codeccommon.NewDMLMessage(1, "test", "t", commitTs, common.RowTypeUpdate, func() *commonEvent.DMLEvent { + return &commonEvent.DMLEvent{ + PhysicalTableID: 1, + CommitTs: commitTs, + RowTypes: []common.RowType{common.RowTypeUpdate}, + Rows: chunk.NewChunkWithCapacity(nil, 0), + TableInfo: &common.TableInfo{ + TableName: common.TableName{Schema: "test", Table: "t", TableID: 1}, + }, + } + }) +} + +type fakePulsarMessage struct { + key string + payload []byte +} + +func (m fakePulsarMessage) Topic() string { + return "" +} + +func (m fakePulsarMessage) ProducerName() string { + return "" +} + +func (m fakePulsarMessage) Properties() map[string]string { + return nil +} + +func (m fakePulsarMessage) Payload() []byte { + return m.payload +} + +func (m fakePulsarMessage) ID() pulsar.MessageID { + return nil +} + +func (m fakePulsarMessage) PublishTime() time.Time { + return time.Time{} +} + +func (m fakePulsarMessage) EventTime() time.Time { + return time.Time{} +} + +func (m fakePulsarMessage) Key() string { + return m.key +} + +func (m fakePulsarMessage) OrderingKey() string { + return "" +} + +func (m fakePulsarMessage) RedeliveryCount() uint32 { + return 0 +} + +func (m fakePulsarMessage) IsReplicated() bool { + return false +} + +func (m fakePulsarMessage) GetReplicatedFrom() string { + return "" +} + +func (m fakePulsarMessage) GetSchemaValue(any) error { + return nil +} + +func (m fakePulsarMessage) SchemaVersion() []byte { + return nil +} + +func (m fakePulsarMessage) GetEncryptionContext() *pulsar.EncryptionContext { + return nil +} + +func (m fakePulsarMessage) Index() *uint64 { + return nil +} + +func (m fakePulsarMessage) BrokerPublishTime() *time.Time { + return nil +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index af913fdac6..a35d0f8ab1 100644 --- a/cmd/storage-consumer/consumer.go +++ b/cmd/storage-consumer/consumer.go @@ -250,8 +250,13 @@ func (c *consumer) appendRow2Group(dml *event.DMLEvent, enableTableAcrossNodes b c.eventsGroup[tableID] = group } if commitTs >= group.HighWatermark { +<<<<<<< HEAD group.Append(dml, false) log.Info("DML event append to the group", +======= + group.AppendMessage(message) + log.Debug("DML event append to the group", +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", dml.RowTypes[0])) @@ -261,8 +266,13 @@ func (c *consumer) appendRow2Group(dml *event.DMLEvent, enableTableAcrossNodes b log.Warn("DML events fallback, but enableTableAcrossNodes is true, still append it", zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), +<<<<<<< HEAD zap.Stringer("eventType", dml.RowTypes[0])) group.Append(dml, true) +======= + zap.Stringer("eventType", message.RowType)) + group.AppendMessage(message) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) return } log.Warn("dml event commit ts fallback, ignore", diff --git a/cmd/util/event_group.go b/cmd/util/event_group.go index 95ff621510..72297db001 100644 --- a/cmd/util/event_group.go +++ b/cmd/util/event_group.go @@ -14,7 +14,11 @@ package util import ( +<<<<<<< HEAD "slices" +======= + "math" +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) "sort" "github.com/pingcap/log" @@ -53,6 +57,7 @@ func NewEventsGroup(partition int32, tableID int64) *EventsGroup { } } +<<<<<<< HEAD // Append will append an event to event groups. func (g *EventsGroup) Append(row *commonEvent.DMLEvent, force bool) { if row.CommitTs > g.HighWatermark { @@ -144,13 +149,88 @@ func (g *EventsGroup) ResolveInto(resolve uint64, dst []*commonEvent.DMLEvent) [ zap.Int32("partition", g.Partition), zap.Int64("tableID", g.tableID), zap.Int("resolved", i), zap.Int("remained", len(g.events)), zap.Uint64("resolveTs", resolve), zap.Uint64("firstCommitTs", g.events[0].CommitTs)) +======= +// AppendMessage appends a message to event groups. +func (g *EventsGroup) AppendMessage(message *codeccommon.DMLMessage) { + commitTs := message.GetCommitTs() + if commitTs > g.HighWatermark { + g.HighWatermark = commitTs + } + g.messages = append(g.messages, message) +} + +// ResolveInto appends all messages with CommitTs <= resolve into dst in commit-ts order and removes +// them from the group. ResolveInto copies pointers into dst first, then clears the resolved messages +// so Go GC can reclaim them once downstream is done with them. +func (g *EventsGroup) ResolveInto(resolve uint64, dst []*codeccommon.DMLMessage) []*codeccommon.DMLMessage { + if len(g.messages) == 0 { + return dst + } + + original := g.messages + remaining := g.messages[:0] + resolved := make([]*codeccommon.DMLMessage, 0, len(g.messages)) + + var ( + lastCommitTs uint64 + outOfOrder bool + outOfOrderLastTs uint64 + outOfOrderCommitTs uint64 + ) + for _, message := range g.messages { + commitTs := message.GetCommitTs() + if commitTs > resolve { + remaining = append(remaining, message) + continue + } + if len(resolved) > 0 && commitTs < lastCommitTs && !outOfOrder { + outOfOrder = true + outOfOrderLastTs = lastCommitTs + outOfOrderCommitTs = commitTs + } + lastCommitTs = commitTs + resolved = append(resolved, message) + } + if len(resolved) == 0 { + return dst + } + + if outOfOrder { + log.Warn("DML events are out of order before flush, sort them", + zap.Int32("partition", g.Partition), + zap.Int64("tableID", g.tableID), + zap.Uint64("resolveTs", resolve), + zap.Int("resolved", len(resolved)), + zap.Uint64("lastCommitTs", outOfOrderLastTs), + zap.Uint64("commitTs", outOfOrderCommitTs)) + sort.SliceStable(resolved, func(i, j int) bool { + return resolved[i].GetCommitTs() < resolved[j].GetCommitTs() + }) + } + + dst = append(dst, resolved...) + clear(original[len(remaining):]) + g.messages = remaining + if len(g.messages) != 0 { + firstCommitTs := g.messages[0].GetCommitTs() + log.Debug("not all events resolved", + zap.Int32("partition", g.Partition), zap.Int64("tableID", g.tableID), + zap.Int("resolved", len(resolved)), zap.Int("remained", len(g.messages)), + zap.Uint64("resolveTs", resolve), zap.Uint64("firstCommitTs", firstCommitTs)) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } return dst } +<<<<<<< HEAD // GetAllEvents will get all events. func (g *EventsGroup) GetAllEvents() []*commonEvent.DMLEvent { result := g.events g.events = nil return result +======= +// GetAllMessages gets all messages. +func (g *EventsGroup) GetAllMessages() []*codeccommon.DMLMessage { + return g.ResolveInto(math.MaxUint64, nil) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } diff --git a/cmd/util/event_group_test.go b/cmd/util/event_group_test.go index 5b8816ea50..805a8ccd8a 100644 --- a/cmd/util/event_group_test.go +++ b/cmd/util/event_group_test.go @@ -64,17 +64,18 @@ func TestEventsGroupAppendForceMergesExistingCommitTs(t *testing.T) { require.Len(t, dst[0].RowTypes, 2) } -func TestEventsGroupResolveIntoAppendsAndClearsResolvedPrefix(t *testing.T) { - // Scenario: A consumer resolves a prefix of events by watermark/commit-ts and appends them - // into a downstream batch slice. We must clear the resolved prefix in the group's backing - // array to avoid retaining already-flushed events and causing unbounded memory growth. +func TestEventsGroupResolveIntoAppendsAndClearsResolvedMessages(t *testing.T) { + // Scenario: A consumer resolves events by watermark/commit-ts and appends them into a downstream + // batch slice. We must clear resolved messages in the group's backing array to avoid retaining + // already-flushed events and causing unbounded memory growth. // // Steps: // 1. Append 3 events with increasing CommitTs. // 2. Call ResolveInto with resolve=2 and a nil dst. // 3. Verify (a) returned events are correct, (b) group keeps only the remaining event, - // (c) the resolved prefix in the original backing slice is cleared (nil'd). + // (c) resolved messages in the original backing slice are cleared (nil'd). group := NewEventsGroup(0, 1) +<<<<<<< HEAD e1 := &commonEvent.DMLEvent{CommitTs: 1} e2 := &commonEvent.DMLEvent{CommitTs: 2} e3 := &commonEvent.DMLEvent{CommitTs: 3} @@ -85,6 +86,18 @@ func TestEventsGroupResolveIntoAppendsAndClearsResolvedPrefix(t *testing.T) { // Keep a reference to the original slice header so we can validate that ResolveInto clears // the resolved prefix in-place (this is what prevents GC retention of flushed events). original := group.events +======= + m1 := newTestDMLMessage(1) + m2 := newTestDMLMessage(2) + m3 := newTestDMLMessage(3) + group.AppendMessage(m1) + group.AppendMessage(m2) + group.AppendMessage(m3) + + // Keep a reference to the original slice header so we can validate that ResolveInto clears + // resolved messages in-place (this is what prevents GC retention of flushed events). + original := group.messages +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) var dst []*commonEvent.DMLEvent dst = group.ResolveInto(2, dst) @@ -96,21 +109,32 @@ func TestEventsGroupResolveIntoAppendsAndClearsResolvedPrefix(t *testing.T) { require.Len(t, group.events, 1) require.Same(t, e3, group.events[0]) - // The resolved prefix must be nil so the group doesn't keep flushed events alive via its - // backing array (classic Go slice memory retention pitfall). - require.Nil(t, original[0]) + // The unresolved event is compacted to the front, and the tail is cleared so the group + // doesn't keep flushed events alive via its backing array. + require.Same(t, m3, original[0]) require.Nil(t, original[1]) +<<<<<<< HEAD require.Same(t, e3, original[2]) +======= + require.Nil(t, original[2]) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } func TestEventsGroupResolveIntoNoopWhenNothingResolved(t *testing.T) { // Scenario: resolveTs is behind all buffered events. // Expectation: ResolveInto should be a no-op (dst unchanged, group unchanged). group := NewEventsGroup(0, 1) +<<<<<<< HEAD e1 := &commonEvent.DMLEvent{CommitTs: 10} e2 := &commonEvent.DMLEvent{CommitTs: 20} group.Append(e1, false) group.Append(e2, false) +======= + m1 := newTestDMLMessage(10) + m2 := newTestDMLMessage(20) + group.AppendMessage(m1) + group.AppendMessage(m2) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) original := group.events dst := make([]*commonEvent.DMLEvent, 0, 1) @@ -130,10 +154,17 @@ func TestEventsGroupResolveIntoClearsAllWhenFullyResolved(t *testing.T) { // Scenario: resolveTs advances beyond all buffered events. // Expectation: group is emptied and all backing-array pointers for resolved events are cleared. group := NewEventsGroup(0, 1) +<<<<<<< HEAD e1 := &commonEvent.DMLEvent{CommitTs: 1} e2 := &commonEvent.DMLEvent{CommitTs: 2} group.Append(e1, false) group.Append(e2, false) +======= + m1 := newTestDMLMessage(1) + m2 := newTestDMLMessage(2) + group.AppendMessage(m1) + group.AppendMessage(m2) +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) original := group.events var dst []*commonEvent.DMLEvent @@ -147,3 +178,98 @@ func TestEventsGroupResolveIntoClearsAllWhenFullyResolved(t *testing.T) { require.Nil(t, original[0]) require.Nil(t, original[1]) } +<<<<<<< HEAD +======= + +func TestEventsGroupResolveIntoSortsOutOfOrderResolvedMessages(t *testing.T) { + group := NewEventsGroup(0, 1) + m1 := newTestDMLMessage(20) + m2 := newTestDMLMessage(10) + m3 := newTestDMLMessage(30) + group.AppendMessage(m1) + group.AppendMessage(m2) + group.AppendMessage(m3) + + original := group.messages + var dst []*codeccommon.DMLMessage + dst = group.ResolveInto(25, dst) + + require.Len(t, dst, 2) + require.Same(t, m2, dst[0]) + require.Same(t, m1, dst[1]) + + require.Len(t, group.messages, 1) + require.Same(t, m3, group.messages[0]) + require.Same(t, m3, original[0]) + require.Nil(t, original[1]) + require.Nil(t, original[2]) +} + +func TestEventsGroupResolveIntoKeepsSameCommitTsStable(t *testing.T) { + group := NewEventsGroup(0, 1) + m1 := newTestDMLMessage(20) + m2 := newTestDMLMessage(10) + m3 := newTestDMLMessage(20) + group.AppendMessage(m1) + group.AppendMessage(m2) + group.AppendMessage(m3) + + var dst []*codeccommon.DMLMessage + dst = group.ResolveInto(20, dst) + + require.Len(t, dst, 3) + require.Same(t, m2, dst[0]) + require.Same(t, m1, dst[1]) + require.Same(t, m3, dst[2]) + require.Empty(t, group.messages) +} + +func TestEventsGroupGetAllMessagesSortsOutOfOrderMessages(t *testing.T) { + group := NewEventsGroup(0, 1) + m1 := newTestDMLMessage(20) + m2 := newTestDMLMessage(10) + m3 := newTestDMLMessage(30) + group.AppendMessage(m1) + group.AppendMessage(m2) + group.AppendMessage(m3) + + messages := group.GetAllMessages() + + require.Len(t, messages, 3) + require.Same(t, m2, messages[0]) + require.Same(t, m1, messages[1]) + require.Same(t, m3, messages[2]) + require.Empty(t, group.messages) +} + +func TestAppendOrMergeDMLEventMergesSameCommitTs(t *testing.T) { + var flushed []int + e1 := newTestDMLEvent(10, common.RowTypeInsert) + e1.AddPostFlushFunc(func() { flushed = append(flushed, 1) }) + e2 := newTestDMLEvent(10, common.RowTypeDelete) + e2.AddPostFlushFunc(func() { flushed = append(flushed, 2) }) + + events := AppendOrMergeDMLEvent(nil, e1) + events = AppendOrMergeDMLEvent(events, e2) + + require.Len(t, events, 1) + require.Same(t, e1, events[0]) + require.Equal(t, int32(2), events[0].Length) + require.Equal(t, []common.RowType{common.RowTypeInsert, common.RowTypeDelete}, events[0].RowTypes) + + events[0].PostFlush() + require.Equal(t, []int{1, 2}, flushed) +} + +func TestAppendOrMergeDMLEventAppendsDifferentCommitTs(t *testing.T) { + e1 := newTestDMLEvent(10, common.RowTypeInsert) + e2 := newTestDMLEvent(20, common.RowTypeDelete) + + events := AppendOrMergeDMLEvent(nil, e1) + events = AppendOrMergeDMLEvent(events, e2) + + require.Len(t, events, 2) + require.Same(t, e1, events[0]) + require.Same(t, e2, events[1]) +} +>>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) From cd3cd48b5bcd630d5c555bab0562e7ed4e40aa11 Mon Sep 17 00:00:00 2001 From: wk989898 Date: Mon, 3 Aug 2026 08:01:20 +0000 Subject: [PATCH 2/2] update Signed-off-by: wk989898 --- cmd/kafka-consumer/writer.go | 83 ++---------------- cmd/kafka-consumer/writer_test.go | 136 ++++++++++++++++------------- cmd/pulsar-consumer/writer.go | 66 +------------- cmd/pulsar-consumer/writer_test.go | 92 ++++++------------- cmd/storage-consumer/consumer.go | 10 --- cmd/util/event_group.go | 98 ++++++--------------- cmd/util/event_group_test.go | 80 +++-------------- 7 files changed, 148 insertions(+), 417 deletions(-) diff --git a/cmd/kafka-consumer/writer.go b/cmd/kafka-consumer/writer.go index 73018922fa..fa68a1d36e 100644 --- a/cmd/kafka-consumer/writer.go +++ b/cmd/kafka-consumer/writer.go @@ -64,10 +64,8 @@ func (p *partitionProgress) updateWatermark(newWatermark uint64, offset kafka.Of zap.Uint64("watermark", newWatermark)) return } - readOldOffset := true - if offset > p.watermarkOffset { - readOldOffset = false - } + readOldOffset := offset <= p.watermarkOffset + log.Warn("partition resolved ts fall back, ignore it", zap.Bool("readOldOffset", readOldOffset), zap.Int32("partition", p.partition), @@ -121,7 +119,8 @@ func newWriter(ctx context.Context, o *option) *writer { w.progresses[i] = newPartitionProgress(int32(i), decoder) } - eventRouter, err := eventrouter.NewEventRouter(o.sinkConfig, o.topic, false, o.protocol == config.ProtocolAvro) + isAvroLike := o.protocol == config.ProtocolAvro + eventRouter, err := eventrouter.NewEventRouter(o.sinkConfig, o.topic, false, isAvroLike) if err != nil { log.Panic("initialize the event router failed", zap.Any("protocol", o.protocol), zap.Any("topic", o.topic), @@ -158,12 +157,6 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *event.DDLEvent) error { tableIDs := w.getBlockTableIDs(ddl) commitTs := ddl.GetCommitTs() resolvedEvents := make([]*event.DMLEvent, 0) - // resolvedGroups records which EventsGroup has flushed events so we can - // advance its AppliedWatermark after the flush is fully finished. - resolvedGroups := make([]struct { - group *util.EventsGroup - maxCommitTs uint64 - }, 0) for tableID := range tableIDs { for _, progress := range w.progresses { g, ok := progress.eventsGroup[tableID] @@ -204,11 +197,6 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *event.DDLEvent) error { log.Info("flush DML events before DDL done", zap.Uint64("DDLCommitTs", commitTs), zap.Int("total", total), zap.Duration("duration", time.Since(start)), zap.Any("tables", tableIDs)) - for _, item := range resolvedGroups { - if item.maxCommitTs > item.group.AppliedWatermark { - item.group.AppliedWatermark = item.maxCommitTs - } - } return w.mysqlSink.WriteBlockEvent(ddl) case <-ticker.C: log.Warn("DML events cannot be flushed in time", @@ -285,12 +273,6 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error { watermark := w.globalWatermark() resolvedEvents := make([]*event.DMLEvent, 0) - // resolvedGroups records which EventsGroup has flushed events so we can - // advance its AppliedWatermark after the flush is fully finished. - resolvedGroups := make([]struct { - group *util.EventsGroup - maxCommitTs uint64 - }, 0) for _, p := range w.progresses { for _, group := range p.eventsGroup { messages := group.ResolveInto(watermark, nil) @@ -327,11 +309,6 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error { case <-done: log.Info("flush DML events done", zap.Uint64("watermark", watermark), zap.Int("total", total), zap.Duration("duration", time.Since(start))) - for _, item := range resolvedGroups { - if item.maxCommitTs > item.group.AppliedWatermark { - item.group.AppliedWatermark = item.maxCommitTs - } - } return nil case <-ticker.C: log.Warn("DML events cannot be flushed in time", zap.Uint64("watermark", watermark), @@ -627,10 +604,10 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti // if the kafka cluster is normal, this should not hit. // else if the cluster is abnormal, the consumer may consume old message, then cause the watermark fallback. var ( - tableID = dml.GetTableID() - schema = dml.TableInfo.GetSchemaName() - table = dml.TableInfo.GetTableName() - commitTs = dml.GetCommitTs() + tableID = message.TableID + schema = message.Schema + table = message.Table + commitTs = message.GetCommitTs() ) globalWatermark := w.globalWatermark() if commitTs < globalWatermark { @@ -651,49 +628,6 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti group = util.NewEventsGroup(progress.partition, tableID) progress.eventsGroup[tableID] = group } -<<<<<<< HEAD - // IMPORTANT: Kafka offsets are append-only, but CommitTs can go backwards after - // a TiCDC restart/retry (at-least-once replay). We must not drop such events - // solely based on a "seen" watermark (e.g. HighWatermark). The only safe - // ignore condition is "already flushed to downstream". - if commitTs <= group.AppliedWatermark { - log.Warn("DML event replayed after applied, ignore it", - zap.Int64("tableID", tableID), zap.Int32("partition", group.Partition), - zap.Uint64("commitTs", commitTs), zap.Any("offset", offset), - zap.Uint64("appliedWatermark", group.AppliedWatermark), zap.Uint64("highWatermark", group.HighWatermark), - zap.Uint64("partitionWatermark", progress.watermark), zap.Any("watermarkOffset", progress.watermarkOffset), - zap.String("schema", schema), zap.String("table", table), zap.Any("protocol", w.protocol)) - return - } - if commitTs >= group.HighWatermark { - message = w.messageWithPartitionCheck(message, progress.partition, offset) - group.AppendMessage(message, false) - log.Debug("DML event append to the group", - zap.Int32("partition", group.Partition), zap.Any("offset", offset), - zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), - zap.Uint64("appliedWatermark", group.AppliedWatermark), - zap.Uint64("partitionWatermark", progress.watermark), zap.Any("watermarkOffset", progress.watermarkOffset), - zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), - zap.Stringer("eventType", message.RowType)) - return - } - if w.enableTableAcrossNodes { - log.Warn("DML events fallback, but enableTableAcrossNodes is true, still append it", - zap.Int32("partition", group.Partition), zap.Any("offset", offset), - zap.Uint64("commitTs", commitTs), zap.Uint64("HighWatermark", group.HighWatermark), - zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), - zap.Stringer("eventType", message.RowType)) - group.AppendMessage(w.messageWithPartitionCheck(message, progress.partition, offset), true) - return - } - group.Append(dml, false) - log.Info("DML event append to the group", - zap.Int32("partition", group.Partition), zap.Any("offset", offset), - zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), - zap.Uint64("appliedWatermark", group.AppliedWatermark), - zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), - zap.Stringer("eventType", dml.RowTypes[0])) -======= message = w.messageWithPartitionCheck(message, progress.partition, offset) group.AppendMessage(message) if commitTs < progress.watermark { @@ -722,7 +656,6 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", message.RowType), zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } func openDB(ctx context.Context, dsn string) (*sql.DB, error) { diff --git a/cmd/kafka-consumer/writer_test.go b/cmd/kafka-consumer/writer_test.go index 2c78822bc3..e5845efb03 100644 --- a/cmd/kafka-consumer/writer_test.go +++ b/cmd/kafka-consumer/writer_test.go @@ -18,9 +18,10 @@ import ( "testing" "github.com/confluentinc/confluent-kafka-go/v2/kafka" + "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/cmd/util" - "github.com/pingcap/ticdc/downstreamadapter/sink" "github.com/pingcap/ticdc/downstreamadapter/sink/eventrouter" + sinkmock "github.com/pingcap/ticdc/downstreamadapter/sink/mock" "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/config" @@ -30,40 +31,23 @@ import ( "github.com/stretchr/testify/require" ) -// recordingSink is a minimal sink.Sink implementation that records which DDLs are executed. -// -// It lets unit tests validate consumer-side DDL flushing behavior without requiring a real downstream. -type recordingSink struct { - ddls []string -} - -var _ sink.Sink = (*recordingSink)(nil) +func newMockSink(t *testing.T) (*sinkmock.MockSink, *[]string) { + t.Helper() -func (s *recordingSink) SinkType() common.SinkType { return common.MysqlSinkType } -func (s *recordingSink) IsNormal() bool { return true } -func (s *recordingSink) AddDMLEvent(_ *commonEvent.DMLEvent) { -} - -func (s *recordingSink) FlushDMLBeforeBlock(_ commonEvent.BlockEvent) error { - return nil -} - -func (s *recordingSink) WriteBlockEvent(event commonEvent.BlockEvent) error { - if ddl, ok := event.(*commonEvent.DDLEvent); ok { - s.ddls = append(s.ddls, ddl.Query) - } - return nil -} - -func (s *recordingSink) AddCheckpointTs(_ uint64) { -} + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + ddls := make([]string, 0) -func (s *recordingSink) SetTableSchemaStore(_ *commonEvent.TableSchemaStore) { -} + s.EXPECT().AddDMLEvent(gomock.Any()).AnyTimes() + s.EXPECT().WriteBlockEvent(gomock.Any()).DoAndReturn(func(event commonEvent.BlockEvent) error { + if ddl, ok := event.(*commonEvent.DDLEvent); ok { + ddls = append(ddls, ddl.Query) + } + return nil + }).AnyTimes() -func (s *recordingSink) Close() { + return s, &ddls } -func (s *recordingSink) Run(_ context.Context) error { return nil } func TestWriterWrite_executesIndependentCreateTableWithoutWatermark(t *testing.T) { // Scenario: In some integration tests the upstream intentionally pauses dispatcher creation, which can @@ -75,7 +59,7 @@ func TestWriterWrite_executesIndependentCreateTableWithoutWatermark(t *testing.T // 2) Call writer.Write and expect the DDL is executed to advance downstream schema even without the // watermark catching up. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) w := &writer{ progresses: []*partitionProgress{ {partition: 0, watermark: 0}, @@ -100,7 +84,7 @@ func TestWriterWrite_executesIndependentCreateTableWithoutWatermark(t *testing.T w.Write(ctx, codeccommon.MessageTypeDDL) - require.Equal(t, []string{"CREATE TABLE `test`.`t` (`id` INT PRIMARY KEY)"}, s.ddls) + require.Equal(t, []string{"CREATE TABLE `test`.`t` (`id` INT PRIMARY KEY)"}, *ddls) require.Empty(t, w.ddlList) } @@ -113,7 +97,7 @@ func TestWriterWrite_preservesOrderWhenBlockedDDLNotReady(t *testing.T) { // 2) Call writer.Write and expect nothing executes. // 3) Advance watermark beyond the first DDL and expect both execute in order. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) p := &partitionProgress{partition: 0, watermark: 0} w := &writer{ progresses: []*partitionProgress{p}, @@ -145,7 +129,7 @@ func TestWriterWrite_preservesOrderWhenBlockedDDLNotReady(t *testing.T) { } w.Write(ctx, codeccommon.MessageTypeDDL) - require.Empty(t, s.ddls) + require.Empty(t, *ddls) require.Len(t, w.ddlList, 2) p.watermark = 200 @@ -153,7 +137,7 @@ func TestWriterWrite_preservesOrderWhenBlockedDDLNotReady(t *testing.T) { require.Equal(t, []string{ "ALTER TABLE `test`.`t` ADD COLUMN `c2` INT", "CREATE TABLE `test`.`t2` (`id` INT PRIMARY KEY)", - }, s.ddls) + }, *ddls) require.Empty(t, w.ddlList) } @@ -166,7 +150,7 @@ func TestWriterWrite_doesNotBypassWatermarkForCreateTableLike(t *testing.T) { // 2) Call writer.Write and expect the DDL is NOT executed. // 3) Advance watermark beyond the DDL commitTs and expect the DDL executes. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) p := &partitionProgress{partition: 0, watermark: 0} w := &writer{ progresses: []*partitionProgress{p}, @@ -191,12 +175,12 @@ func TestWriterWrite_doesNotBypassWatermarkForCreateTableLike(t *testing.T) { } w.Write(ctx, codeccommon.MessageTypeDDL) - require.Empty(t, s.ddls) + require.Empty(t, *ddls) require.Len(t, w.ddlList, 1) p.watermark = 200 w.Write(ctx, codeccommon.MessageTypeDDL) - require.Equal(t, []string{"CREATE TABLE `test`.`t2` LIKE `test`.`t1`"}, s.ddls) + require.Equal(t, []string{"CREATE TABLE `test`.`t2` LIKE `test`.`t1`"}, *ddls) require.Empty(t, w.ddlList) } @@ -210,7 +194,7 @@ func TestWriterWrite_handlesOutOfOrderDDLsByCommitTs(t *testing.T) { // 2) Call writer.Write and expect all DDLs with commitTs <= watermark execute (in commit-ts order), // and only the truly "future" DDL remains pending. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) p := &partitionProgress{partition: 0, watermark: 944040962} w := &writer{ progresses: []*partitionProgress{p}, @@ -279,22 +263,11 @@ func TestWriterWrite_handlesOutOfOrderDDLsByCommitTs(t *testing.T) { "ALTER TABLE `common_1`.`add_and_drop_columns` ADD COLUMN `col1` INT NULL, ADD COLUMN `col2` INT NULL, ADD COLUMN `col3` INT NULL", "ALTER TABLE `common_1`.`add_and_drop_columns` DROP COLUMN `col1`, DROP COLUMN `col2`", "CREATE DATABASE `common`", - }, s.ddls) + }, *ddls) require.Len(t, w.ddlList, 1) require.Equal(t, "CREATE TABLE `common_1`.`a` (`a` BIGINT PRIMARY KEY,`b` INT)", w.ddlList[0].Query) } -<<<<<<< HEAD -func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) { - // Scenario: - // 1) TiCDC writes DML messages to Kafka in commitTs order. - // 2) Under network partition / changefeed restart, TiCDC may replay older commitTs, - // which will be appended to Kafka at a larger offset (commitTs appears to go backwards). - // - // The kafka-consumer must not drop these "fallback commitTs" events unless they have - // already been flushed to downstream (AppliedWatermark), otherwise the replay cannot - // heal the missing window. -======= func TestWriterWrite_sortsOutOfOrderDMLByWatermark(t *testing.T) { ctx := context.Background() ctrl := gomock.NewController(t) @@ -386,7 +359,6 @@ func TestAppendMessageKeepsFallbackDMLAboveGlobalWatermark(t *testing.T) { } func TestOnDDLMarksRoutedCreateTableLikePartitionTableForAvro(t *testing.T) { ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) replicaCfg := config.GetDefaultReplicaConfig() eventRouter, err := eventrouter.NewEventRouter(replicaCfg.Sink, "test-topic", false, true) require.NoError(t, err) @@ -434,14 +406,56 @@ func TestOnDDLMarksRoutedCreateTableLikePartitionTableForAvro(t *testing.T) { resolved := progress.eventsGroup[1].ResolveInto(150, nil) require.Len(t, resolved, 1) - require.Equal(t, uint64(100), resolved[0].CommitTs) - - // Step 3: once downstream has flushed beyond commitTs=100, the replay is safe to ignore. - resolvedEvents = make([]*commonEvent.DMLEvent, 0) - group.AppliedWatermark = 200 - w.appendRow2Group(newDMLEvent(1, 100), progress, kafka.Offset(12)) - resolved = group.ResolveInto(150, resolvedEvents) - require.Empty(t, resolved) + require.Equal(t, uint64(100), resolved[0].GetCommitTs()) +} + +func TestAppendRow2GroupKeepsDebeziumPartitionTableFallback(t *testing.T) { + for _, protocol := range []config.Protocol{ + config.ProtocolDebezium, + } { + t.Run(protocol.String(), func(t *testing.T) { + replicaCfg := config.GetDefaultReplicaConfig() + eventRouter, err := eventrouter.NewEventRouter(replicaCfg.Sink, "test-topic", false, false) + require.NoError(t, err) + + w := &writer{ + progresses: []*partitionProgress{{partition: 0, eventsGroup: make(map[int64]*util.EventsGroup)}}, + eventRouter: eventRouter, + protocol: protocol, + partitionTableAccessor: codeccommon.NewPartitionTableAccessor(), + } + + w.partitionTableAccessor.Add("target", "src") + ddl := &commonEvent.DDLEvent{ + Query: "CREATE TABLE `target`.`dst` LIKE `target`.`src`", + SchemaName: "target", + TableName: "dst", + Type: byte(timodel.ActionCreateTable), + } + w.onDDL(ddl) + require.True(t, w.partitionTableAccessor.IsPartitionTable("target", "dst")) + + newDMLEvent := func(commitTs uint64) *commonEvent.DMLEvent { + return &commonEvent.DMLEvent{ + PhysicalTableID: 1, + CommitTs: commitTs, + RowTypes: []common.RowType{common.RowTypeUpdate}, + Rows: chunk.NewChunkWithCapacity(nil, 0), + TableInfo: &common.TableInfo{ + TableName: common.TableName{Schema: "target", Table: "dst"}, + }, + } + } + + progress := w.progresses[0] + w.appendMessage2Group(codeccommon.NewDMLMessageFromEvent(newDMLEvent(200)), progress, kafka.Offset(10)) + w.appendMessage2Group(codeccommon.NewDMLMessageFromEvent(newDMLEvent(100)), progress, kafka.Offset(11)) + + resolved := progress.eventsGroup[1].ResolveInto(150, nil) + require.Len(t, resolved, 1) + require.Equal(t, uint64(100), resolved[0].GetCommitTs()) + }) + } } func newDMLMessageForWriterTest(commitTs uint64) *codeccommon.DMLMessage { diff --git a/cmd/pulsar-consumer/writer.go b/cmd/pulsar-consumer/writer.go index 017572665a..afd0aeddca 100644 --- a/cmd/pulsar-consumer/writer.go +++ b/cmd/pulsar-consumer/writer.go @@ -149,12 +149,6 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *commonEvent.DDLEvent) e tableIDs := w.getBlockTableIDs(ddl) commitTs := ddl.GetCommitTs() resolvedEvents := make([]*commonEvent.DMLEvent, 0) - // resolvedGroups records which EventsGroup has flushed events so we can - // advance its AppliedWatermark after the flush is fully finished. - resolvedGroups := make([]struct { - group *util.EventsGroup - maxCommitTs uint64 - }, 0) for tableID := range tableIDs { for _, progress := range w.progresses { g, ok := progress.eventsGroup[tableID] @@ -195,11 +189,6 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *commonEvent.DDLEvent) e log.Info("flush DML events before DDL done", zap.Uint64("DDLCommitTs", commitTs), zap.Int("total", total), zap.Duration("duration", time.Since(start)), zap.Any("tables", tableIDs)) - for _, item := range resolvedGroups { - if item.maxCommitTs > item.group.AppliedWatermark { - item.group.AppliedWatermark = item.maxCommitTs - } - } return w.mysqlSink.WriteBlockEvent(ddl) case <-ticker.C: log.Warn("DML events cannot be flushed in time", @@ -276,12 +265,6 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error { watermark := w.globalWatermark() resolvedEvents := make([]*commonEvent.DMLEvent, 0) - // resolvedGroups records which EventsGroup has flushed events so we can - // advance its AppliedWatermark after the flush is fully finished. - resolvedGroups := make([]struct { - group *util.EventsGroup - maxCommitTs uint64 - }, 0) for _, p := range w.progresses { for _, group := range p.eventsGroup { messages := group.ResolveInto(watermark, nil) @@ -316,11 +299,6 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error { case <-done: log.Info("flush DML events done", zap.Uint64("watermark", watermark), zap.Int("total", total), zap.Duration("duration", time.Since(start))) - for _, item := range resolvedGroups { - if item.maxCommitTs > item.group.AppliedWatermark { - item.group.AppliedWatermark = item.maxCommitTs - } - } return nil case <-ticker.C: log.Warn("DML events cannot be flushed in time", zap.Uint64("watermark", watermark), @@ -515,10 +493,10 @@ func (w *writer) addPartitionTable(schema, table string) { func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *partitionProgress) { var ( - tableID = dml.GetTableID() - schema = dml.TableInfo.GetSchemaName() - table = dml.TableInfo.GetTableName() - commitTs = dml.GetCommitTs() + tableID = message.TableID + schema = message.Schema + table = message.Table + commitTs = message.GetCommitTs() ) globalWatermark := w.globalWatermark() if commitTs < globalWatermark { @@ -538,21 +516,6 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti group = util.NewEventsGroup(progress.partition, tableID) progress.eventsGroup[tableID] = group } -<<<<<<< HEAD - if commitTs <= group.AppliedWatermark { - log.Warn("DML event replayed after applied, ignore it", - zap.Int64("tableID", tableID), zap.Int32("partition", group.Partition), - zap.Uint64("commitTs", commitTs), - zap.Uint64("appliedWatermark", group.AppliedWatermark), zap.Uint64("highWatermark", group.HighWatermark), - zap.Uint64("partitionWatermark", progress.watermark), - zap.String("schema", schema), zap.String("table", table), zap.Any("protocol", w.protocol)) - return - } - forceInsert := commitTs < group.HighWatermark || commitTs < progress.watermark || w.enableTableAcrossNodes - if forceInsert { - log.Warn("DML event commit ts fallback, append with forceInsert", - zap.Int32("partition", group.Partition), -======= group.AppendMessage(message) if commitTs < progress.watermark { log.Warn("DML event fallback row, since less than the partition watermark, append it and sort before flush", @@ -566,31 +529,11 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti } if commitTs >= group.HighWatermark { log.Debug("DML event append to the group", ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) - zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), - zap.Uint64("appliedWatermark", group.AppliedWatermark), - zap.Uint64("partitionWatermark", progress.watermark), - zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), - zap.Stringer("eventType", message.RowType)) - return - } - if w.enableTableAcrossNodes { - log.Warn("DML events fallback, but enableTableAcrossNodes is true, still append it", zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", message.RowType)) - group.AppendMessage(message, true) return } -<<<<<<< HEAD - group.Append(dml, false) - log.Info("DML event append to the group", - zap.Int32("partition", group.Partition), - zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), - zap.Uint64("appliedWatermark", group.AppliedWatermark), - zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), - zap.Stringer("eventType", dml.RowTypes[0])) -======= log.Warn("DML event commit ts fallback, append it and sort before flush", zap.Int32("partition", progress.partition), zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), @@ -598,5 +541,4 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", message.RowType), zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes)) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } diff --git a/cmd/pulsar-consumer/writer_test.go b/cmd/pulsar-consumer/writer_test.go index 1c6b731bc9..35c9c037e3 100644 --- a/cmd/pulsar-consumer/writer_test.go +++ b/cmd/pulsar-consumer/writer_test.go @@ -21,50 +21,33 @@ import ( "github.com/apache/pulsar-client-go/pulsar" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/cmd/util" - "github.com/pingcap/ticdc/downstreamadapter/sink" sinkmock "github.com/pingcap/ticdc/downstreamadapter/sink/mock" "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/config" codeccommon "github.com/pingcap/ticdc/pkg/sink/codec/common" timodel "github.com/pingcap/tidb/pkg/meta/model" + "github.com/pingcap/tidb/pkg/util/chunk" "github.com/stretchr/testify/require" ) -// recordingSink is a minimal sink.Sink implementation that records which DDLs are executed. -// -// It lets unit tests validate consumer-side DDL flushing behavior without requiring a real downstream. -type recordingSink struct { - ddls []string -} - -var _ sink.Sink = (*recordingSink)(nil) - -func (s *recordingSink) SinkType() common.SinkType { return common.MysqlSinkType } -func (s *recordingSink) IsNormal() bool { return true } -func (s *recordingSink) AddDMLEvent(_ *commonEvent.DMLEvent) { -} +func newMockSink(t *testing.T) (*sinkmock.MockSink, *[]string) { + t.Helper() -func (s *recordingSink) FlushDMLBeforeBlock(_ commonEvent.BlockEvent) error { - return nil -} - -func (s *recordingSink) WriteBlockEvent(event commonEvent.BlockEvent) error { - if ddl, ok := event.(*commonEvent.DDLEvent); ok { - s.ddls = append(s.ddls, ddl.Query) - } - return nil -} - -func (s *recordingSink) AddCheckpointTs(_ uint64) { -} + ctrl := gomock.NewController(t) + s := sinkmock.NewMockSink(ctrl) + ddls := make([]string, 0) -func (s *recordingSink) SetTableSchemaStore(_ *commonEvent.TableSchemaStore) { -} + s.EXPECT().AddDMLEvent(gomock.Any()).AnyTimes() + s.EXPECT().WriteBlockEvent(gomock.Any()).DoAndReturn(func(event commonEvent.BlockEvent) error { + if ddl, ok := event.(*commonEvent.DDLEvent); ok { + ddls = append(ddls, ddl.Query) + } + return nil + }).AnyTimes() -func (s *recordingSink) Close() { + return s, &ddls } -func (s *recordingSink) Run(_ context.Context) error { return nil } func TestWriterWrite_executesIndependentCreateTableWithoutWatermark(t *testing.T) { // Scenario: If upstream resolved-ts is held back (e.g. failpoints in integration tests), the consumer @@ -75,7 +58,7 @@ func TestWriterWrite_executesIndependentCreateTableWithoutWatermark(t *testing.T // 1) Enqueue an independent CREATE TABLE DDL with commitTs > watermark. // 2) Call writer.Write and expect the DDL is executed even without watermark catching up. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) w := &writer{ progresses: []*partitionProgress{ {partition: 0, watermark: 0}, @@ -100,7 +83,7 @@ func TestWriterWrite_executesIndependentCreateTableWithoutWatermark(t *testing.T w.Write(ctx, codeccommon.MessageTypeDDL) - require.Equal(t, []string{"CREATE TABLE `test`.`t` (`id` INT PRIMARY KEY)"}, s.ddls) + require.Equal(t, []string{"CREATE TABLE `test`.`t` (`id` INT PRIMARY KEY)"}, *ddls) require.Empty(t, w.ddlList) } @@ -113,7 +96,7 @@ func TestWriterWrite_preservesOrderWhenBlockedDDLNotReady(t *testing.T) { // 2) Call writer.Write and expect nothing executes. // 3) Advance watermark beyond the first DDL and expect both execute in order. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) p := &partitionProgress{partition: 0, watermark: 0} w := &writer{ progresses: []*partitionProgress{p}, @@ -145,7 +128,7 @@ func TestWriterWrite_preservesOrderWhenBlockedDDLNotReady(t *testing.T) { } w.Write(ctx, codeccommon.MessageTypeDDL) - require.Empty(t, s.ddls) + require.Empty(t, *ddls) require.Len(t, w.ddlList, 2) p.watermark = 200 @@ -153,7 +136,7 @@ func TestWriterWrite_preservesOrderWhenBlockedDDLNotReady(t *testing.T) { require.Equal(t, []string{ "ALTER TABLE `test`.`t` ADD COLUMN `c2` INT", "CREATE TABLE `test`.`t2` (`id` INT PRIMARY KEY)", - }, s.ddls) + }, *ddls) require.Empty(t, w.ddlList) } @@ -166,7 +149,7 @@ func TestWriterWrite_doesNotBypassWatermarkForCreateTableLike(t *testing.T) { // 2) Call writer.Write and expect the DDL is NOT executed. // 3) Advance watermark beyond the DDL commitTs and expect the DDL executes. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) p := &partitionProgress{partition: 0, watermark: 0} w := &writer{ progresses: []*partitionProgress{p}, @@ -191,12 +174,12 @@ func TestWriterWrite_doesNotBypassWatermarkForCreateTableLike(t *testing.T) { } w.Write(ctx, codeccommon.MessageTypeDDL) - require.Empty(t, s.ddls) + require.Empty(t, *ddls) require.Len(t, w.ddlList, 1) p.watermark = 200 w.Write(ctx, codeccommon.MessageTypeDDL) - require.Equal(t, []string{"CREATE TABLE `test`.`t2` LIKE `test`.`t1`"}, s.ddls) + require.Equal(t, []string{"CREATE TABLE `test`.`t2` LIKE `test`.`t1`"}, *ddls) require.Empty(t, w.ddlList) } @@ -211,7 +194,7 @@ func TestWriterWrite_handlesOutOfOrderDDLsByCommitTs(t *testing.T) { // 2) Call writer.Write and expect all DDLs with commitTs <= watermark execute (in commit-ts order), // and only the truly "future" DDL remains pending. ctx := context.Background() - s := &recordingSink{} + s, ddls := newMockSink(t) p := &partitionProgress{partition: 0, watermark: 944040962} w := &writer{ progresses: []*partitionProgress{p}, @@ -280,22 +263,11 @@ func TestWriterWrite_handlesOutOfOrderDDLsByCommitTs(t *testing.T) { "ALTER TABLE `common_1`.`add_and_drop_columns` ADD COLUMN `col1` INT NULL, ADD COLUMN `col2` INT NULL, ADD COLUMN `col3` INT NULL", "ALTER TABLE `common_1`.`add_and_drop_columns` DROP COLUMN `col1`, DROP COLUMN `col2`", "CREATE DATABASE `common`", - }, s.ddls) + }, *ddls) require.Len(t, w.ddlList, 1) require.Equal(t, "CREATE TABLE `common_1`.`a` (`a` BIGINT PRIMARY KEY,`b` INT)", w.ddlList[0].Query) } -<<<<<<< HEAD -func TestAppendRow2Group_DoesNotDropCommitTsFallbackBeforeApplied(t *testing.T) { - // Scenario: - // 1) TiCDC writes DML messages to Pulsar in commitTs order. - // 2) Under network partition / changefeed restart, TiCDC may replay older commitTs - // at a later time (commitTs appears to go backwards). - // - // The pulsar-consumer must not drop these "fallback commitTs" events unless they - // have already been flushed to downstream (AppliedWatermark), otherwise replayed - // messages cannot heal missing windows. -======= func TestWriterWrite_sortsOutOfOrderDMLByWatermark(t *testing.T) { ctx := context.Background() ctrl := gomock.NewController(t) @@ -383,13 +355,9 @@ func TestAppendMessageKeepsFallbackDMLAboveGlobalWatermark(t *testing.T) { } func TestOnDDLMarksRoutedCreateTableLikePartitionTable(t *testing.T) { ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) w := &writer{ progresses: []*partitionProgress{ - { - partition: 0, - eventsGroup: make(map[int64]*util.EventsGroup), - }, + {partition: 0, eventsGroup: make(map[int64]*util.EventsGroup)}, }, protocol: config.ProtocolCanalJSON, partitionTableAccessor: codeccommon.NewPartitionTableAccessor(), @@ -423,15 +391,6 @@ func TestOnDDLMarksRoutedCreateTableLikePartitionTable(t *testing.T) { resolved := progress.eventsGroup[1].ResolveInto(150, nil) require.Len(t, resolved, 1) -<<<<<<< HEAD - require.Equal(t, uint64(100), resolved[0].CommitTs) - - // Step 3: once downstream has flushed beyond commitTs=100, replay is safe to ignore. - group.AppliedWatermark = 200 - w.appendRow2Group(newDMLEvent(1, 100), progress) - resolved = group.ResolveInto(150, nil) - require.Empty(t, resolved) -======= require.Equal(t, uint64(100), resolved[0].GetCommitTs()) } @@ -603,5 +562,4 @@ func (m fakePulsarMessage) Index() *uint64 { func (m fakePulsarMessage) BrokerPublishTime() *time.Time { return nil ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } diff --git a/cmd/storage-consumer/consumer.go b/cmd/storage-consumer/consumer.go index b4bc16ecdc..05ec67cf4b 100644 --- a/cmd/storage-consumer/consumer.go +++ b/cmd/storage-consumer/consumer.go @@ -250,13 +250,8 @@ func (c *consumer) appendMessage2Group(message *common.DMLMessage, enableTableAc c.eventsGroup[tableID] = group } if commitTs >= group.HighWatermark { -<<<<<<< HEAD - group.Append(dml, false) - log.Info("DML event append to the group", -======= group.AppendMessage(message) log.Debug("DML event append to the group", ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), zap.Stringer("eventType", message.RowType)) @@ -266,13 +261,8 @@ func (c *consumer) appendMessage2Group(message *common.DMLMessage, enableTableAc log.Warn("DML events fallback, but enableTableAcrossNodes is true, still append it", zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark), zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID), -<<<<<<< HEAD - zap.Stringer("eventType", dml.RowTypes[0])) - group.Append(dml, true) -======= zap.Stringer("eventType", message.RowType)) group.AppendMessage(message) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) return } log.Warn("dml event commit ts fallback, ignore", diff --git a/cmd/util/event_group.go b/cmd/util/event_group.go index 059276d2e8..2391215f03 100644 --- a/cmd/util/event_group.go +++ b/cmd/util/event_group.go @@ -14,11 +14,7 @@ package util import ( -<<<<<<< HEAD - "slices" -======= "math" ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) "sort" "github.com/pingcap/log" @@ -34,13 +30,6 @@ type EventsGroup struct { messages []*codeccommon.DMLMessage HighWatermark uint64 - // AppliedWatermark is the maximum CommitTs that has been successfully flushed - // to the downstream for this group. - // - // It is used to distinguish "safe to ignore" replays (CommitTs <= - // AppliedWatermark) from "still needed" events that arrive late due to sink - // retries / restarts. - AppliedWatermark uint64 } // NewEventsGroup will create new event group. @@ -52,58 +41,6 @@ func NewEventsGroup(partition int32, tableID int64) *EventsGroup { } } -<<<<<<< HEAD -// Append will append an event to event groups. -func (g *EventsGroup) Append(row *commonEvent.DMLEvent, force bool) { - if row.CommitTs > g.HighWatermark { - g.HighWatermark = row.CommitTs - } - - var lastMessage *codeccommon.DMLMessage - if len(g.messages) > 0 { - lastMessage = g.messages[len(g.messages)-1] - } - - if lastMessage == nil || lastMessage.GetCommitTs() <= commitTs { - g.messages = append(g.messages, message) - return - } - - if force { - i := sort.Search(len(g.messages), func(i int) bool { - return g.messages[i].GetCommitTs() > commitTs - }) - g.messages = append(g.messages, nil) - copy(g.messages[i+1:], g.messages[i:]) - g.messages[i] = message - return - } - log.Panic("append event with smaller commit ts", - zap.Int32("partition", g.Partition), zap.Int64("tableID", g.tableID), - zap.Uint64("lastCommitTs", lastMessage.GetCommitTs()), zap.Uint64("commitTs", commitTs)) -} - -// ResolveInto appends all messages with CommitTs <= resolve into dst and removes them from the group. -// ResolveInto copies pointers into dst first, then clears the resolved prefix so Go GC can reclaim -// resolved messages once downstream is done with them. -func (g *EventsGroup) ResolveInto(resolve uint64, dst []*codeccommon.DMLMessage) []*codeccommon.DMLMessage { - i := sort.Search(len(g.messages), func(i int) bool { - return g.messages[i].GetCommitTs() > resolve - }) - if i == 0 { - return dst - } - - // Copy pointers out first so we can safely clear the group's slice without affecting callers. - dst = append(dst, g.messages[:i]...) - clear(g.messages[:i]) - g.messages = g.messages[i:] - if len(g.messages) != 0 { - log.Debug("not all events resolved", - zap.Int32("partition", g.Partition), zap.Int64("tableID", g.tableID), - zap.Int("resolved", i), zap.Int("remained", len(g.events)), - zap.Uint64("resolveTs", resolve), zap.Uint64("firstCommitTs", g.events[0].CommitTs)) -======= // AppendMessage appends a message to event groups. func (g *EventsGroup) AppendMessage(message *codeccommon.DMLMessage) { commitTs := message.GetCommitTs() @@ -171,20 +108,37 @@ func (g *EventsGroup) ResolveInto(resolve uint64, dst []*codeccommon.DMLMessage) zap.Int32("partition", g.Partition), zap.Int64("tableID", g.tableID), zap.Int("resolved", len(resolved)), zap.Int("remained", len(g.messages)), zap.Uint64("resolveTs", resolve), zap.Uint64("firstCommitTs", firstCommitTs)) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } return dst } -<<<<<<< HEAD -// GetAllEvents will get all events. -func (g *EventsGroup) GetAllEvents() []*commonEvent.DMLEvent { - result := g.events - g.events = nil - return result -======= // GetAllMessages gets all messages. func (g *EventsGroup) GetAllMessages() []*codeccommon.DMLMessage { return g.ResolveInto(math.MaxUint64, nil) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) +} + +// AppendOrMergeDMLEvent appends a DML event, or merges it into the previous event +// when both events belong to the same table group and have the same commit-ts. +func AppendOrMergeDMLEvent(events []*commonEvent.DMLEvent, row *commonEvent.DMLEvent) []*commonEvent.DMLEvent { + var lastDMLEvent *commonEvent.DMLEvent + if len(events) > 0 { + lastDMLEvent = events[len(events)-1] + } + + if lastDMLEvent == nil || lastDMLEvent.GetCommitTs() < row.GetCommitTs() { + return append(events, row) + } + + if lastDMLEvent.GetCommitTs() == row.GetCommitTs() { + lastDMLEvent.Rows.Append(row.Rows, 0, row.Rows.NumRows()) + lastDMLEvent.RowTypes = append(lastDMLEvent.RowTypes, row.RowTypes...) + lastDMLEvent.Length += row.Length + lastDMLEvent.PostTxnFlushed = append(lastDMLEvent.PostTxnFlushed, row.PostTxnFlushed...) + return events + } + + log.Panic("append event with smaller commit ts", + zap.Int64("tableID", row.GetTableID()), + zap.Uint64("lastCommitTs", lastDMLEvent.GetCommitTs()), zap.Uint64("commitTs", row.GetCommitTs())) + return events } diff --git a/cmd/util/event_group_test.go b/cmd/util/event_group_test.go index e5bf8655b1..da2cbec9c4 100644 --- a/cmd/util/event_group_test.go +++ b/cmd/util/event_group_test.go @@ -23,44 +23,18 @@ import ( "github.com/stretchr/testify/require" ) -func TestEventsGroupAppendForceMergesExistingCommitTs(t *testing.T) { - // Scenario: - // 1) An upstream transaction (commitTs=100) is split into multiple messages. - // 2) Due to sink retry/restart, a later transaction (commitTs=200) is observed first. - // 3) A "late" fragment of the commitTs=100 transaction arrives afterwards. - // - // The EventsGroup must merge the late fragment into the existing commitTs=100 event, - // instead of turning it into a second commitTs=100 item (which would split one upstream - // transaction into multiple downstream transactions). - group := NewEventsGroup(0, 1) +func newTestDMLMessage(commitTs uint64) *codeccommon.DMLMessage { + return codeccommon.NewDMLMessage(1, "test", "t", commitTs, common.RowTypeInsert, nil) +} - newDMLEvent := func(commitTs uint64) *commonEvent.DMLEvent { - return &commonEvent.DMLEvent{ - CommitTs: commitTs, - RowTypes: []common.RowType{common.RowTypeUpdate}, - Rows: chunk.NewChunkWithCapacity(nil, 0), - Length: 0, - TableInfo: common.NewTableInfo4Decoder("test", &timodel.TableInfo{ - ID: 100, - Name: parser_model.NewCIStr("t"), - Columns: []*timodel.ColumnInfo{ - {Name: parser_model.NewCIStr("a")}, - }, - }), - } +func newTestDMLEvent(commitTs uint64, rowTypes ...common.RowType) *commonEvent.DMLEvent { + return &commonEvent.DMLEvent{ + PhysicalTableID: 1, + CommitTs: commitTs, + Length: int32(len(rowTypes)), + RowTypes: rowTypes, + Rows: chunk.NewChunkWithCapacity(nil, 0), } - - group.Append(newDMLEvent(100), false) - group.Append(newDMLEvent(200), false) - group.Append(newDMLEvent(100), true) - - require.Equal(t, uint64(200), group.HighWatermark) - - var dst []*commonEvent.DMLEvent - dst = group.ResolveInto(150, dst) - require.Len(t, dst, 1) - require.Equal(t, uint64(100), dst[0].CommitTs) - require.Len(t, dst[0].RowTypes, 2) } func TestEventsGroupResolveIntoAppendsAndClearsResolvedMessages(t *testing.T) { @@ -74,18 +48,6 @@ func TestEventsGroupResolveIntoAppendsAndClearsResolvedMessages(t *testing.T) { // 3. Verify (a) returned events are correct, (b) group keeps only the remaining event, // (c) resolved messages in the original backing slice are cleared (nil'd). group := NewEventsGroup(0, 1) -<<<<<<< HEAD - e1 := &commonEvent.DMLEvent{CommitTs: 1} - e2 := &commonEvent.DMLEvent{CommitTs: 2} - e3 := &commonEvent.DMLEvent{CommitTs: 3} - group.Append(e1, false) - group.Append(e2, false) - group.Append(e3, false) - - // Keep a reference to the original slice header so we can validate that ResolveInto clears - // the resolved prefix in-place (this is what prevents GC retention of flushed events). - original := group.events -======= m1 := newTestDMLMessage(1) m2 := newTestDMLMessage(2) m3 := newTestDMLMessage(3) @@ -96,7 +58,6 @@ func TestEventsGroupResolveIntoAppendsAndClearsResolvedMessages(t *testing.T) { // Keep a reference to the original slice header so we can validate that ResolveInto clears // resolved messages in-place (this is what prevents GC retention of flushed events). original := group.messages ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) var dst []*codeccommon.DMLMessage dst = group.ResolveInto(2, dst) @@ -112,28 +73,17 @@ func TestEventsGroupResolveIntoAppendsAndClearsResolvedMessages(t *testing.T) { // doesn't keep flushed events alive via its backing array. require.Same(t, m3, original[0]) require.Nil(t, original[1]) -<<<<<<< HEAD - require.Same(t, e3, original[2]) -======= require.Nil(t, original[2]) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) } func TestEventsGroupResolveIntoNoopWhenNothingResolved(t *testing.T) { // Scenario: resolveTs is behind all buffered events. // Expectation: ResolveInto should be a no-op (dst unchanged, group unchanged). group := NewEventsGroup(0, 1) -<<<<<<< HEAD - e1 := &commonEvent.DMLEvent{CommitTs: 10} - e2 := &commonEvent.DMLEvent{CommitTs: 20} - group.Append(e1, false) - group.Append(e2, false) -======= m1 := newTestDMLMessage(10) m2 := newTestDMLMessage(20) group.AppendMessage(m1) group.AppendMessage(m2) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) original := group.messages dst := make([]*codeccommon.DMLMessage, 0, 1) @@ -153,17 +103,10 @@ func TestEventsGroupResolveIntoClearsAllWhenFullyResolved(t *testing.T) { // Scenario: resolveTs advances beyond all buffered events. // Expectation: group is emptied and all backing-array pointers for resolved events are cleared. group := NewEventsGroup(0, 1) -<<<<<<< HEAD - e1 := &commonEvent.DMLEvent{CommitTs: 1} - e2 := &commonEvent.DMLEvent{CommitTs: 2} - group.Append(e1, false) - group.Append(e2, false) -======= m1 := newTestDMLMessage(1) m2 := newTestDMLMessage(2) group.AppendMessage(m1) group.AppendMessage(m2) ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824)) original := group.messages var dst []*codeccommon.DMLMessage @@ -177,8 +120,6 @@ func TestEventsGroupResolveIntoClearsAllWhenFullyResolved(t *testing.T) { require.Nil(t, original[0]) require.Nil(t, original[1]) } -<<<<<<< HEAD -======= func TestEventsGroupResolveIntoSortsOutOfOrderResolvedMessages(t *testing.T) { group := NewEventsGroup(0, 1) @@ -271,4 +212,3 @@ func TestAppendOrMergeDMLEventAppendsDifferentCommitTs(t *testing.T) { require.Same(t, e1, events[0]) require.Same(t, e2, events[1]) } ->>>>>>> af33cc193 (consumer: sort fallback DML before flush (#5824))