Skip to content
Merged
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
92 changes: 57 additions & 35 deletions cmd/kafka-consumer/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,6 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *event.DDLEvent) error {
var (
done = make(chan struct{}, 1)

total int
flushed atomic.Int64
)

Expand All @@ -164,12 +163,16 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *event.DDLEvent) error {
if !ok {
continue
}
before := len(resolvedEvents)
resolvedEvents = g.ResolveInto(commitTs, resolvedEvents)
total += len(resolvedEvents) - before
messages := g.ResolveInto(commitTs, nil)
events := make([]*event.DMLEvent, 0, len(messages))
for _, message := range messages {
events = util.AppendOrMergeDMLEvent(events, message.ToDMLEvent())
}
resolvedEvents = append(resolvedEvents, events...)
}
}

total := len(resolvedEvents)
if total == 0 {
return w.mysqlSink.WriteBlockEvent(ddl)
}
Expand Down Expand Up @@ -265,19 +268,22 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error {
var (
done = make(chan struct{}, 1)

total int
flushed atomic.Int64
)

watermark := w.globalWatermark()
resolvedEvents := make([]*event.DMLEvent, 0)
for _, p := range w.progresses {
for _, group := range p.eventsGroup {
before := len(resolvedEvents)
resolvedEvents = group.ResolveInto(watermark, resolvedEvents)
total += len(resolvedEvents) - before
messages := group.ResolveInto(watermark, nil)
events := make([]*event.DMLEvent, 0, len(messages))
for _, message := range messages {
events = util.AppendOrMergeDMLEvent(events, message.ToDMLEvent())
}
resolvedEvents = append(resolvedEvents, events...)
}
}
total := len(resolvedEvents)
if total == 0 {
return nil
}
Expand Down Expand Up @@ -344,12 +350,12 @@ func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool
ddl := progress.decoder.NextDDLEvent()

if dec, ok := progress.decoder.(*simple.Decoder); ok {
cachedEvents := dec.GetCachedEvents()
for _, row := range cachedEvents {
cachedMessages := dec.GetCachedMessages()
for _, dmlMessage := range cachedMessages {
log.Info("simple protocol cached event resolved, append to the group",
zap.Int64("tableID", row.GetTableID()), zap.Uint64("commitTs", row.CommitTs),
zap.Int64("tableID", dmlMessage.TableID), zap.Uint64("commitTs", dmlMessage.GetCommitTs()),
zap.Int32("partition", partition), zap.Any("offset", offset))
w.appendRow2Group(row, progress, offset)
w.appendMessage2Group(dmlMessage, progress, offset)
}
}

Expand All @@ -373,25 +379,33 @@ func (w *writer) WriteMessage(ctx context.Context, message *kafka.Message) bool
needFlush = true
case common.MessageTypeRow:
var counter int
row := progress.decoder.NextDMLEvent()
if row == nil {
dmlMessage := progress.decoder.NextDMLMessage()
if dmlMessage == nil {
if w.protocol != config.ProtocolSimple {
log.Panic("DML event is nil, it's not expected",
log.Panic("DML message is nil, it's not expected",
zap.Int32("partition", partition), zap.Any("offset", offset))
}
log.Debug("DML event is nil, it's cached", zap.Int32("partition", partition), zap.Any("offset", offset))
log.Debug("DML message is nil, it's cached", zap.Int32("partition", partition), zap.Any("offset", offset))
break
}

w.appendRow2Group(row, progress, offset)
w.appendMessage2Group(dmlMessage, progress, offset)
counter++
for {
_, hasNext = progress.decoder.HasNext()
if !hasNext {
break
}
row = progress.decoder.NextDMLEvent()
w.appendRow2Group(row, progress, offset)
dmlMessage = progress.decoder.NextDMLMessage()
if dmlMessage == nil {
if w.protocol != config.ProtocolSimple {
log.Panic("DML message is nil, it's not expected",
zap.Int32("partition", partition), zap.Any("offset", offset))
}
log.Debug("DML message is nil, it's cached", zap.Int32("partition", partition), zap.Any("offset", offset))
break
}
w.appendMessage2Group(dmlMessage, progress, offset)
counter++
}
// If the message containing only one event exceeds the length limit, CDC will allow it and issue a warning.
Expand Down Expand Up @@ -578,15 +592,22 @@ func (w *writer) checkPartition(row *event.DMLEvent, partition int32, offset kaf
}
}

func (w *writer) appendRow2Group(dml *event.DMLEvent, progress *partitionProgress, offset kafka.Offset) {
w.checkPartition(dml, progress.partition, offset)
func (w *writer) messageWithPartitionCheck(message *common.DMLMessage, partition int32, offset kafka.Offset) *common.DMLMessage {
return common.NewDMLMessage(message.TableID, message.Schema, message.Table, message.GetCommitTs(), message.RowType, func() *event.DMLEvent {
row := message.ToDMLEvent()
w.checkPartition(row, partition, offset)
return row
})
}

func (w *writer) appendMessage2Group(message *common.DMLMessage, progress *partitionProgress, offset kafka.Offset) {
// 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()
)
group := progress.eventsGroup[tableID]
if group == nil {
Expand All @@ -602,21 +623,22 @@ func (w *writer) appendRow2Group(dml *event.DMLEvent, progress *partitionProgres
return
}
if commitTs >= group.HighWatermark {
group.Append(dml, false)
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.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID),
zap.Stringer("eventType", dml.RowTypes[0]))
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", dml.RowTypes[0]))
group.Append(dml, true)
zap.Stringer("eventType", message.RowType))
group.AppendMessage(w.messageWithPartitionCheck(message, progress.partition, offset), true)
return
}
switch w.protocol {
Expand All @@ -631,9 +653,9 @@ func (w *writer) appendRow2Group(dml *event.DMLEvent, progress *partitionProgres
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", dml.RowTypes[0]),
zap.Stringer("eventType", message.RowType),
// zap.Any("columns", row.Columns), zap.Any("preColumns", row.PreColumns),
zap.Any("protocol", w.protocol), zap.Bool("IsPartition", dml.TableInfo.TableName.IsPartition))
zap.Any("protocol", w.protocol))
case config.ProtocolCanalJSON, config.ProtocolOpen, config.ProtocolAvro,
config.ProtocolDebezium, config.ProtocolDebeziumAvro:
// for partition table, these protocols cannot assign physical table id to each dml message,
Expand All @@ -643,18 +665,18 @@ func (w *writer) appendRow2Group(dml *event.DMLEvent, progress *partitionProgres
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", dml.RowTypes[0]), zap.Any("protocol", w.protocol))
group.Append(dml, true)
zap.Stringer("eventType", message.RowType), zap.Any("protocol", w.protocol))
group.AppendMessage(w.messageWithPartitionCheck(message, progress.partition, offset), true)
return
}
log.Warn("DML event fallback row, since less than the group high watermark, ignore it",
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", dml.RowTypes[0]),
zap.Stringer("eventType", message.RowType),
// zap.Any("columns", row.Columns), zap.Any("preColumns", row.PreColumns),
zap.Any("protocol", w.protocol), zap.Bool("IsPartition", dml.TableInfo.TableName.IsPartition))
zap.Any("protocol", w.protocol))
default:
log.Panic("unknown protocol", zap.Any("protocol", w.protocol))
}
Expand Down
12 changes: 6 additions & 6 deletions cmd/kafka-consumer/writer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -311,12 +311,12 @@ func TestOnDDLMarksRoutedCreateTableLikePartitionTableForAvro(t *testing.T) {
}

progress := w.progresses[0]
w.appendRow2Group(newDMLEvent(200), progress, kafka.Offset(10))
w.appendRow2Group(newDMLEvent(100), progress, kafka.Offset(11))
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].CommitTs)
require.Equal(t, uint64(100), resolved[0].GetCommitTs())
}

func TestAppendRow2GroupKeepsDebeziumPartitionTableFallback(t *testing.T) {
Expand Down Expand Up @@ -359,12 +359,12 @@ func TestAppendRow2GroupKeepsDebeziumPartitionTableFallback(t *testing.T) {
}

progress := w.progresses[0]
w.appendRow2Group(newDMLEvent(200), progress, kafka.Offset(10))
w.appendRow2Group(newDMLEvent(100), progress, kafka.Offset(11))
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].CommitTs)
require.Equal(t, uint64(100), resolved[0].GetCommitTs())
})
}
}
2 changes: 1 addition & 1 deletion cmd/pulsar-consumer/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,7 @@ func (c *consumer) readMessage(ctx context.Context) error {
if !needCommit {
continue
}
err := c.pulsarConsumer.AckID(consumerMsg.ID())
err := c.pulsarConsumer.AckIDCumulative(consumerMsg.ID())
if err != nil {
log.Panic("Error ack message", zap.Error(err))
}
Expand Down
61 changes: 34 additions & 27 deletions cmd/pulsar-consumer/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -143,7 +143,6 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *commonEvent.DDLEvent) e
var (
done = make(chan struct{}, 1)

total int
flushed atomic.Int64
)

Expand All @@ -156,12 +155,16 @@ func (w *writer) flushDDLEvent(ctx context.Context, ddl *commonEvent.DDLEvent) e
if !ok {
continue
}
before := len(resolvedEvents)
resolvedEvents = g.ResolveInto(commitTs, resolvedEvents)
total += len(resolvedEvents) - before
messages := g.ResolveInto(commitTs, nil)
events := make([]*commonEvent.DMLEvent, 0, len(messages))
for _, message := range messages {
events = util.AppendOrMergeDMLEvent(events, message.ToDMLEvent())
}
resolvedEvents = append(resolvedEvents, events...)
}
}

total := len(resolvedEvents)
if total == 0 {
return w.mysqlSink.WriteBlockEvent(ddl)
}
Expand Down Expand Up @@ -257,19 +260,22 @@ func (w *writer) flushDMLEventsByWatermark(ctx context.Context) error {
var (
done = make(chan struct{}, 1)

total int
flushed atomic.Int64
)

watermark := w.globalWatermark()
resolvedEvents := make([]*commonEvent.DMLEvent, 0)
for _, p := range w.progresses {
for _, group := range p.eventsGroup {
before := len(resolvedEvents)
resolvedEvents = group.ResolveInto(watermark, resolvedEvents)
total += len(resolvedEvents) - before
messages := group.ResolveInto(watermark, nil)
events := make([]*commonEvent.DMLEvent, 0, len(messages))
for _, message := range messages {
events = util.AppendOrMergeDMLEvent(events, message.ToDMLEvent())
}
resolvedEvents = append(resolvedEvents, events...)
}
}
total := len(resolvedEvents)
if total == 0 {
return nil
}
Expand Down Expand Up @@ -341,12 +347,11 @@ func (w *writer) WriteMessage(ctx context.Context, message pulsar.Message) bool
zap.Any("blockedTables", ddl.GetBlockedTables()))
needFlush = true
case common.MessageTypeRow:
row := progress.decoder.NextDMLEvent()
if row == nil {
log.Panic("DML event is nil, it's not expected")
dmlMessage := progress.decoder.NextDMLMessage()
if dmlMessage == nil {
log.Panic("DML message is nil, it's not expected")
}

w.appendRow2Group(row, progress)
w.appendMessage2Group(dmlMessage, progress)
default:
log.Panic("unknown message type", zap.Any("messageType", messageType))
}
Expand Down Expand Up @@ -486,12 +491,12 @@ func (w *writer) addPartitionTable(schema, table string) {
w.partitionTableAccessor.Add(schema, table)
}

func (w *writer) appendRow2Group(dml *commonEvent.DMLEvent, progress *partitionProgress) {
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()
)
group := progress.eventsGroup[tableID]
if group == nil {
Expand All @@ -506,39 +511,41 @@ func (w *writer) appendRow2Group(dml *commonEvent.DMLEvent, progress *partitionP
return
}
if commitTs >= group.HighWatermark {
group.Append(dml, false)
group.AppendMessage(message, false)
log.Debug("DML event append to the group",
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]))
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", dml.RowTypes[0]))
group.Append(dml, true)
zap.Stringer("eventType", message.RowType))
group.AppendMessage(message, true)
return
}
switch w.protocol {
case config.ProtocolCanalJSON:
// for partition table, the canal-json message cannot assign physical table id to each dml message,
// we cannot distinguish whether it's a real fallback event or not, still append it.
if w.partitionTableAccessor.IsPartitionTable(schema, table) {
isPartitionTable := w.partitionTableAccessor != nil &&
w.partitionTableAccessor.IsPartitionTable(schema, table)
if isPartitionTable {
log.Warn("DML events fallback, but it's canal-json and partition table, 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", dml.RowTypes[0]))
group.Append(dml, true)
zap.Stringer("eventType", message.RowType))
group.AppendMessage(message, true)
return
}
log.Warn("DML event fallback row, since less than the group high watermark, ignore it",
zap.Uint64("commitTs", commitTs), zap.Uint64("highWatermark", group.HighWatermark),
zap.Any("partitionWatermark", progress.watermark), zap.Any("watermark", progress.watermark),
zap.String("schema", schema), zap.String("table", table), zap.Int64("tableID", tableID),
zap.Stringer("eventType", dml.RowTypes[0]),
zap.Any("protocol", w.protocol), zap.Bool("IsPartition", dml.TableInfo.TableName.IsPartition))
zap.Stringer("eventType", message.RowType),
zap.Any("protocol", w.protocol), zap.Bool("IsPartition", isPartitionTable))
default:
log.Panic("unknown protocol", zap.Any("protocol", w.protocol))
}
Expand Down
Loading
Loading