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), 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/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index cc104292ea..673d3e9993 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,77 @@ 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) + } + + 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) + } + + encoder, err := codec.NewEventEncoder(ctx, encoderConfig) + if err != nil { + return errors.Trace(err) + } + encoder.Clean() + + return nil } func New( 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: 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