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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 11 additions & 2 deletions cmd/kafka-consumer/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,10 @@ func (c *consumer) readMessage(ctx context.Context) error {
log.Error("read message failed, just continue to retry", zap.Error(err))
continue
}
needCommit := c.writer.WriteMessage(ctx, msg)
needCommit, err := c.writer.WriteMessage(ctx, msg)
if err != nil {
return err
}
if !needCommit {
continue
}
Expand All @@ -169,7 +172,13 @@ func (c *consumer) readMessage(ctx context.Context) error {
}

// Run the consumer, read data and write to the downstream target.
func (c *consumer) Run(ctx context.Context) error {
func (c *consumer) Run(ctx context.Context) (err error) {
defer func() {
if cleanupErr := c.writer.cleanupEventsGroups(); err == nil && cleanupErr != nil {
err = cleanupErr
}
}()

g, ctx := errgroup.WithContext(ctx)
g.Go(func() error {
return c.writer.run(ctx)
Expand Down
73 changes: 52 additions & 21 deletions cmd/kafka-consumer/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,22 @@ func (w *writer) run(ctx context.Context) error {
return w.mysqlSink.Run(ctx)
}

func (w *writer) cleanupEventsGroups() error {
var cleanupErr error
for _, progress := range w.progresses {
for _, group := range progress.eventsGroup {
if err := group.Cleanup(); err != nil {
log.Warn("cleanup events group spill file failed",
zap.Int32("partition", progress.partition), zap.Error(err))
if cleanupErr == nil {
cleanupErr = err
}
}
}
}
return cleanupErr
}

func (w *writer) flushDDLEvent(ctx context.Context, ddl *event.DDLEvent) error {
var (
done = make(chan struct{}, 1)
Expand All @@ -163,7 +179,10 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *event.DDLEvent) error {
if !ok {
continue
}
messages := g.ResolveInto(commitTs, nil)
messages, err := g.ResolveInto(commitTs, nil)
if err != nil {
return err
}
events := make([]*event.DMLEvent, 0, len(messages))
for _, message := range messages {
events = util.AppendOrMergeDMLEvent(events, message.ToDMLEvent())
Expand Down Expand Up @@ -275,7 +294,10 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error {
resolvedEvents := make([]*event.DMLEvent, 0)
for _, p := range w.progresses {
for _, group := range p.eventsGroup {
messages := group.ResolveInto(watermark, nil)
messages, err := group.ResolveInto(watermark, nil)
if err != nil {
return err
}
events := make([]*event.DMLEvent, 0, len(messages))
for _, message := range messages {
events = util.AppendOrMergeDMLEvent(events, message.ToDMLEvent())
Expand Down Expand Up @@ -320,7 +342,7 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error {
// WriteMessage is to decode kafka message to event.
// return true if the message is flushed to the downstream.
// return error if flush messages failed.
func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool {
func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) (bool, error) {
var (
partition = message.TopicPartition.Partition
offset = message.TopicPartition.Offset
Expand Down Expand Up @@ -355,19 +377,21 @@ func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool
log.Info("simple protocol cached event resolved, append to the group",
zap.Int64("tableID", dmlMessage.TableID), zap.Uint64("commitTs", dmlMessage.GetCommitTs()),
zap.Int32("partition", partition), zap.Any("offset", offset))
w.appendMessage2Group(dmlMessage, progress, offset)
if err := w.appendMessage2Group(dmlMessage, progress, offset); err != nil {
return false, err
}
}
}

w.onDDL(ddl)
// DDL is broadcast to all partitions, but only handle the DDL from partition-0.
if partition != 0 {
return false
return false, nil
}

// the Query maybe empty if using simple protocol, it's comes from `bootstrap` event, no need to handle it.
if ddl.Query == "" {
return false
return false, nil
}
w.appendDDL(ddl)
log.Info("DDL event received",
Expand All @@ -389,7 +413,9 @@ func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool
break
}

w.appendMessage2Group(dmlMessage, progress, offset)
if err := w.appendMessage2Group(dmlMessage, progress, offset); err != nil {
return false, err
}
counter++
for {
_, hasNext = progress.decoder.HasNext()
Expand All @@ -405,7 +431,9 @@ func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool
log.Debug("DML message is nil, it's cached", zap.Int32("partition", partition), zap.Any("offset", offset))
break
}
w.appendMessage2Group(dmlMessage, progress, offset)
if err := w.appendMessage2Group(dmlMessage, progress, offset); err != nil {
return false, err
}
counter++
}
// If the message containing only one event exceeds the length limit, CDC will allow it and issue a warning.
Expand All @@ -427,11 +455,11 @@ func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool
if needFlush {
return w.Write(ctx, messageType)
}
return false
return false, nil
}

// Write will synchronously write data downstream
func (w *writer) Write(ctx context.Context, messageType common.MessageType) bool {
func (w *writer) Write(ctx context.Context, messageType common.MessageType) (bool, error) {
// DDL events can be received out of commit-ts order (e.g. due to protocol-level broadcasting and
// buffering differences between DDL kinds). We must execute DDLs in commit-ts order; otherwise a
// "future" DDL that is not yet eligible (commitTs > watermark) can block executing earlier DDLs
Expand Down Expand Up @@ -479,16 +507,15 @@ func (w *writer) Write(ctx context.Context, messageType common.MessageType) bool
break
}
if err := w.flushDDLEvent(ctx, todoDDL); err != nil {
log.Panic("write DDL event failed", zap.Error(err),
zap.String("DDL", todoDDL.Query), zap.Uint64("commitTs", todoDDL.GetCommitTs()))
return false, err
}
}

if messageType == common.MessageTypeResolved {
// since watermark is broadcast to all partitions, so that each partition can flush events individually.
err := w.flushDMLEventsByWatermark(ctx)
if err != nil {
log.Panic("flush dml events by the watermark failed", zap.Error(err))
return false, err
}
}

Expand All @@ -498,9 +525,9 @@ func (w *writer) Write(ctx context.Context, messageType common.MessageType) bool
log.Info("some DDL events will be flushed in the future",
zap.Uint64("watermark", watermark),
zap.Int("length", len(w.ddlList)))
return false
return false, nil
}
return true
return true, nil
}

func (w *writer) onDDL(ddl *event.DDLEvent) {
Expand Down Expand Up @@ -600,7 +627,7 @@ func (w *writer) messageWithPartitionCheck(message *common.DMLMessage, partition
})
}

func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *partitionProgress, offset kafka.Offset) {
func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *partitionProgress, offset kafka.Offset) error {
// 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 (
Expand All @@ -620,16 +647,19 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti
zap.String("schema", schema), zap.String("table", table),
zap.Stringer("eventType", message.RowType),
zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes))
return
return nil
}

group := progress.eventsGroup[tableID]
if group == nil {
group = util.NewEventsGroup(progress.partition, tableID)
progress.eventsGroup[tableID] = group
}
message = w.messageWithPartitionCheck(message, progress.partition, offset)
group.AppendMessage(message)
if err := group.AppendMessageWithPostRestore(message, func(message *common.DMLMessage) *common.DMLMessage {
return w.messageWithPartitionCheck(message, progress.partition, offset)
}); err != nil {
return err
}
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),
Expand All @@ -639,15 +669,15 @@ func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *parti
zap.String("schema", schema), zap.String("table", table),
zap.Stringer("eventType", message.RowType),
zap.Any("protocol", w.protocol), zap.Bool("enableTableAcrossNodes", w.enableTableAcrossNodes))
return
return nil
}
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
return nil
}
log.Warn("DML event commit ts fallback, append it and sort before flush",
zap.Int32("partition", progress.partition), zap.Any("offset", offset),
Expand All @@ -656,6 +686,7 @@ 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))
return nil
}

func openDB(ctx context.Context, dsn string) (*sql.DB, error) {
Expand Down
19 changes: 14 additions & 5 deletions cmd/kafka-consumer/writer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -273,8 +273,10 @@ func TestWriterWrite_sortsOutOfOrderDMLByWatermark(t *testing.T) {
ctrl := gomock.NewController(t)
s := sinkmock.NewMockSink(ctrl)
flushedCommitTs := make([]uint64, 0)
flushedRowTypeCounts := make([]int, 0)
s.EXPECT().AddDMLEvent(gomock.Any()).Do(func(event *commonEvent.DMLEvent) {
flushedCommitTs = append(flushedCommitTs, event.GetCommitTs())
flushedRowTypeCounts = append(flushedRowTypeCounts, len(event.RowTypes))
event.PostFlush()
}).Times(2)

Expand All @@ -299,8 +301,11 @@ func TestWriterWrite_sortsOutOfOrderDMLByWatermark(t *testing.T) {
w.appendMessage2Group(newDMLMessageForWriterTest(20), p, kafka.Offset(3))

p.watermark = 20
require.True(t, w.Write(ctx, codeccommon.MessageTypeResolved))
needCommit, err := w.Write(ctx, codeccommon.MessageTypeResolved)
require.NoError(t, err)
require.True(t, needCommit)
require.Equal(t, []uint64{10, 20}, flushedCommitTs)
require.Equal(t, []int{1, 2}, flushedRowTypeCounts)
}

func TestWriteMessageIgnoresFallbackDMLBelowGlobalWatermark(t *testing.T) {
Expand All @@ -323,9 +328,10 @@ func TestWriteMessageIgnoresFallbackDMLBelowGlobalWatermark(t *testing.T) {
maxMessageBytes: 1,
}

needCommit := w.WriteMessage(ctx, &kafka.Message{
needCommit, err := w.WriteMessage(ctx, &kafka.Message{
TopicPartition: kafka.TopicPartition{Partition: 0, Offset: kafka.Offset(10)},
})
require.NoError(t, err)

require.False(t, needCommit)
require.Nil(t, progress.eventsGroup[1])
Expand Down Expand Up @@ -353,7 +359,8 @@ func TestAppendMessageKeepsFallbackDMLAboveGlobalWatermark(t *testing.T) {
w.appendMessage2Group(newDMLMessageForWriterTest(10), progress, kafka.Offset(10))

require.NotNil(t, progress.eventsGroup[1])
resolved := progress.eventsGroup[1].ResolveInto(20, nil)
resolved, err := progress.eventsGroup[1].ResolveInto(20, nil)
require.NoError(t, err)
require.Len(t, resolved, 1)
require.Equal(t, uint64(10), resolved[0].GetCommitTs())
}
Expand Down Expand Up @@ -404,7 +411,8 @@ func TestOnDDLMarksRoutedCreateTableLikePartitionTableForAvro(t *testing.T) {
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)
resolved, err := progress.eventsGroup[1].ResolveInto(150, nil)
require.NoError(t, err)
require.Len(t, resolved, 1)
require.Equal(t, uint64(100), resolved[0].GetCommitTs())
}
Expand Down Expand Up @@ -452,7 +460,8 @@ func TestAppendRow2GroupKeepsDebeziumPartitionTableFallback(t *testing.T) {
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)
resolved, err := progress.eventsGroup[1].ResolveInto(150, nil)
require.NoError(t, err)
require.Len(t, resolved, 1)
require.Equal(t, uint64(100), resolved[0].GetCommitTs())
})
Expand Down
13 changes: 11 additions & 2 deletions cmd/pulsar-consumer/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,10 @@ func (c *consumer) readMessage(ctx context.Context) error {
return errors.Trace(ctx.Err())
case consumerMsg := <-msgChan:
log.Debug("Received message", zap.Stringer("msgId", consumerMsg.ID()), zap.ByteString("content", consumerMsg.Payload()))
needCommit := c.writer.WriteMessage(ctx, consumerMsg)
needCommit, writeErr := c.writer.WriteMessage(ctx, consumerMsg)
if writeErr != nil {
return writeErr
}
if !needCommit {
continue
}
Expand All @@ -123,7 +126,13 @@ func (c *consumer) readMessage(ctx context.Context) error {
}

// Run the consumer, read data and write to the downstream target.
func (c *consumer) Run(ctx context.Context) error {
func (c *consumer) Run(ctx context.Context) (err error) {
defer func() {
if cleanupErr := c.writer.cleanupEventsGroups(); err == nil && cleanupErr != nil {
err = cleanupErr
}
}()

g, ctx := errgroup.WithContext(ctx)
g.Go(func() error {
return c.writer.run(ctx)
Expand Down
Loading
Loading