From 4fb02776d890a29bc47afdbe11968b87ac435121 Mon Sep 17 00:00:00 2001 From: nhsmw Date: Tue, 21 Jul 2026 12:48:15 +0800 Subject: [PATCH] This is an automated cherry-pick of #5696 Signed-off-by: ti-chi-bot --- .../sink/topicmanager/kafka_topic_manager.go | 46 ++++- .../topicmanager/kafka_topic_manager_test.go | 192 ++++++++++++++++++ pkg/sink/kafka/admin.go | 6 + 3 files changed, 239 insertions(+), 5 deletions(-) diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 8e92167327..09af904230 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -290,19 +290,29 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( // which means we should create the topic later. topicDetails, err := m.admin.GetTopicsMeta([]string{topicName}, true) if err != nil { + if kafka.IsAdminAuthorizationFailed(err) { + return m.useConfiguredPartitionNum(topicName, err), nil + } return 0, errors.Trace(err) } - if detail, ok := topicDetails[topicName]; ok { - numPartition := detail.NumPartitions - if topicName == m.defaultTopic { - numPartition = m.cfg.PartitionNum + if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { + return numPartition, nil + } + + topicDetails, err = m.admin.GetTopicsMeta([]string{topicName}, false) + if err != nil { + if kafka.IsAdminAuthorizationFailed(err) { + return m.useConfiguredPartitionNum(topicName, err), nil } - m.tryUpdatePartitionsAndLogging(topicName, numPartition) + } else if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { return numPartition, nil } partitionNum, err := m.createTopic(ctx, topicName) if err != nil { + if kafka.IsAdminAuthorizationFailed(err) { + return m.useConfiguredPartitionNum(topicName, err), nil + } return 0, errors.Trace(err) } @@ -314,6 +324,32 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( return partitionNum, nil } +func (m *kafkaTopicManager) tryStoreTopicMeta( + topicName string, topicDetails map[string]kafka.TopicDetail, +) (int32, bool) { + detail, ok := topicDetails[topicName] + if !ok { + return 0, false + } + numPartition := detail.NumPartitions + if topicName == m.defaultTopic { + numPartition = m.cfg.PartitionNum + } + m.tryUpdatePartitionsAndLogging(topicName, numPartition) + return numPartition, true +} + +func (m *kafkaTopicManager) useConfiguredPartitionNum(topicName string, cause error) int32 { + log.Warn("skip Kafka topic creation because topic authorization failed", + zap.String("keyspace", m.changefeedID.Keyspace()), + zap.String("changefeed", m.changefeedID.Name()), + zap.String("topic", topicName), + zap.Int32("partitionNumber", m.cfg.PartitionNum), + zap.Error(cause)) + m.tryUpdatePartitionsAndLogging(topicName, m.cfg.PartitionNum) + return m.cfg.PartitionNum +} + // Close exits the background goroutine. func (m *kafkaTopicManager) Close() { m.cancel() diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index bf02658b24..a1cb89d7ba 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -22,6 +22,58 @@ import ( "github.com/stretchr/testify/require" ) +<<<<<<< HEAD +======= +const kafkaTopicManagerTestTopic = "mock_topic" + +type mockAdminClientWithDeniedDescribe struct { + *kafka.MockClusterAdminClient + createTopicCalled bool + describeCount int +} + +func (m *mockAdminClientWithDeniedDescribe) GetTopicsMeta( + topics []string, + ignoreTopicError bool, +) (map[string]kafka.TopicDetail, error) { + m.describeCount++ + if ignoreTopicError { + return map[string]kafka.TopicDetail{}, nil + } + return nil, sarama.ErrTopicAuthorizationFailed +} + +func (m *mockAdminClientWithDeniedDescribe) CreateTopic( + detail *kafka.TopicDetail, + validateOnly bool, +) error { + m.createTopicCalled = true + return nil +} + +type mockAdminClientWithDeniedCreate struct { + *kafka.MockClusterAdminClient + createTopicCalled bool + describeCount int +} + +func (m *mockAdminClientWithDeniedCreate) GetTopicsMeta( + topics []string, + ignoreTopicError bool, +) (map[string]kafka.TopicDetail, error) { + m.describeCount++ + return map[string]kafka.TopicDetail{}, nil +} + +func (m *mockAdminClientWithDeniedCreate) CreateTopic( + detail *kafka.TopicDetail, + validateOnly bool, +) error { + m.createTopicCalled = true + return sarama.ErrClusterAuthorizationFailed +} + +>>>>>>> 86f31cf6b (kafka: avoid create changefeed failures if can't get a topic from broker (#5696)) func TestCreateTopic(t *testing.T) { t.Parallel() @@ -35,7 +87,56 @@ func TestCreateTopic(t *testing.T) { changefeedID := common.NewChangefeedID4Test("test", "test") ctx := context.Background() +<<<<<<< HEAD manager := newKafkaTopicManager(ctx, kafka.DefaultMockTopicName, changefeedID, adminClient, cfg) +======= + var gotNewTopicDetail *kafka.TopicDetail + var gotNewTopicValidateOnly bool + var gotFailedTopicDetail *kafka.TopicDetail + var gotFailedTopicValidateOnly bool + gomock.InOrder( + adminClient.EXPECT().GetTopicsMeta([]string{kafkaTopicManagerTestTopic}, true).Return( + map[string]kafka.TopicDetail{ + kafkaTopicManagerTestTopic: { + Name: kafkaTopicManagerTestTopic, + NumPartitions: 2, + }, + }, nil), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, true).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().CreateTopic(gomock.Any(), false).DoAndReturn( + func(detail *kafka.TopicDetail, validateOnly bool) error { + gotNewTopicDetail = detail + gotNewTopicValidateOnly = validateOnly + return nil + }), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic"}, false).Return( + map[string]kafka.TopicDetail{ + "new-topic": { + Name: "new-topic", + NumPartitions: 2, + }, + }, nil), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic2"}, true).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic2"}, false).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic-failed"}, true).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().GetTopicsMeta([]string{"new-topic-failed"}, false).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().CreateTopic(gomock.Any(), false).DoAndReturn( + func(detail *kafka.TopicDetail, validateOnly bool) error { + gotFailedTopicDetail = detail + gotFailedTopicValidateOnly = validateOnly + return sarama.ErrInvalidReplicationFactor + }), + ) + + manager := newKafkaTopicManager(ctx, kafkaTopicManagerTestTopic, changefeedID, adminClient, cfg) +>>>>>>> 86f31cf6b (kafka: avoid create changefeed failures if can't get a topic from broker (#5696)) defer manager.Close() partitionNum, err := manager.CreateTopicAndWaitUntilVisible(ctx, kafka.DefaultMockTopicName) require.NoError(t, err) @@ -88,7 +189,40 @@ func TestCreateTopicWithDelay(t *testing.T) { ReplicationFactor: 1, } +<<<<<<< HEAD topic := "new_topic" +======= + topic := "delayed-topic" + gomock.InOrder( + adminClient.EXPECT().GetTopicsMeta([]string{topic}, true).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( + map[string]kafka.TopicDetail{}, nil), + adminClient.EXPECT().CreateTopic(gomock.Any(), false).DoAndReturn( + func(detail *kafka.TopicDetail, validateOnly bool) error { + require.Equal(t, &kafka.TopicDetail{ + Name: topic, + NumPartitions: 2, + ReplicationFactor: 1, + }, detail) + require.False(t, validateOnly) + return nil + }), + adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( + nil, sarama.ErrUnknownTopicOrPartition), + adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( + nil, sarama.ErrUnknownTopicOrPartition), + adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( + map[string]kafka.TopicDetail{ + topic: { + Name: topic, + NumPartitions: 2, + }, + }, nil), + ) + + ctx := context.Background() +>>>>>>> 86f31cf6b (kafka: avoid create changefeed failures if can't get a topic from broker (#5696)) changefeedID := common.NewChangefeedID4Test("test", "test") ctx := context.Background() manager := newKafkaTopicManager(ctx, topic, changefeedID, adminClient, cfg) @@ -101,3 +235,61 @@ func TestCreateTopicWithDelay(t *testing.T) { require.NoError(t, err) require.Equal(t, int32(2), partitionNum) } + +func TestCreateTopicWithTopicDescribeDenied(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := &mockAdminClientWithDeniedDescribe{ + MockClusterAdminClient: kafka.NewMockClusterAdminClient(ctrl), + } + cfg := &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + } + + changefeedID := common.NewChangefeedID4Test("test", "test") + ctx := context.Background() + manager := newKafkaTopicManager(ctx, "precreated-topic", changefeedID, adminClient, cfg) + defer manager.Close() + + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(ctx, "precreated-topic") + require.NoError(t, err) + require.Equal(t, int32(2), partitionNum) + require.False(t, adminClient.createTopicCalled) + require.Equal(t, 2, adminClient.describeCount) + + partitions, ok := manager.topics.Load("precreated-topic") + require.True(t, ok) + require.Equal(t, int32(2), partitions) +} + +func TestCreateTopicWithCreateDenied(t *testing.T) { + t.Parallel() + + ctrl := gomock.NewController(t) + adminClient := &mockAdminClientWithDeniedCreate{ + MockClusterAdminClient: kafka.NewMockClusterAdminClient(ctrl), + } + cfg := &kafka.AutoCreateTopicConfig{ + AutoCreate: true, + PartitionNum: 2, + ReplicationFactor: 1, + } + + changefeedID := common.NewChangefeedID4Test("test", "test") + ctx := context.Background() + manager := newKafkaTopicManager(ctx, "precreated-topic", changefeedID, adminClient, cfg) + defer manager.Close() + + partitionNum, err := manager.CreateTopicAndWaitUntilVisible(ctx, "precreated-topic") + require.NoError(t, err) + require.Equal(t, int32(2), partitionNum) + require.True(t, adminClient.createTopicCalled) + require.Equal(t, 2, adminClient.describeCount) + + partitions, ok := manager.topics.Load("precreated-topic") + require.True(t, ok) + require.Equal(t, int32(2), partitions) +} diff --git a/pkg/sink/kafka/admin.go b/pkg/sink/kafka/admin.go index b8cfd1cfc3..16087c889c 100644 --- a/pkg/sink/kafka/admin.go +++ b/pkg/sink/kafka/admin.go @@ -153,6 +153,12 @@ func (a *saramaAdminClient) GetTopicsMeta(topics []string, ignoreTopicError bool return result, nil } +// IsAdminAuthorizationFailed checks whether err is an authorization failure from Kafka admin APIs. +func IsAdminAuthorizationFailed(err error) bool { + return errors.Is(err, sarama.ErrTopicAuthorizationFailed) || + errors.Is(err, sarama.ErrClusterAuthorizationFailed) +} + func (a *saramaAdminClient) GetTopicsPartitionsNum(topics []string) (map[string]int32, error) { result := make(map[string]int32, len(topics)) for _, topic := range topics {