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..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 @@ -565,6 +565,34 @@ protected CompletableFuture internalGetPartitionedMeta }); } + protected CompletableFuture> internalGetPropertiesAsync(boolean authoritative) { + return validateTopicOwnershipAsync(topicName, authoritative) + .thenCompose(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_METADATA)) + .thenCompose(__ -> { + if (topicName.isPartitioned()) { + return getPropertiesAsync(); + } + 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 -> { 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..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 @@ -918,6 +918,40 @@ public void getPartitionedMetadata( }); } + @GET + @Path("/{tenant}/{namespace}/{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 = "Topic does not exist"), + @ApiResponse(code = 409, message = "Concurrent modification"), + @ApiResponse(code = 412, message = "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); + internalGetPropertiesAsync(authoritative) + .thenAccept(asyncResponse::resume) + .exceptionally(ex -> { + if (!isRedirectException(ex)) { + log.error("[{}] Failed to get 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..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 @@ -858,6 +858,27 @@ public void testPersistentTopicList() throws Exception { assertEquals(topicsInNs.size(), 0); } + @Test + public void testCreateAndGetTopicProperties() throws Exception { + final String namespace = "prop-xyz/ns2"; + final String nonPartitionedTopicName = "persistent://" + namespace + "/non-partitioned-TopicProperties"; + admin.namespaces().createNamespace(namespace, 20); + 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 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..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 @@ -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 @@ -516,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 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..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 @@ -717,6 +717,22 @@ void updatePartitionedTopic(String topic, int numPartitions, boolean updateLocal */ CompletableFuture getPartitionedTopicMetadataAsync(String topic); + /** + * Get properties of a topic. + * @param topic + * Topic name + * @return Topic properties + */ + Map getProperties(String topic) throws PulsarAdminException; + + /** + * Get properties of a topic asynchronously. + * @param topic + * Topic name + * @return a future that can be used to track when the 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..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 @@ -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 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().getProperties(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 "