diff --git a/downstreamadapter/sink/kafka/sink.go b/downstreamadapter/sink/kafka/sink.go index 8cd9ed3be3..dae2b4ebd8 100644 --- a/downstreamadapter/sink/kafka/sink.go +++ b/downstreamadapter/sink/kafka/sink.go @@ -102,6 +102,10 @@ func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url. } defer claimCheck.Close() + if _, err = codec.NewEventEncoder(ctx, encoderConfig, claimCheck); 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) @@ -144,11 +148,6 @@ func Verify(ctx context.Context, changefeedID commonType.ChangeFeedID, uri *url. if err != nil { return errors.WrapError(errors.ErrKafkaCreateTopic, err) } - - _, err = codec.NewEventEncoder(ctx, encoderConfig, claimCheck) - if err != nil { - return errors.Trace(err) - } return nil } diff --git a/downstreamadapter/sink/kafka/sink_test.go b/downstreamadapter/sink/kafka/sink_test.go index bfed3536c0..8e93386f34 100644 --- a/downstreamadapter/sink/kafka/sink_test.go +++ b/downstreamadapter/sink/kafka/sink_test.go @@ -37,7 +37,7 @@ import ( const kafkaSinkTestTopic = "mock_topic" -func TestVerifyValidatesEncoderConfigBeforeKafkaConnection(t *testing.T) { +func TestVerifyInvalidEncoderConfig(t *testing.T) { openProtocol := config.ProtocolOpen.String() sinkConfig := &config.SinkConfig{Protocol: &openProtocol} sinkURI, err := url.Parse("kafka://127.0.0.1:1/" + kafkaSinkTestTopic + "?max-batch-size=0") @@ -50,6 +50,21 @@ func TestVerifyValidatesEncoderConfigBeforeKafkaConnection(t *testing.T) { require.ErrorContains(t, err, "invalid max-batch-size 0") } +func TestVerifyEncoderInitialization(t *testing.T) { + avroProtocol := config.ProtocolAvro.String() + schemaRegistry := "http://127.0.0.1:1" + sinkConfig := &config.SinkConfig{ + Protocol: &avroProtocol, + SchemaRegistry: &schemaRegistry, + } + sinkURI, err := url.Parse("kafka://127.0.0.1:1/" + kafkaSinkTestTopic) + require.NoError(t, err) + + changefeedID := common.NewChangefeedID4Test("test", "verify-encoder") + err = Verify(context.Background(), changefeedID, sinkURI, sinkConfig) + require.ErrorContains(t, err, "ErrAvroSchemaAPIError") +} + func newKafkaSinkForTestWithProducers(ctx context.Context, t *testing.T, ctrl *gomock.Controller, diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index 16087c889c..0c58d8933f 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -146,13 +146,27 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool continue } result[meta.Name] = TopicDetail{ - Name: meta.Name, - NumPartitions: int32(len(meta.Partitions)), + Name: meta.Name, + NumPartitions: int32(len(meta.Partitions)), + ReplicationFactor: minReplicationFactor(meta.Partitions), } } return result, nil } +func minReplicationFactor(partitions []*sarama.PartitionMetadata) int16 { + minReplicas := 0 + for i, partition := range partitions { + if partition == nil { + return 0 + } + if i == 0 || len(partition.Replicas) < minReplicas { + minReplicas = len(partition.Replicas) + } + } + return int16(minReplicas) +} + // IsAdminAuthorizationFailed checks whether err is an authorization failure from Kafka admin APIs. func IsAdminAuthorizationFailed(err error) bool { return errors.Is(err, sarama.ErrTopicAuthorizationFailed) || diff --git a/pkg/sink/kafka/admin_test.go b/pkg/sink/kafka/admin_test.go index c2e3f90e37..493e0c1bc3 100644 --- a/pkg/sink/kafka/admin_test.go +++ b/pkg/sink/kafka/admin_test.go @@ -16,6 +16,7 @@ package kafka import ( "testing" + "github.com/IBM/sarama" "github.com/golang/mock/gomock" "github.com/pingcap/ticdc/pkg/common" "github.com/stretchr/testify/require" @@ -62,3 +63,26 @@ func TestAdminClientClose(t *testing.T) { }) } } + +func TestTopicReplicationFactor(t *testing.T) { + ctrl := gomock.NewController(t) + admin := NewMocksaramaClusterAdmin(ctrl) + admin.EXPECT().DescribeTopics([]string{"test-topic"}).Return([]*sarama.TopicMetadata{ + { + Name: "test-topic", + Err: sarama.ErrNoError, + Partitions: []*sarama.PartitionMetadata{ + {Replicas: []int32{1, 2, 3}}, + {Replicas: []int32{1, 2}}, + }, + }, + }, nil) + + client := &saramaAdminClient{ + changefeed: common.NewChangeFeedIDWithName("test", "default"), + admin: admin, + } + topics, err := client.GetTopicsMeta([]string{"test-topic"}, false) + require.NoError(t, err) + require.Equal(t, int16(2), topics["test-topic"].ReplicationFactor) +} diff --git a/pkg/sink/kafka/options.go b/pkg/sink/kafka/options.go index 301471917c..2eb7ec550a 100644 --- a/pkg/sink/kafka/options.go +++ b/pkg/sink/kafka/options.go @@ -633,7 +633,11 @@ func validateRequiredAcks( if options.RequiredAcks != WaitForAll { return nil } - return validateMinInsyncReplicas(ctx, admin, topics, topic, int(options.ReplicationFactor)) + replicationFactor := options.ReplicationFactor + if info, exists := topics[topic]; exists { + replicationFactor = info.ReplicationFactor + } + return validateMinInsyncReplicas(ctx, admin, topics, topic, int(replicationFactor)) } func adjustExistingTopicOption( diff --git a/pkg/sink/kafka/options_test.go b/pkg/sink/kafka/options_test.go index ea5d2e7147..d33a9f21da 100644 --- a/pkg/sink/kafka/options_test.go +++ b/pkg/sink/kafka/options_test.go @@ -83,7 +83,11 @@ func newKafkaAdminFixture(t *testing.T) *kafkaAdminFixture { } func (f *kafkaAdminFixture) addTopic(name string, partitionNum int32) { - f.topics[name] = TopicDetail{Name: name, NumPartitions: partitionNum} + f.topics[name] = TopicDetail{ + Name: name, + NumPartitions: partitionNum, + ReplicationFactor: mockClusterReplicationFactor, + } } func (f *kafkaAdminFixture) getTopicsMeta( @@ -450,8 +454,9 @@ func TestAdjustConfigFallsBackToBrokerMessageMaxBytesWhenTopicConfigMissing(t *t adminClient := adminFixture.admin detail := &TopicDetail{ - Name: topicName, - NumPartitions: 3, + Name: topicName, + NumPartitions: 3, + ReplicationFactor: mockClusterReplicationFactor, } err := adminClient.CreateTopic(detail, false) require.NoError(t, err) @@ -540,9 +545,16 @@ func TestAdjustConfigMinInsyncReplicas(t *testing.T) { err = adjustOptions(ctx, changefeedID, adminClient, options, topicName) require.Nil(t, err) - // topic found, and have `min.insync.replicas`, but set to 2, larger than `replication-factor`. + // Existing topics use their actual replication factor rather than the option + // used only when creating a topic. adminFixture.setMinInsyncReplicas("2") err = adjustOptions(ctx, changefeedID, adminClient, options, defaultMockTopicName) + require.NoError(t, err) + + topicDetail := adminFixture.topics[defaultMockTopicName] + topicDetail.ReplicationFactor = 1 + adminFixture.topics[defaultMockTopicName] = topicDetail + err = adjustOptions(ctx, changefeedID, adminClient, options, defaultMockTopicName) require.Regexp(t, ".*`replication-factor` 1 is smaller than the `min.insync.replicas` 2 of topic.*", errors.Cause(err),