From 9dc0839d62301f2c63a98fe8b2025bdca8309267 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Fri, 10 Jul 2026 16:45:07 +0800 Subject: [PATCH 1/6] simplify the kafka verify --- downstreamadapter/sink/kafka/sink.go | 76 ++++++++++++++++++++++- downstreamadapter/sink/kafka/sink_test.go | 12 ++++ 2 files changed, 85 insertions(+), 3 deletions(-) diff --git a/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index 2f5f7dd005..10faa4f121 100644 --- a/downstreamadapter/sink/kafka/sink.go +++ b/downstreamadapter/sink/kafka/sink.go @@ -19,12 +19,15 @@ import ( "time" "github.com/pingcap/log" + "github.com/pingcap/ticdc/downstreamadapter/sink/columnselector" + "github.com/pingcap/ticdc/downstreamadapter/sink/eventrouter" "github.com/pingcap/ticdc/downstreamadapter/sink/helper" commonType "github.com/pingcap/ticdc/pkg/common" commonEvent "github.com/pingcap/ticdc/pkg/common/event" "github.com/pingcap/ticdc/pkg/config" "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/metrics" + "github.com/pingcap/ticdc/pkg/sink/codec" "github.com/pingcap/ticdc/pkg/sink/codec/common" "github.com/pingcap/ticdc/pkg/sink/kafka" "github.com/pingcap/ticdc/pkg/util" @@ -68,9 +71,76 @@ func (s *sink) SinkType() commonType.SinkType { } func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url.URL, sinkConfig *config.SinkConfig) error { - comp, _, err := newKafkaSinkComponent(ctx, changefeedID, uri, sinkConfig) - defer comp.close() - return err + protocol, err := helper.GetProtocol(util.GetOrZero(sinkConfig.Protocol)) + if err != nil { + return errors.Trace(err) + } + + topic, err := helper.GetTopic(uri) + if err != nil { + return errors.Trace(err) + } + + options := kafka.NewOptions() + if err = options.Apply(changefeedID, uri, sinkConfig); err != nil { + return errors.WrapError(errors.ErrKafkaInvalidConfig, err) + } + options.Topic = topic + + encoderConfig, err := helper.GetEncoderConfig(changefeedID, uri, protocol, sinkConfig, options.MaxMessageBytes) + if err != nil { + return errors.Trace(err) + } + + isAvroLike := protocol == config.ProtocolAvro || protocol == config.ProtocolDebeziumAvro + if _, err = eventrouter.NewEventRouter(sinkConfig, topic, false, isAvroLike); err != nil { + return errors.Trace(err) + } + + if _, err = columnselector.New(sinkConfig); err != nil { + return errors.Trace(err) + } + + encoder, err := codec.NewEventEncoder(ctx, encoderConfig) + if err != nil { + return errors.Trace(err) + } + encoder.Clean() + + factory, err := kafka.NewSaramaFactory(ctx, options, changefeedID) + if err != nil { + return errors.WrapError(errors.ErrKafkaNewProducer, err) + } + + adminClient, err := factory.AdminClient(ctx) + if err != nil { + return errors.WrapError(errors.ErrKafkaNewProducer, err) + } + defer adminClient.Close() + + topics, err := adminClient.GetTopicsMeta([]string{topic}, false) + if err != nil { + return errors.Trace(err) + } + if _, exists := topics[topic]; exists { + return nil + } + + topicConfig := options.DeriveTopicConfig() + if !topicConfig.AutoCreate { + return errors.ErrKafkaInvalidConfig.GenWithStack("`auto-create-topic` is false, and %s not found", topic) + } + + // the topic is not created, only validate. + err = adminClient.CreateTopic(&kafka.TopicDetail{ + Name: topic, + NumPartitions: topicConfig.PartitionNum, + ReplicationFactor: topicConfig.ReplicationFactor, + }, true) + if err != nil { + return errors.WrapError(errors.ErrKafkaCreateTopic, err) + } + return nil } func New( diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index 7d46286e27..06175797de 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -37,6 +37,18 @@ import ( const kafkaSinkTestTopic = "mock_topic" +func TestVerifyRejectsKafkaUnsupportedProtocolBeforeConnecting(t *testing.T) { + changefeedID := common.NewChangefeedID4Test("test", "verify") + protocol := config.ProtocolCsv.String() + sinkConfig := &config.SinkConfig{Protocol: &protocol} + sinkURI, err := url.Parse("kafka://127.0.0.1:1/test-topic?protocol=csv&kafka-version=2.4.0") + require.NoError(t, err) + + err = Verify(context.Background(), changefeedID, sinkURI, sinkConfig) + require.Error(t, err) + require.Contains(t, err.Error(), "csv") +} + func newKafkaSinkForTestWithProducers(ctx context.Context, t *testing.T, ctrl *gomock.Controller, From 599cdc4f97283531dbb190ac266b9804518dab20 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Fri, 10 Jul 2026 16:51:35 +0800 Subject: [PATCH 2/6] simplify the kafka verify --- downstreamadapter/sink/kafka/helper.go | 15 +++------------ pkg/sink/kafka/factory.go | 4 ---- 2 files changed, 3 insertions(+), 16 deletions(-) diff --git a/downstreamadapter/sink/kafka/helper.go b/downstreamadapter/sink/kafka/helper.go index 4b83310595..6665b4c973 100644 --- a/downstreamadapter/sink/kafka/helper.go +++ b/downstreamadapter/sink/kafka/helper.go @@ -49,11 +49,11 @@ func (c components) close() { } } -func newKafkaSinkComponentWithFactory(ctx context.Context, +func newKafkaSinkComponent( + ctx context.Context, changefeedID commonType.ChangeFeedID, sinkURI *url.URL, sinkConfig *config.SinkConfig, - factoryCreator kafka.FactoryCreator, ) (components, config.Protocol, error) { kafkaComponent := components{} protocol, err := helper.GetProtocol(utils.GetOrZero(sinkConfig.Protocol)) @@ -72,7 +72,7 @@ func newKafkaSinkComponentWithFactory(ctx context.Context, } options.Topic = topic - kafkaComponent.factory, err = factoryCreator(ctx, options, changefeedID) + kafkaComponent.factory, err = kafka.NewSaramaFactory(ctx, options, changefeedID) if err != nil { return kafkaComponent, protocol, errors.WrapError(errors.ErrKafkaNewProducer, err) } @@ -129,12 +129,3 @@ func newKafkaSinkComponentWithFactory(ctx context.Context, } return kafkaComponent, protocol, nil } - -func newKafkaSinkComponent( - ctx context.Context, - changefeedID commonType.ChangeFeedID, - sinkURI *url.URL, - sinkConfig *config.SinkConfig, -) (components, config.Protocol, error) { - return newKafkaSinkComponentWithFactory(ctx, changefeedID, sinkURI, sinkConfig, kafka.NewSaramaFactory) -} diff --git a/pkg/sink/kafka/factory.go b/pkg/sink/kafka/factory.go index 72a458a508..14d83b390c 100644 --- a/pkg/sink/kafka/factory.go +++ b/pkg/sink/kafka/factory.go @@ -16,7 +16,6 @@ package kafka import ( "context" - commonType "github.com/pingcap/ticdc/pkg/common" "github.com/pingcap/ticdc/pkg/sink/codec/common" ) @@ -32,9 +31,6 @@ type Factory interface { MetricsCollector(adminClient ClusterAdminClient) MetricsCollector } -// FactoryCreator defines the type of factory creator. -type FactoryCreator func(context.Context, *options, commonType.ChangeFeedID) (Factory, error) - // SyncProducer is the kafka sync producer type SyncProducer interface { // SendMessage produces a given message, and returns only when it either has From 709fe63bbaacb912243a2c0fdfebd6ca37eaff14 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Fri, 10 Jul 2026 19:25:40 +0800 Subject: [PATCH 3/6] fix consumer panic --- cmd/kafka-consumer/consumer.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/cmd/kafka-consumer/consumer.go b/cmd/kafka-consumer/consumer.go index 143627c124..2bf26f1aee 100644 --- a/cmd/kafka-consumer/consumer.go +++ b/cmd/kafka-consumer/consumer.go @@ -62,7 +62,9 @@ func getPartitionNum(o *option) (int32, error) { } return 0, errors.Trace(err) } - if topicDetail, ok := resp.Topics[topic]; ok { + + topicDetail, ok := resp.Topics[topic] + if ok && topicDetail.Error.Code() == kafka.ErrNoError { numPartitions := int32(len(topicDetail.Partitions)) log.Info("get partition number of topic", zap.String("topic", topic), From cacbca595ae01c75afd5d39ddbb553846efea2c6 Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Mon, 20 Jul 2026 18:36:15 +0800 Subject: [PATCH 4/6] create the encoder at the last --- downstreamadapter/sink/kafka/sink.go | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index daf1236efd..673d3e9993 100644 --- a/downstreamadapter/sink/kafka/sink.go +++ b/downstreamadapter/sink/kafka/sink.go @@ -101,12 +101,6 @@ func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url. return errors.Trace(err) } - encoder, err := codec.NewEventEncoder(ctx, encoderConfig) - if err != nil { - return errors.Trace(err) - } - encoder.Clean() - factory, err := kafka.NewSaramaFactory(ctx, options, changefeedID) if err != nil { return errors.WrapError(errors.ErrKafkaNewProducer, err) @@ -140,6 +134,13 @@ func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url. if err != nil { return errors.WrapError(errors.ErrKafkaCreateTopic, err) } + + encoder, err := codec.NewEventEncoder(ctx, encoderConfig) + if err != nil { + return errors.Trace(err) + } + encoder.Clean() + return nil } From 78d9d90cf2edbf4c2c85f0433193da3ff65ecdfd Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Mon, 20 Jul 2026 18:37:16 +0800 Subject: [PATCH 5/6] create the encoder at the last --- pkg/sink/codec/builder.go | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/pkg/sink/codec/builder.go b/pkg/sink/codec/builder.go index e003710bfb..e6415d06f6 100644 --- a/pkg/sink/codec/builder.go +++ b/pkg/sink/codec/builder.go @@ -20,7 +20,6 @@ import ( "github.com/pingcap/log" "github.com/pingcap/ticdc/pkg/config" "github.com/pingcap/ticdc/pkg/errors" - cerror "github.com/pingcap/ticdc/pkg/errors" "github.com/pingcap/ticdc/pkg/sink/codec/avro" "github.com/pingcap/ticdc/pkg/sink/codec/canal" "github.com/pingcap/ticdc/pkg/sink/codec/common" @@ -62,7 +61,7 @@ func NewEventDecoder( case config.ProtocolAvro: schemaM, err := avro.NewConfluentSchemaManager(ctx, codecConfig.AvroConfluentSchemaRegistry, nil) if err != nil { - return nil, cerror.Trace(err) + return nil, errors.Trace(err) } return avro.NewDecoder(codecConfig, idx, schemaM, topic, upstreamTiDB), nil case config.ProtocolSimple: From 35105960d1f85b506ccd5276c98423021f94f16b Mon Sep 17 00:00:00 2001 From: 3AceShowHand Date: Mon, 20 Jul 2026 19:08:34 +0800 Subject: [PATCH 6/6] remove one unit test --- downstreamadapter/sink/kafka/sink_test.go | 12 ------------ 1 file changed, 12 deletions(-) diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index 66d0e0f861..440806501c 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -37,18 +37,6 @@ import ( const kafkaSinkTestTopic = "mock_topic" -func TestVerifyRejectsKafkaUnsupportedProtocolBeforeConnecting(t *testing.T) { - changefeedID := common.NewChangefeedID4Test("test", "verify") - protocol := config.ProtocolCsv.String() - sinkConfig := &config.SinkConfig{Protocol: &protocol} - sinkURI, err := url.Parse("kafka://127.0.0.1:1/test-topic?protocol=csv&kafka-version=2.4.0") - require.NoError(t, err) - - err = Verify(context.Background(), changefeedID, sinkURI, sinkConfig) - require.Error(t, err) - require.Contains(t, err.Error(), "csv") -} - func newKafkaSinkForTestWithProducers(ctx context.Context, t *testing.T, ctrl *gomock.Controller,