diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java index e019b556b5..83c3794415 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/AdminManager.java @@ -143,13 +143,22 @@ CompletableFuture> describeC admin.topics().getPartitionedTopicMetadataAsync(kopTopic.getFullName()) .whenComplete((metadata, e) -> { if (e != null) { - future.complete(new DescribeConfigsResponse.Config( - ApiError.fromThrowable(e), Collections.emptyList())); + if (e instanceof PulsarAdminException.NotFoundException) { + final ApiError error = new ApiError( + Errors.UNKNOWN_TOPIC_OR_PARTITION, + "Topic " + kopTopic.getOriginalName() + " doesn't exist"); + future.complete(new DescribeConfigsResponse.Config( + error, Collections.emptyList())); + } else { + future.complete(new DescribeConfigsResponse.Config( + ApiError.fromThrowable(e), Collections.emptyList())); + } } else if (metadata.partitions > 0) { future.complete(defaultTopicConfig); } else { - final ApiError error = new ApiError(Errors.UNKNOWN_TOPIC_OR_PARTITION, - "Topic " + kopTopic.getOriginalName() + " doesn't exist"); + final ApiError error = new ApiError(Errors.INVALID_TOPIC_EXCEPTION, + "Topic " + kopTopic.getOriginalName() + + " is non-partitioned"); future.complete(new DescribeConfigsResponse.Config( error, Collections.emptyList())); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 93c9a0fc57..8379ad62b6 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -153,6 +153,7 @@ import org.apache.pulsar.broker.service.Producer; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdmin; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.api.proto.MarkerType; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.naming.NamespaceName; @@ -498,32 +499,12 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, getPartitionedTopicMetadataAsync(fullTopicName) .whenComplete((partitionedTopicMetadata, throwable) -> { if (throwable != null) { - // Failed get partitions. - allTopicMetadata.add( - new TopicMetadata( - Errors.UNKNOWN_TOPIC_OR_PARTITION, - topic, - false, - Collections.emptyList())); - log.warn("[{}] Request {}: Failed to get partitioned pulsar topic {} metadata: {}", - ctx.channel(), metadataHar.getHeader(), - fullTopicName, throwable.getMessage()); - completeOneTopic.run(); - } else { - if (partitionedTopicMetadata.partitions > 0) { - if (log.isDebugEnabled()) { - log.debug("Topic {} has {} partitions", - topic, partitionedTopicMetadata.partitions); - } - addTopicPartition.accept(topic, partitionedTopicMetadata.partitions); - } else { + if (throwable instanceof PulsarAdminException.NotFoundException) { if (kafkaConfig.isAllowAutoTopicCreation() && metadataRequest.allowAutoTopicCreation()) { - if (log.isDebugEnabled()) { - log.debug("[{}] Request {}: Topic {} has single partition, " - + "auto create partitioned topic", - ctx.channel(), metadataHar.getHeader(), topic); - } + log.info("[{}] Request {}: Topic {} doesn't exist, auto create it with {} " + + "partitions", ctx.channel(), metadataHar.getHeader(), + topic, defaultNumPartitions); admin.topics().createPartitionedTopicAsync(fullTopicName, defaultNumPartitions) .whenComplete((ignored, e) -> { if (e == null) { @@ -535,21 +516,47 @@ protected void handleTopicMetadataRequest(KafkaHeaderAndRequest metadataHar, } }); } else { - // NOTE: Currently no matter topic is a non-partitioned topic or topic doesn't - // exist, the queried partitions from broker are both 0. - // See https://github.com/apache/pulsar/issues/8813 for details. log.error("[{}] Request {}: Topic {} doesn't exist and it's not allowed to" + "auto create partitioned topic", ctx.channel(), metadataHar.getHeader(), topic); // not allow to auto create topic, return unknown topic allTopicMetadata.add( - new TopicMetadata( - Errors.UNKNOWN_TOPIC_OR_PARTITION, - topic, - false, - Collections.emptyList())); + new TopicMetadata( + Errors.UNKNOWN_TOPIC_OR_PARTITION, + topic, + false, + Collections.emptyList())); completeOneTopic.run(); } + } else { + // Failed get partitions. + allTopicMetadata.add( + new TopicMetadata( + Errors.UNKNOWN_TOPIC_OR_PARTITION, + topic, + false, + Collections.emptyList())); + log.warn("[{}] Request {}: Failed to get partitioned pulsar topic {} metadata: {}", + ctx.channel(), metadataHar.getHeader(), + fullTopicName, throwable.getMessage()); + completeOneTopic.run(); + } + } else { // the topic already existed + if (partitionedTopicMetadata.partitions > 0) { + if (log.isDebugEnabled()) { + log.debug("Topic {} has {} partitions", + topic, partitionedTopicMetadata.partitions); + } + addTopicPartition.accept(topic, partitionedTopicMetadata.partitions); + } else { + log.error("Topic {} is a non-partitioned topic", topic); + allTopicMetadata.add( + new TopicMetadata( + Errors.INVALID_TOPIC_EXCEPTION, + topic, + false, + Collections.emptyList())); + completeOneTopic.run(); } } }); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java index c191bdee80..0c66cd5b69 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java @@ -15,10 +15,8 @@ import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; -import java.util.HashSet; import java.util.List; import java.util.Set; -import java.util.stream.Collectors; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.common.internals.Topic; import org.apache.pulsar.client.admin.Clusters; @@ -27,8 +25,6 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.admin.PulsarAdminException.ConflictException; import org.apache.pulsar.client.admin.Tenants; -import org.apache.pulsar.client.admin.Topics; -import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; @@ -118,13 +114,13 @@ private static void createKafkaMetadataIfMissing(PulsarAdmin pulsarAdmin, if (!tenants.getTenants().contains(kafkaMetadataTenant)) { log.info("Tenant: {} does not exist, creating it ...", kafkaMetadataTenant); tenants.createTenant(kafkaMetadataTenant, - new TenantInfo(Sets.newHashSet(conf.getSuperUserRoles()), Sets.newHashSet(cluster))); + new TenantInfo(Sets.newHashSet(conf.getSuperUserRoles()), Sets.newHashSet(cluster))); } else { TenantInfo kafkaMetadataTenantInfo = tenants.getTenantInfo(kafkaMetadataTenant); Set allowedClusters = kafkaMetadataTenantInfo.getAllowedClusters(); if (!allowedClusters.contains(cluster)) { log.info("Tenant: {} exists but cluster: {} is not in the allowedClusters list, updating it ...", - kafkaMetadataTenant, cluster); + kafkaMetadataTenant, cluster); allowedClusters.add(cluster); tenants.updateTenant(kafkaMetadataTenant, kafkaMetadataTenantInfo); } @@ -135,7 +131,7 @@ private static void createKafkaMetadataIfMissing(PulsarAdmin pulsarAdmin, Namespaces namespaces = pulsarAdmin.namespaces(); if (!namespaces.getNamespaces(kafkaMetadataTenant).contains(kafkaMetadataNamespace)) { log.info("Namespaces: {} does not exist in tenant: {}, creating it ...", - kafkaMetadataNamespace, kafkaMetadataTenant); + kafkaMetadataNamespace, kafkaMetadataTenant); Set replicationClusters = Sets.newHashSet(cluster); namespaces.createNamespace(kafkaMetadataNamespace, replicationClusters); namespaces.setNamespaceReplicationClusters(kafkaMetadataNamespace, replicationClusters); @@ -143,7 +139,7 @@ private static void createKafkaMetadataIfMissing(PulsarAdmin pulsarAdmin, List replicationClusters = namespaces.getNamespaceReplicationClusters(kafkaMetadataNamespace); if (!replicationClusters.contains(cluster)) { log.info("Namespace: {} exists but cluster: {} is not in the replicationClusters list," - + "updating it ...", kafkaMetadataNamespace, cluster); + + "updating it ...", kafkaMetadataNamespace, cluster); Set newReplicationClusters = Sets.newHashSet(replicationClusters); newReplicationClusters.add(cluster); namespaces.setNamespaceReplicationClusters(kafkaMetadataNamespace, newReplicationClusters); @@ -171,40 +167,7 @@ private static void createKafkaMetadataIfMissing(PulsarAdmin pulsarAdmin, namespaceExists = true; // Check if the offsets topic exists and create it if not - Topics topics = pulsarAdmin.topics(); - PartitionedTopicMetadata topicMetadata = - topics.getPartitionedTopicMetadata(kopTopic.getFullName()); - - Set partitionSet = new HashSet<>(partitionNum); - for (int i = 0; i < partitionNum; i++) { - partitionSet.add(kopTopic.getPartitionName(i)); - } - - if (topicMetadata.partitions <= 0) { - log.info("Kafka group metadata topic {} doesn't exist. Creating it ...", kopTopic.getFullName()); - - topics.createPartitionedTopic( - kopTopic.getFullName(), - partitionNum - ); - - log.info("Successfully created kop metadata topic {} with {} partitions.", - kopTopic.getFullName(), partitionNum); - } else { - // Check to see if the partitions all exist - partitionSet.removeAll( - topics.getList(kafkaMetadataNamespace).stream() - .filter((topic) -> topic.startsWith(kopTopic.getFullName())) - .collect(Collectors.toList()) - ); - - if (!partitionSet.isEmpty()) { - log.info("Identified missing kop metadata topic {} partitions: {}", kopTopic, partitionSet); - for (String offsetPartition : partitionSet) { - topics.createNonPartitionedTopic(offsetPartition); - } - } - } + createTopicIfNotExist(pulsarAdmin, kopTopic.getFullName(), partitionNum); offsetsTopicExists = true; } catch (PulsarAdminException e) { if (e instanceof ConflictException) { @@ -222,4 +185,19 @@ private static void createKafkaMetadataIfMissing(PulsarAdmin pulsarAdmin, kopTopic.getOriginalName(), offsetsTopicExists); } } + + private static void createTopicIfNotExist(final PulsarAdmin admin, + final String topic, + final int numPartitions) throws PulsarAdminException { + try { + admin.topics().createPartitionedTopic(topic, numPartitions); + } catch (PulsarAdminException.ConflictException e) { + log.info("Resources concurrent creating: {}", e.getMessage()); + } + try { + // Ensure all partitions are created + admin.topics().createMissedPartitions(topic); + } catch (PulsarAdminException ignored) { + } + } } diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java index 3337d8d987..0969e4de5d 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java @@ -148,7 +148,7 @@ public void testCreateKafkaMetadataIfMissing() throws Exception { verify(mockTenants, times(1)).updateTenant(eq(conf.getKafkaMetadataTenant()), any(TenantInfo.class)); verify(mockNamespaces, times(2)).setNamespaceReplicationClusters(eq(conf.getKafkaMetadataTenant() + "/" + conf.getKafkaMetadataNamespace()), any(Set.class)); - verify(mockTopics, times(2)).createNonPartitionedTopic(contains(offsetsTopic.getOriginalName())); - verify(mockTopics, times(2)).createNonPartitionedTopic(contains(txnTopic.getOriginalName())); + verify(mockTopics, times(1)).createMissedPartitions(contains(offsetsTopic.getOriginalName())); + verify(mockTopics, times(1)).createMissedPartitions(contains(txnTopic.getOriginalName())); } } diff --git a/pom.xml b/pom.xml index 474e4fc577..67dff40d19 100644 --- a/pom.xml +++ b/pom.xml @@ -48,7 +48,7 @@ 1.18.4 2.22.0 io.streamnative - 2.8.0-rc-202105092228 + 2.8.0-rc-202105182205 1.7.25 3.1.8 1.15.1 diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java index f55c39f86d..3717032e99 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java @@ -64,6 +64,7 @@ import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.config.ConfigResource; +import org.apache.kafka.common.errors.InvalidTopicException; import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.errors.UnknownServerException; import org.apache.kafka.common.errors.UnknownTopicOrPartitionException; @@ -75,6 +76,8 @@ import org.apache.kafka.common.requests.IsolationLevel; import org.apache.kafka.common.requests.ListOffsetRequest; import org.apache.kafka.common.requests.ListOffsetResponse; +import org.apache.kafka.common.requests.MetadataRequest; +import org.apache.kafka.common.requests.MetadataResponse; import org.apache.kafka.common.requests.MetadataResponse.PartitionMetadata; import org.apache.kafka.common.requests.OffsetCommitRequest; import org.apache.kafka.common.requests.RequestHeader; @@ -90,8 +93,9 @@ import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; -import org.testng.annotations.AfterMethod; -import org.testng.annotations.BeforeMethod; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; /** @@ -103,7 +107,12 @@ public class KafkaRequestHandlerTest extends KopProtocolHandlerTestBase { private KafkaRequestHandler handler; private AdminManager adminManager; - @BeforeMethod + @DataProvider(name = "metadataVersions") + public static Object[][] metadataVersions() { + return new Object[][]{ { (short) 0 }, { (short) 1 } }; + } + + @BeforeClass @Override protected void setup() throws Exception { super.internalSetup(); @@ -164,7 +173,7 @@ protected void setup() throws Exception { handler.ctx = mockCtx; } - @AfterMethod + @AfterClass @Override protected void cleanup() throws Exception { adminManager.shutdown(); @@ -326,7 +335,11 @@ private void verifyTopicsDeletedByPulsarAdmin(Map topicToNumPar throws PulsarAdminException { for (Map.Entry entry : topicToNumPartitions.entrySet()) { final String topic = entry.getKey(); - assertEquals(this.admin.topics().getPartitionedTopicMetadata(topic).partitions, 0); + try { + admin.topics().getPartitionedTopicMetadata(topic); + fail("getPartitionedTopicMetadata should fail if topic doesn't exist"); + } catch (PulsarAdminException.NotFoundException expected) { + } } } @@ -454,6 +467,15 @@ public void testDescribeConfigs() throws Exception { assertTrue(e.getCause() instanceof UnknownTopicOrPartitionException); assertTrue(e.getMessage().contains("Topic " + invalidTopic + " doesn't exist")); } + + admin.topics().createNonPartitionedTopic(invalidTopic); + try { + kafkaAdmin.describeConfigs(Collections.singletonList( + new ConfigResource(ConfigResource.Type.TOPIC, invalidTopic))).all().get(); + } catch (ExecutionException e) { + assertTrue(e.getCause() instanceof InvalidTopicException); + assertTrue(e.getMessage().contains("Topic " + invalidTopic + " is non-partitioned")); + } } @Test(timeOut = 10000) @@ -638,4 +660,21 @@ public void testListOffsetsForNotExistedTopic() throws Exception { assertTrue(response.responseData().containsKey(topicPartition)); assertEquals(response.responseData().get(topicPartition).error, Errors.UNKNOWN_TOPIC_OR_PARTITION); } + + @Test(timeOut = 10000, dataProvider = "metadataVersions") + public void testMetadataForNonPartitionedTopic(short version) throws Exception { + final String topic = "testMetadataForNonPartitionedTopic-" + version; + admin.topics().createNonPartitionedTopic(topic); + + final RequestHeader header = new RequestHeader(ApiKeys.METADATA, version, "client", 0); + final MetadataRequest request = new MetadataRequest(Collections.singletonList(topic), false, version); + final CompletableFuture responseFuture = new CompletableFuture<>(); + handler.handleTopicMetadataRequest( + new KafkaHeaderAndRequest(header, request, PulsarByteBufAllocator.DEFAULT.heapBuffer(), null), + responseFuture); + final MetadataResponse response = (MetadataResponse) responseFuture.get(); + assertEquals(response.topicMetadata().size(), 1); + assertEquals(response.errors().size(), 1); + assertEquals(response.errors().get(topic), Errors.INVALID_TOPIC_EXCEPTION); + } }