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))