From 28f329e2328814dfa93634e2aec22a1df8d562de Mon Sep 17 00:00:00 2001 From: wk989898 Date: Mon, 10 Aug 2026 09:22:37 +0000 Subject: [PATCH] init Signed-off-by: wk989898 --- cmd/util/event_group.go | 60 ++++++++++++++---------------------- cmd/util/event_group_test.go | 50 ++++++++++++++++++++++++++++++ 2 files changed, 73 insertions(+), 37 deletions(-) diff --git a/cmd/util/event_group.go b/cmd/util/event_group.go index 2391215f03..d4d324282a 100644 --- a/cmd/util/event_group.go +++ b/cmd/util/event_group.go @@ -29,6 +29,7 @@ type EventsGroup struct { tableID int64 messages []*codeccommon.DMLMessage + outOfOrder bool HighWatermark uint64 } @@ -44,6 +45,9 @@ func NewEventsGroup(partition int32, tableID int64) *EventsGroup { // AppendMessage appends a message to event groups. func (g *EventsGroup) AppendMessage(message *codeccommon.DMLMessage) { commitTs := message.GetCommitTs() + if len(g.messages) > 0 && commitTs < g.messages[len(g.messages)-1].GetCommitTs() { + g.outOfOrder = true + } if commitTs > g.HighWatermark { g.HighWatermark = commitTs } @@ -58,55 +62,37 @@ func (g *EventsGroup) ResolveInto(resolve uint64, dst []*codeccommon.DMLMessage) 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 g.outOfOrder { + sort.SliceStable(g.messages, func(i, j int) bool { + return g.messages[i].GetCommitTs() < g.messages[j].GetCommitTs() + }) } - if outOfOrder { + resolvedCount := sort.Search(len(g.messages), func(i int) bool { + return g.messages[i].GetCommitTs() > resolve + }) + if g.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() - }) + zap.Int("resolved", resolvedCount)) + g.outOfOrder = false + } + if resolvedCount == 0 { + return dst } - dst = append(dst, resolved...) - clear(original[len(remaining):]) - g.messages = remaining + dst = append(dst, g.messages[:resolvedCount]...) + remainingCount := len(g.messages) - resolvedCount + copy(g.messages, g.messages[resolvedCount:]) + clear(g.messages[remainingCount:]) + g.messages = g.messages[:remainingCount] 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.Int("resolved", resolvedCount), zap.Int("remained", len(g.messages)), zap.Uint64("resolveTs", resolve), zap.Uint64("firstCommitTs", firstCommitTs)) } return dst diff --git a/cmd/util/event_group_test.go b/cmd/util/event_group_test.go index da2cbec9c4..a756f297bd 100644 --- a/cmd/util/event_group_test.go +++ b/cmd/util/event_group_test.go @@ -16,11 +16,13 @@ package util import ( "testing" + "github.com/pingcap/log" "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" codeccommon "github.com/pingcap/ticdc/pkg/sink/codec/common" "github.com/pingcap/tidb/pkg/util/chunk" "github.com/stretchr/testify/require" + "go.uber.org/zap/zapcore" ) func newTestDMLMessage(commitTs uint64) *codeccommon.DMLMessage { @@ -182,6 +184,54 @@ func TestEventsGroupGetAllMessagesSortsOutOfOrderMessages(t *testing.T) { require.Empty(t, group.messages) } +func BenchmarkEventsGroupResolveInto(b *testing.B) { + const messageCount = 16 * 1024 + + messages := make([]*codeccommon.DMLMessage, messageCount) + for i := range messages { + messages[i] = newTestDMLMessage(uint64(i + 1)) + } + + benchmarks := []struct { + name string + resolveTs uint64 + outOfOrder bool + }{ + {name: "ordered/noop", resolveTs: 0}, + {name: "ordered/half", resolveTs: messageCount / 2}, + {name: "ordered/all", resolveTs: messageCount}, + {name: "out-of-order/all", resolveTs: messageCount, outOfOrder: true}, + } + + oldLogLevel := log.GetLevel() + log.SetLevel(zapcore.FatalLevel) + b.Cleanup(func() { log.SetLevel(oldLogLevel) }) + + for _, benchmark := range benchmarks { + b.Run(benchmark.name, func(b *testing.B) { + source := messages + if benchmark.outOfOrder { + source = append([]*codeccommon.DMLMessage(nil), messages...) + lastIndex := len(source) - 1 + source[lastIndex-1], source[lastIndex] = source[lastIndex], source[lastIndex-1] + } + group := NewEventsGroup(0, 1) + group.messages = make([]*codeccommon.DMLMessage, 0, messageCount) + dst := make([]*codeccommon.DMLMessage, 0, messageCount) + + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + if len(group.messages) != messageCount { + group.messages = append(group.messages[:0], source...) + group.outOfOrder = benchmark.outOfOrder + } + dst = group.ResolveInto(benchmark.resolveTs, dst[:0]) + } + }) + } +} + func TestAppendOrMergeDMLEventMergesSameCommitTs(t *testing.T) { var flushed []int e1 := newTestDMLEvent(10, common.RowTypeInsert)