diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go index 3db0938012..be790a71e4 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager.go @@ -270,7 +270,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( if kafka.IsAdminAuthorizationFailed(err) { return m.useConfiguredPartitionNum(topicName, err), nil } - return 0, err + return 0, errors.Trace(err) } if numPartition, ok := m.tryStoreTopicMeta(topicName, topicDetails); ok { return numPartition, nil @@ -291,7 +291,7 @@ func (m *kafkaTopicManager) CreateTopicAndWaitUntilVisible( if kafka.IsAdminAuthorizationFailed(err) { return m.useConfiguredPartitionNum(topicName, err), nil } - return 0, err + return 0, errors.Trace(err) } err = m.waitUntilTopicVisible(ctx, topicName) @@ -328,11 +328,11 @@ func (m *kafkaTopicManager) tryStoreTopicMeta( } func (m *kafkaTopicManager) useConfiguredPartitionNum(topicName string, cause error) int32 { - log.Warn("kafka topic creation skipped due to authorization failure", + 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("partitionNum", m.cfg.PartitionNum), + zap.Int32("partitionNumber", m.cfg.PartitionNum), zap.Error(cause)) m.tryUpdatePartitionsAndLogging(topicName, m.cfg.PartitionNum) return m.cfg.PartitionNum diff --git a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go index 3708d77668..f5e2529437 100644 --- a/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go +++ b/downstreamadapter/sink/topicmanager/kafka_topic_manager_test.go @@ -88,6 +88,7 @@ func TestCreateTopic(t *testing.T) { changefeedID := common.NewChangefeedID4Test("test", "test") ctx := context.Background() + manager := newKafkaTopicManager(ctx, kafka.DefaultMockTopicName, changefeedID, adminClient, cfg) var gotNewTopicDetail *kafka.TopicDetail var gotNewTopicValidateOnly bool var gotFailedTopicDetail *kafka.TopicDetail @@ -129,7 +130,7 @@ func TestCreateTopic(t *testing.T) { func(detail *kafka.TopicDetail, validateOnly bool) error { gotFailedTopicDetail = detail gotFailedTopicValidateOnly = validateOnly - return errors.WrapError(errors.ErrKafkaAdminAPI, sarama.ErrInvalidReplicationFactor, "create-topic", detail.Name) + return sarama.ErrInvalidReplicationFactor }), ) @@ -231,7 +232,7 @@ func TestCreateTopicWaitsUntilVisible(t *testing.T) { ReplicationFactor: 1, } - topic := "delayed-topic" + topic := "new_topic" gomock.InOrder( adminClient.EXPECT().GetTopicsMeta([]string{topic}, true).Return( map[string]kafka.TopicDetail{}, nil), @@ -248,7 +249,9 @@ func TestCreateTopicWaitsUntilVisible(t *testing.T) { return nil }), adminClient.EXPECT().GetTopicsMeta([]string{topic}, false).Return( - map[string]kafka.TopicDetail{}, nil), + 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: { @@ -325,3 +328,61 @@ func TestCreateTopicWithCreateDenied(t *testing.T) { require.True(t, ok) require.Equal(t, int32(2), partitions) } + +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) +}