From a8055c7277c47f7b312c9fc99d857b4f95b16a4d Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Mon, 6 Jun 2022 13:00:41 +0800 Subject: [PATCH 1/5] Support to get non-partitioned topic properties. --- .../admin/impl/PersistentTopicsBase.java | 12 +++++++ .../broker/admin/v2/PersistentTopics.java | 36 +++++++++++++++++++ .../pulsar/broker/admin/AdminApi2Test.java | 13 +++++++ .../broker/admin/PersistentTopicsTest.java | 11 ++++-- .../apache/pulsar/client/admin/Topics.java | 16 +++++++++ .../client/admin/internal/TopicsImpl.java | 26 ++++++++++++++ .../apache/pulsar/admin/cli/CmdTopics.java | 14 ++++++++ 7 files changed, 126 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 7603e0c2e437a..8cbefe2d4e301 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -565,6 +565,18 @@ protected CompletableFuture internalGetPartitionedMeta }); } + protected CompletableFuture> internalGetPropertiesAsync(boolean authoritative) { + return validateTopicOwnershipAsync(topicName, authoritative) + .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA)) + .thenCompose(__ -> pulsar().getBrokerService().getTopicIfExists(topicName.toString())) + .thenApply(opt -> { + if (!opt.isPresent()) { + throw new RestException(Status.NOT_FOUND, getTopicNotFoundErrorMessage(topicName.toString())); + } + return ((PersistentTopic) opt.get()).getManagedLedger().getProperties(); + }); + } + protected CompletableFuture internalCheckTopicExists(TopicName topicName) { return pulsar().getNamespaceService().checkTopicExists(topicName) .thenAccept(exist -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java index 3261e84b505a1..cbc8b9764c7da 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java @@ -918,6 +918,42 @@ public void getPartitionedMetadata( }); } + @GET + @Path("/{tenant}/{namespace}/{topic}/properties") + @ApiOperation(value = "Get non-partitioned topic properties.") + @ApiResponses(value = { + @ApiResponse(code = 307, message = "Current broker doesn't serve the namespace of this topic"), + @ApiResponse(code = 401, message = "Don't have permission to administrate resources on this tenant"), + @ApiResponse(code = 403, message = "Don't have admin permission"), + @ApiResponse(code = 404, message = "Non-Partitioned topic does not exist"), + @ApiResponse(code = 409, message = "Concurrent modification"), + @ApiResponse(code = 412, message = "Non-Partitioned topic name is invalid"), + @ApiResponse(code = 500, message = "Internal server error") + }) + public void getProperties( + @Suspended final AsyncResponse asyncResponse, + @ApiParam(value = "Specify the tenant", required = true) + @PathParam("tenant") String tenant, + @ApiParam(value = "Specify the namespace", required = true) + @PathParam("namespace") String namespace, + @ApiParam(value = "Specify topic name", required = true) + @PathParam("topic") @Encoded String encodedTopic, + @ApiParam(value = "Is authentication required to perform this operation") + @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { + validatePersistentTopicName(tenant, namespace, encodedTopic); + validatePartitionedTopicName(tenant, namespace, encodedTopic); + internalGetPropertiesAsync(authoritative) + .thenAccept(asyncResponse::resume) + .exceptionally(ex -> { + if (!isRedirectException(ex)) { + log.error("[{}] Failed to get non-partitioned topic {} properties", + clientAppId(), topicName, ex); + } + resumeAsyncResponseExceptionally(asyncResponse, ex); + return null; + }); + } + @DELETE @Path("/{tenant}/{namespace}/{topic}/partitions") @ApiOperation(value = "Delete a partitioned topic.", diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index fa1b7ba1657a3..dbc7279ad585a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -858,6 +858,19 @@ public void testPersistentTopicList() throws Exception { assertEquals(topicsInNs.size(), 0); } + @Test + public void testCreateAndGetNonPartitionedTopicProperties() throws Exception { + final String namespace = "prop-xyz/ns2"; + final String topicName = "persistent://" + namespace + "/testGetNonPartitionedTopicProperties"; + admin.namespaces().createNamespace(namespace, 20); + Map properties = new HashMap<>(); + properties.put("key1", "value1"); + admin.topics().createNonPartitionedTopic(topicName, properties); + Map properties2 = admin.topics().getProperties(topicName); + Assert.assertNotNull(properties2); + Assert.assertEquals(properties2.get("key1"), "value1"); + } + @Test public void testNonPersistentTopics() throws Exception { final String namespace = "prop-xyz/ns2"; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java index 356367f410ae6..3693364879bda 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java @@ -453,7 +453,7 @@ public void testNonPartitionedTopics() { @Test public void testCreateNonPartitionedTopic() { - final String topic = "standard-topic-partition-a"; + final String topic = "testCreateNonPartitionedTopic-a"; TopicName topicName = TopicName.get(TopicDomain.persistent.value(), testTenant, testNamespace, topic); AsyncResponse response = mock(AsyncResponse.class); persistentTopics.createNonPartitionedTopic(response, testTenant, testNamespace, topic, true, null); @@ -476,7 +476,7 @@ public void testCreateNonPartitionedTopic() { response = mock(AsyncResponse.class); metaResponse = mock(AsyncResponse.class); metaResponseCaptor = ArgumentCaptor.forClass(PartitionedTopicMetadata.class); - final String topic2 = "standard-topic-partition-b"; + final String topic2 = "testCreateNonPartitionedTopic-b"; TopicName topicName2 = TopicName.get(TopicDomain.persistent.value(), testTenant, testNamespace, topic2); Map topicMetadata = Maps.newHashMap(); topicMetadata.put("key1", "value1"); @@ -488,6 +488,13 @@ public void testCreateNonPartitionedTopic() { testTenant, testNamespace, topic2, true, false); verify(metaResponse, timeout(5000).times(1)).resume(metaResponseCaptor.capture()); Assert.assertNull(metaResponseCaptor.getValue().properties); + metaResponse = mock(AsyncResponse.class); + ArgumentCaptor metaResponseCaptor2 = ArgumentCaptor.forClass(Map.class); + persistentTopics.getProperties(metaResponse, + testTenant, testNamespace, topic2, true); + verify(metaResponse, timeout(5000).times(1)).resume(metaResponseCaptor2.capture()); + Assert.assertNotNull(metaResponseCaptor2.getValue()); + Assert.assertEquals(metaResponseCaptor2.getValue().get("key1"), "value1"); } @Test diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java index d07c67e11ac92..9e94e830844aa 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java @@ -717,6 +717,22 @@ void updatePartitionedTopic(String topic, int numPartitions, boolean updateLocal */ CompletableFuture getPartitionedTopicMetadataAsync(String topic); + /** + * Get properties of a non-partitioned topic. + * @param topic + * Topic name + * @return Non-partitioned topic properties + */ + Map getProperties(String topic) throws PulsarAdminException; + + /** + * Get properties of a non-partitioned topic asynchronously. + * @param topic + * Topic name + * @return a future that can be used to track when the non-partitioned topic properties is returned + */ + CompletableFuture> getPropertiesAsync(String topic); + /** * Delete a partitioned topic and its schemas. *

diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java index d3682a153178f..938b85f1920a3 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java @@ -482,6 +482,32 @@ public void failed(Throwable throwable) { return future; } + @Override + public Map getProperties(String topic) throws PulsarAdminException { + return sync(() -> getPropertiesAsync(topic)); + } + + @Override + public CompletableFuture> getPropertiesAsync(String topic) { + TopicName tn = validateTopic(topic); + WebTarget path = topicPath(tn, "properties"); + final CompletableFuture> future = new CompletableFuture<>(); + asyncGetRequest(path, + new InvocationCallback>() { + + @Override + public void completed(Map response) { + future.complete(response); + } + + @Override + public void failed(Throwable throwable) { + future.completeExceptionally(getApiException(throwable.getCause())); + } + }); + return future; + } + @Override public void deletePartitionedTopic(String topic) throws PulsarAdminException { deletePartitionedTopic(topic, false); diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java index fb78679c5f0b1..15dadb9fcf451 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java @@ -117,6 +117,7 @@ public CmdTopics(Supplier admin) { jcommander.addCommand("create", new CreateNonPartitionedCmd()); jcommander.addCommand("update-partitioned-topic", new UpdatePartitionedCmd()); jcommander.addCommand("get-partitioned-topic-metadata", new GetPartitionedTopicMetadataCmd()); + jcommander.addCommand("get-properties", new GetPropertiesCmd()); jcommander.addCommand("delete-partitioned-topic", new DeletePartitionedCmd()); jcommander.addCommand("peek-messages", new PeekMessages()); @@ -605,6 +606,19 @@ void run() throws Exception { } } + @Parameters(commandDescription = "Get the non-partitioned topic properties.") + private class GetPropertiesCmd extends CliCommand { + + @Parameter(description = "persistent://tenant/namespace/topic", required = true) + private java.util.List params; + + @Override + void run() throws Exception { + String topic = validateTopicName(params); + print(getTopics().getPartitionedTopicMetadata(topic)); + } + } + @Parameters(commandDescription = "Delete a partitioned topic. " + "It will also delete all the partitions of the topic if it exists." + "And the application is not able to connect to the topic(delete then re-create with same name) again " From fc44118f84081ddb5da8cb02cd501e5d8b50a4c4 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Mon, 6 Jun 2022 16:25:01 +0800 Subject: [PATCH 2/5] support partitioned topic to get properties. --- .../admin/impl/PersistentTopicsBase.java | 21 ++++++++++------ .../broker/admin/v2/PersistentTopics.java | 10 ++++---- .../pulsar/broker/admin/AdminApi2Test.java | 24 ++++++++++++------- .../apache/pulsar/client/admin/Topics.java | 8 +++---- .../apache/pulsar/admin/cli/CmdTopics.java | 4 ++-- 5 files changed, 40 insertions(+), 27 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 8cbefe2d4e301..4d3884a609613 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -566,14 +566,21 @@ protected CompletableFuture internalGetPartitionedMeta } protected CompletableFuture> internalGetPropertiesAsync(boolean authoritative) { - return validateTopicOwnershipAsync(topicName, authoritative) - .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA)) - .thenCompose(__ -> pulsar().getBrokerService().getTopicIfExists(topicName.toString())) - .thenApply(opt -> { - if (!opt.isPresent()) { - throw new RestException(Status.NOT_FOUND, getTopicNotFoundErrorMessage(topicName.toString())); + return pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName) + .thenCompose(metadata -> { + if (metadata.partitions == 0 || topicName.isPartitioned()) { + return validateTopicOwnershipAsync(topicName, authoritative) + .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA)) + .thenCompose(__ -> pulsar().getBrokerService().getTopicIfExists(topicName.toString())) + .thenApply(opt -> { + if (!opt.isPresent()) { + throw new RestException(Status.NOT_FOUND, + getTopicNotFoundErrorMessage(topicName.toString())); + } + return ((PersistentTopic) opt.get()).getManagedLedger().getProperties(); + }); } - return ((PersistentTopic) opt.get()).getManagedLedger().getProperties(); + return CompletableFuture.completedFuture(metadata.properties); }); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java index cbc8b9764c7da..daadb9894bc82 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java @@ -920,14 +920,14 @@ public void getPartitionedMetadata( @GET @Path("/{tenant}/{namespace}/{topic}/properties") - @ApiOperation(value = "Get non-partitioned topic properties.") + @ApiOperation(value = "Get topic properties.") @ApiResponses(value = { @ApiResponse(code = 307, message = "Current broker doesn't serve the namespace of this topic"), @ApiResponse(code = 401, message = "Don't have permission to administrate resources on this tenant"), @ApiResponse(code = 403, message = "Don't have admin permission"), - @ApiResponse(code = 404, message = "Non-Partitioned topic does not exist"), + @ApiResponse(code = 404, message = "Topic does not exist"), @ApiResponse(code = 409, message = "Concurrent modification"), - @ApiResponse(code = 412, message = "Non-Partitioned topic name is invalid"), + @ApiResponse(code = 412, message = "Topic name is invalid"), @ApiResponse(code = 500, message = "Internal server error") }) public void getProperties( @@ -941,13 +941,11 @@ public void getProperties( @ApiParam(value = "Is authentication required to perform this operation") @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { validatePersistentTopicName(tenant, namespace, encodedTopic); - validatePartitionedTopicName(tenant, namespace, encodedTopic); internalGetPropertiesAsync(authoritative) .thenAccept(asyncResponse::resume) .exceptionally(ex -> { if (!isRedirectException(ex)) { - log.error("[{}] Failed to get non-partitioned topic {} properties", - clientAppId(), topicName, ex); + log.error("[{}] Failed to get topic {} properties", clientAppId(), topicName, ex); } resumeAsyncResponseExceptionally(asyncResponse, ex); return null; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index dbc7279ad585a..f568882384a28 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -859,16 +859,24 @@ public void testPersistentTopicList() throws Exception { } @Test - public void testCreateAndGetNonPartitionedTopicProperties() throws Exception { + public void testCreateAndGetTopicProperties() throws Exception { final String namespace = "prop-xyz/ns2"; - final String topicName = "persistent://" + namespace + "/testGetNonPartitionedTopicProperties"; + final String nonPartitionedTopicName = "persistent://" + namespace + "/non-partitioned-TopicProperties"; admin.namespaces().createNamespace(namespace, 20); - Map properties = new HashMap<>(); - properties.put("key1", "value1"); - admin.topics().createNonPartitionedTopic(topicName, properties); - Map properties2 = admin.topics().getProperties(topicName); - Assert.assertNotNull(properties2); - Assert.assertEquals(properties2.get("key1"), "value1"); + Map nonPartitionedTopicProperties = new HashMap<>(); + nonPartitionedTopicProperties.put("key1", "value1"); + admin.topics().createNonPartitionedTopic(nonPartitionedTopicName, nonPartitionedTopicProperties); + Map properties11 = admin.topics().getProperties(nonPartitionedTopicName); + Assert.assertNotNull(properties11); + Assert.assertEquals(properties11.get("key1"), "value1"); + + final String partitionedTopicName = "persistent://" + namespace + "/partitioned-TopicProperties"; + Map partitionedTopicProperties = new HashMap<>(); + partitionedTopicProperties.put("key2", "value2"); + admin.topics().createPartitionedTopic(partitionedTopicName, 2, partitionedTopicProperties); + Map properties22 = admin.topics().getProperties(partitionedTopicName); + Assert.assertNotNull(properties22); + Assert.assertEquals(properties22.get("key2"), "value2"); } @Test diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java index 9e94e830844aa..bcacca00e5966 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java @@ -718,18 +718,18 @@ void updatePartitionedTopic(String topic, int numPartitions, boolean updateLocal CompletableFuture getPartitionedTopicMetadataAsync(String topic); /** - * Get properties of a non-partitioned topic. + * Get properties of a topic. * @param topic * Topic name - * @return Non-partitioned topic properties + * @return Topic properties */ Map getProperties(String topic) throws PulsarAdminException; /** - * Get properties of a non-partitioned topic asynchronously. + * Get properties of a topic asynchronously. * @param topic * Topic name - * @return a future that can be used to track when the non-partitioned topic properties is returned + * @return a future that can be used to track when the topic properties is returned */ CompletableFuture> getPropertiesAsync(String topic); diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java index 15dadb9fcf451..bff613f4ce316 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java @@ -606,7 +606,7 @@ void run() throws Exception { } } - @Parameters(commandDescription = "Get the non-partitioned topic properties.") + @Parameters(commandDescription = "Get the topic properties.") private class GetPropertiesCmd extends CliCommand { @Parameter(description = "persistent://tenant/namespace/topic", required = true) @@ -615,7 +615,7 @@ private class GetPropertiesCmd extends CliCommand { @Override void run() throws Exception { String topic = validateTopicName(params); - print(getTopics().getPartitionedTopicMetadata(topic)); + print(getTopics().getProperties(topic)); } } From b16205c5c474eaac263cc0ca2c38ca0a333a243b Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Mon, 6 Jun 2022 16:29:53 +0800 Subject: [PATCH 3/5] add more test. --- .../apache/pulsar/broker/admin/PersistentTopicsTest.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java index 3693364879bda..23fee7236545d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java @@ -523,6 +523,12 @@ public void testCreatePartitionedTopic() { Assert.assertEquals(responseCaptor2.getValue().properties.size(), 1); Assert.assertEquals(responseCaptor2.getValue().properties, topicMetadata); }); + AsyncResponse response3 = mock(AsyncResponse.class); + ArgumentCaptor metaResponseCaptor2 = ArgumentCaptor.forClass(Map.class); + persistentTopics.getProperties(response3, testTenant, testNamespace, topicName2, true); + verify(response3, timeout(5000).times(1)).resume(metaResponseCaptor2.capture()); + Assert.assertNotNull(metaResponseCaptor2.getValue()); + Assert.assertEquals(metaResponseCaptor2.getValue().get("key1"), "value1"); } @Test From 5c24a61afc2a498456bf2569009bdee2ea275b82 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Tue, 7 Jun 2022 12:32:44 +0800 Subject: [PATCH 4/5] add auth. --- .../apache/pulsar/broker/admin/impl/PersistentTopicsBase.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 4d3884a609613..ae0c31b7792e7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -566,11 +566,11 @@ protected CompletableFuture internalGetPartitionedMeta } protected CompletableFuture> internalGetPropertiesAsync(boolean authoritative) { - return pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName) + return validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA) + .thenCompose(__ -> pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName)) .thenCompose(metadata -> { if (metadata.partitions == 0 || topicName.isPartitioned()) { return validateTopicOwnershipAsync(topicName, authoritative) - .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA)) .thenCompose(__ -> pulsar().getBrokerService().getTopicIfExists(topicName.toString())) .thenApply(opt -> { if (!opt.isPresent()) { From 3229ecff1075e726c8e98464d5ae0c756599c498 Mon Sep 17 00:00:00 2001 From: Jiwe Guo Date: Tue, 7 Jun 2022 14:42:41 +0800 Subject: [PATCH 5/5] apply comment. --- .../admin/impl/PersistentTopicsBase.java | 37 ++++++++++++------- 1 file changed, 23 insertions(+), 14 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index ae0c31b7792e7..6387e9f18fad8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -566,24 +566,33 @@ protected CompletableFuture internalGetPartitionedMeta } protected CompletableFuture> internalGetPropertiesAsync(boolean authoritative) { - return validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA) - .thenCompose(__ -> pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName)) - .thenCompose(metadata -> { - if (metadata.partitions == 0 || topicName.isPartitioned()) { - return validateTopicOwnershipAsync(topicName, authoritative) - .thenCompose(__ -> pulsar().getBrokerService().getTopicIfExists(topicName.toString())) - .thenApply(opt -> { - if (!opt.isPresent()) { - throw new RestException(Status.NOT_FOUND, - getTopicNotFoundErrorMessage(topicName.toString())); - } - return ((PersistentTopic) opt.get()).getManagedLedger().getProperties(); - }); + return validateTopicOwnershipAsync(topicName, authoritative) + .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA)) + .thenCompose(__ -> { + if (topicName.isPartitioned()) { + return getPropertiesAsync(); } - return CompletableFuture.completedFuture(metadata.properties); + return pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName) + .thenCompose(metadata -> { + if (metadata.partitions == 0) { + return getPropertiesAsync(); + } + return CompletableFuture.completedFuture(metadata.properties); + }); }); } + private CompletableFuture> getPropertiesAsync() { + return pulsar().getBrokerService().getTopicIfExists(topicName.toString()) + .thenApply(opt -> { + if (!opt.isPresent()) { + throw new RestException(Status.NOT_FOUND, + getTopicNotFoundErrorMessage(topicName.toString())); + } + return ((PersistentTopic) opt.get()).getManagedLedger().getProperties(); + }); + } + protected CompletableFuture internalCheckTopicExists(TopicName topicName) { return pulsar().getNamespaceService().checkTopicExists(topicName) .thenAccept(exist -> {