From f85b039118f538ec6458d31e1b9ac01d27b66fad Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Thu, 28 Apr 2022 16:50:26 +0800 Subject: [PATCH 1/3] [improve][broker] Make internalGetStats to async Signed-off-by: Zixuan Liu --- .../admin/impl/PersistentTopicsBase.java | 27 ++++++------- .../broker/admin/v1/NonPersistentTopics.java | 23 ----------- .../broker/admin/v1/PersistentTopics.java | 18 +++++++-- .../broker/admin/v2/NonPersistentTopics.java | 38 ------------------- .../broker/admin/v2/PersistentTopics.java | 15 ++++++-- .../broker/admin/PersistentTopicsTest.java | 32 +++++++++++++--- 6 files changed, 66 insertions(+), 87 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 72d7fcfe3cbc1..5f214de35f496 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 @@ -1167,22 +1167,23 @@ private void internalGetSubscriptionsForNonPartitionedTopic(AsyncResponse asyncR }); } - protected TopicStats internalGetStats(boolean authoritative, boolean getPreciseBacklog, - boolean subscriptionBacklogSize, boolean getEarliestTimeInBacklog) { + protected CompletableFuture internalGetStatsAsync(boolean authoritative, + boolean getPreciseBacklog, + boolean subscriptionBacklogSize, + boolean getEarliestTimeInBacklog) { + CompletableFuture future; + if (topicName.isGlobal()) { - validateGlobalNamespaceOwnership(namespaceName); + future = validateGlobalNamespaceOwnershipAsync(namespaceName); + } else { + future = CompletableFuture.completedFuture(null); } - validateTopicOwnership(topicName, authoritative); - validateTopicOperation(topicName, TopicOperation.GET_STATS); - Topic topic = getTopicReference(topicName); - try { - return topic.asyncGetStats(getPreciseBacklog, subscriptionBacklogSize, getEarliestTimeInBacklog).get(); - } catch (InterruptedException | ExecutionException e) { - log.error("[{}] Failed to get stats for {}", clientAppId(), topicName, e); - throw new RestException(Status.INTERNAL_SERVER_ERROR, - (e instanceof ExecutionException) ? e.getCause().getMessage() : e.getMessage()); - } + return future.thenCompose(__ -> validateTopicOwnershipAsync(topicName, authoritative)) + .thenComposeAsync(__ -> validateTopicOperationAsync(topicName, TopicOperation.GET_STATS)) + .thenCompose(__ -> getTopicReferenceAsync(topicName)) + .thenCompose(topic -> topic.asyncGetStats(getPreciseBacklog, subscriptionBacklogSize, + getEarliestTimeInBacklog)); } protected PersistentTopicInternalStats internalGetInternalStats(boolean authoritative, boolean metadata) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java index dc466be86d67b..f63a83045c609 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java @@ -45,7 +45,6 @@ import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.service.Topic; -import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamespaceBundle; @@ -53,7 +52,6 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.NamespaceOperation; -import org.apache.pulsar.common.policies.data.NonPersistentTopicStats; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TopicOperation; @@ -89,27 +87,6 @@ public PartitionedTopicMetadata getPartitionedMetadata(@PathParam("property") St return getPartitionedTopicMetadata(topicName, authoritative, checkAllowAutoCreation); } - @GET - @Path("{property}/{cluster}/{namespace}/{topic}/stats") - @ApiOperation(hidden = true, value = "Get the stats for the topic.") - @ApiResponses(value = { - @ApiResponse(code = 307, message = "Current broker doesn't serve the namespace of this topic"), - @ApiResponse(code = 403, message = "Don't have admin permission"), - @ApiResponse(code = 404, message = "Topic does not exist")}) - public NonPersistentTopicStats getStats(@PathParam("property") String property, - @PathParam("cluster") String cluster, - @PathParam("namespace") String namespace, - @PathParam("topic") @Encoded String encodedTopic, - @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, - @QueryParam("getPreciseBacklog") @DefaultValue("false") - boolean getPreciseBacklog) { - validateTopicName(property, cluster, namespace, encodedTopic); - validateTopicOwnership(topicName, authoritative); - validateTopicOperation(topicName, TopicOperation.GET_STATS); - Topic topic = getTopicReference(topicName); - return ((NonPersistentTopic) topic).getStats(getPreciseBacklog, false, false); - } - @GET @Path("{property}/{cluster}/{namespace}/{topic}/internalStats") @ApiOperation(hidden = true, value = "Get the internal stats for the topic.") diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java index 986ef9bdb2421..5ea49be6d8575 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java @@ -52,7 +52,6 @@ import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.PersistentOfflineTopicStats; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; -import org.apache.pulsar.common.policies.data.TopicStats; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -364,13 +363,24 @@ public void getSubscriptions(@Suspended final AsyncResponse asyncResponse, @Path @ApiResponses(value = { @ApiResponse(code = 307, message = "Current broker doesn't serve the namespace of this topic"), @ApiResponse(code = 403, message = "Don't have admin permission"), - @ApiResponse(code = 404, message = "Topic does not exist") }) - public TopicStats getStats(@PathParam("property") String property, @PathParam("cluster") String cluster, + @ApiResponse(code = 404, message = "Topic does not exist")}) + public void getStats( + @Suspended final AsyncResponse asyncResponse, + @PathParam("property") String property, @PathParam("cluster") String cluster, @PathParam("namespace") String namespace, @PathParam("topic") @Encoded String encodedTopic, @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog) { validateTopicName(property, cluster, namespace, encodedTopic); - return internalGetStats(authoritative, getPreciseBacklog, false, false); + internalGetStatsAsync(authoritative, getPreciseBacklog, false, false) + .thenAccept(asyncResponse::resume) + .exceptionally(ex -> { + // If the exception is not redirect exception we need to log it. + if (!isRedirectException(ex)) { + log.error("[{}] Failed to get stats for {}", clientAppId(), topicName, ex); + } + resumeAsyncResponseExceptionally(asyncResponse, ex); + return null; + }); } @GET diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index 66867b68426a0..245ca766d631c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java @@ -50,13 +50,11 @@ import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.service.Topic; -import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.NamespaceOperation; -import org.apache.pulsar.common.policies.data.NonPersistentTopicStats; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TopicOperation; @@ -102,42 +100,6 @@ public PartitionedTopicMetadata getPartitionedMetadata( return super.getPartitionedMetadata(tenant, namespace, encodedTopic, authoritative, checkAllowAutoCreation); } - @GET - @Path("{tenant}/{namespace}/{topic}/stats") - @ApiOperation(value = "Get the stats for the topic.") - @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 manage resources on this tenant"), - @ApiResponse(code = 403, message = "Don't have admin permission"), - @ApiResponse(code = 404, message = "The tenant/namespace/topic does not exist"), - @ApiResponse(code = 412, message = "Topic name is not valid"), - @ApiResponse(code = 500, message = "Internal server error"), - }) - public NonPersistentTopicStats getStats( - @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, - @ApiParam(value = "If return precise backlog or imprecise backlog") - @QueryParam("getPreciseBacklog") @DefaultValue("false") boolean getPreciseBacklog, - @ApiParam(value = "If return backlog size for each subscription, require locking on ledger so be careful " - + "not to use when there's heavy traffic.") - @QueryParam("subscriptionBacklogSize") @DefaultValue("false") boolean subscriptionBacklogSize, - @ApiParam(value = "If return time of the earliest message in backlog") - @QueryParam("getEarliestTimeInBacklog") @DefaultValue("false") boolean getEarliestTimeInBacklog) { - validateTopicName(tenant, namespace, encodedTopic); - validateTopicOwnership(topicName, authoritative); - validateTopicOperation(topicName, TopicOperation.GET_STATS); - - Topic topic = getTopicReference(topicName); - return ((NonPersistentTopic) topic).getStats(getPreciseBacklog, subscriptionBacklogSize, - getEarliestTimeInBacklog); - } - @GET @Path("{tenant}/{namespace}/{topic}/internalStats") @ApiOperation(value = "Get the internal stats for the topic.") 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 f1c17502870c8..41b7b8ebaa4fc 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 @@ -68,7 +68,6 @@ import org.apache.pulsar.common.policies.data.SchemaCompatibilityStrategy; import org.apache.pulsar.common.policies.data.SubscribeRate; import org.apache.pulsar.common.policies.data.TopicPolicies; -import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.common.policies.data.impl.BacklogQuotaImpl; import org.apache.pulsar.common.policies.data.impl.DispatchRateImpl; import org.slf4j.Logger; @@ -1025,7 +1024,8 @@ public void getSubscriptions( @ApiResponse(code = 412, message = "Topic name is not valid"), @ApiResponse(code = 500, message = "Internal server error"), @ApiResponse(code = 503, message = "Failed to validate global cluster configuration") }) - public TopicStats getStats( + public void getStats( + @Suspended final AsyncResponse asyncResponse, @ApiParam(value = "Specify the tenant", required = true) @PathParam("tenant") String tenant, @ApiParam(value = "Specify the namespace", required = true) @@ -1042,7 +1042,16 @@ public TopicStats getStats( @ApiParam(value = "If return time of the earliest message in backlog") @QueryParam("getEarliestTimeInBacklog") @DefaultValue("false") boolean getEarliestTimeInBacklog) { validateTopicName(tenant, namespace, encodedTopic); - return internalGetStats(authoritative, getPreciseBacklog, subscriptionBacklogSize, getEarliestTimeInBacklog); + internalGetStatsAsync(authoritative, getPreciseBacklog, subscriptionBacklogSize, getEarliestTimeInBacklog) + .thenAccept(asyncResponse::resume) + .exceptionally(ex -> { + // If the exception is not redirect exception we need to log it. + if (!isRedirectException(ex)) { + log.error("[{}] Failed to get stats for {}", clientAppId(), topicName, ex); + } + resumeAsyncResponseExceptionally(asyncResponse, ex); + return null; + }); } @GET 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 bb56f9a59c67f..bc11beb1ed328 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 @@ -276,7 +276,12 @@ public void testCreateSubscriptions() throws Exception{ ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); - TopicStats topicStats = persistentTopics.getStats(testTenant, testNamespace, testLocalTopicName, true, true, false, false); + + response = mock(AsyncResponse.class); + persistentTopics.getStats(response, testTenant, testNamespace, testLocalTopicName, true, true, false, false); + ArgumentCaptor statCaptor = ArgumentCaptor.forClass(TopicStats.class); + verify(response, timeout(5000).times(1)).resume(statCaptor.capture()); + TopicStats topicStats = statCaptor.getValue(); long msgBacklog = topicStats.getSubscriptions().get(SUB_EARLIEST).getMsgBacklog(); System.out.println("Message back log for " + SUB_EARLIEST + " is :" + msgBacklog); Assert.assertEquals(msgBacklog, numberOfMessages); @@ -289,7 +294,12 @@ public void testCreateSubscriptions() throws Exception{ responseCaptor = ArgumentCaptor.forClass(Response.class); verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); - topicStats = persistentTopics.getStats(testTenant, testNamespace, testLocalTopicName, true, true, false, false); + + response = mock(AsyncResponse.class); + persistentTopics.getStats(response, testTenant, testNamespace, testLocalTopicName, true, true, false, false); + statCaptor = ArgumentCaptor.forClass(TopicStats.class); + verify(response, timeout(5000).times(1)).resume(statCaptor.capture()); + topicStats = statCaptor.getValue(); msgBacklog = topicStats.getSubscriptions().get(SUB_LATEST).getMsgBacklog(); System.out.println("Message back log for " + SUB_LATEST + " is :" + msgBacklog); Assert.assertEquals(msgBacklog, 0); @@ -302,7 +312,12 @@ public void testCreateSubscriptions() throws Exception{ responseCaptor = ArgumentCaptor.forClass(Response.class); verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); - topicStats = persistentTopics.getStats(testTenant, testNamespace, testLocalTopicName, true, true, false, false); + + response = mock(AsyncResponse.class); + persistentTopics.getStats(response, testTenant, testNamespace, testLocalTopicName, true, true, false, false); + statCaptor = ArgumentCaptor.forClass(TopicStats.class); + verify(response, timeout(5000).times(1)).resume(statCaptor.capture()); + topicStats = statCaptor.getValue(); msgBacklog = topicStats.getSubscriptions().get(SUB_NONE_MESSAGE_ID).getMsgBacklog(); System.out.println("Message back log for " + SUB_NONE_MESSAGE_ID + " is :" + msgBacklog); Assert.assertEquals(msgBacklog, 0); @@ -315,9 +330,14 @@ public void testCreateSubscriptions() throws Exception{ responseCaptor = ArgumentCaptor.forClass(Response.class); verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); - TopicStats stats = persistentTopics.getStats(testTenant, testNamespace, testLocalTopicName, true, true, false, false); - Assert.assertNotNull(stats.getSubscriptions().get(replicateSubName)); - Assert.assertTrue(stats.getSubscriptions().get(replicateSubName).isReplicated()); + + response = mock(AsyncResponse.class); + persistentTopics.getStats(response,testTenant, testNamespace, testLocalTopicName, true, true, false, false); + statCaptor = ArgumentCaptor.forClass(TopicStats.class); + verify(response, timeout(5000).times(1)).resume(statCaptor.capture()); + topicStats = statCaptor.getValue(); + Assert.assertNotNull(topicStats.getSubscriptions().get(replicateSubName)); + Assert.assertTrue(topicStats.getSubscriptions().get(replicateSubName).isReplicated()); producer.close(); } From e8d7d6ef0e56796fd9331d21df28fece5a8fc59e Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Sat, 7 May 2022 16:37:00 +0800 Subject: [PATCH 2/3] Add test Signed-off-by: Zixuan Liu --- .../broker/admin/AdminTopicApiTest.java | 54 ++++++++++++++++++- 1 file changed, 52 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java index 15dc6906adda2..b13d971d589ac 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.broker.admin; +import java.util.UUID; import lombok.Cleanup; import org.apache.pulsar.client.api.Consumer; @@ -25,31 +26,42 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.common.naming.TopicDomain; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.PersistentTopicStats; +import org.apache.pulsar.common.policies.data.TopicStats; +import org.apache.pulsar.common.policies.data.stats.NonPersistentTopicStatsImpl; +import org.apache.pulsar.common.policies.data.stats.TopicStatsImpl; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeClass; import org.testng.annotations.BeforeMethod; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; import java.util.List; import java.util.concurrent.TimeUnit; import static java.nio.charset.StandardCharsets.UTF_8; +import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertTrue; @Test(groups = "broker-admin") public class AdminTopicApiTest extends ProducerConsumerBase { private static final Logger log = LoggerFactory.getLogger(AdminTopicApiTest.class); @Override - @BeforeMethod + @BeforeClass(alwaysRun = true) protected void setup() throws Exception { super.internalSetup(); super.producerBaseSetup(); } @Override - @AfterMethod(alwaysRun = true) + @BeforeClass(alwaysRun = true) protected void cleanup() throws Exception { super.internalCleanup(); } @@ -97,4 +109,42 @@ public void testPeekMessages() throws Exception { Assert.assertEquals(new String(messages.get(3).getValue(), UTF_8), "value-3"); Assert.assertEquals(new String(messages.get(4).getValue(), UTF_8), "value-4"); } + + @DataProvider + public Object[] getStatsDataProvider() { + return new Object[]{ + // v1 topic + TopicDomain.persistent + "://my-property/test/my-ns/" + UUID.randomUUID(), + TopicDomain.non_persistent+ "://my-property/test/my-ns/" + UUID.randomUUID(), + //v2 topic + TopicDomain.persistent+ "://my-property/my-ns/" + UUID.randomUUID(), + TopicDomain.non_persistent+ "://my-property/my-ns/" + UUID.randomUUID(), + }; + } + + @Test(dataProvider = "getStatsDataProvider") + public void testGetStats(String topic) throws Exception { + admin.topics().createNonPartitionedTopic(topic); + + @Cleanup + PulsarClient newPulsarClient = PulsarClient.builder() + .serviceUrl(lookupUrl.toString()) + .build(); + + final String subscriptionName = "my-sub"; + @Cleanup + Consumer consumer = newPulsarClient.newConsumer() + .topic(topic) + .subscriptionName(subscriptionName) + .subscribe(); + + TopicStats stats = admin.topics().getStats(topic); + assertNotNull(stats); + if (topic.startsWith(TopicDomain.non_persistent.value())) { + assertTrue(stats instanceof NonPersistentTopicStatsImpl); + } else { + assertTrue(stats instanceof TopicStatsImpl); + } + assertTrue(stats.getSubscriptions().containsKey(subscriptionName)); + } } From 43c955c663d58d6f807572df51920603da83296e Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Sat, 7 May 2022 21:50:42 +0800 Subject: [PATCH 3/3] Fix CI Signed-off-by: Zixuan Liu --- .../pulsar/broker/admin/AdminTopicApiTest.java | 18 +++++------------- 1 file changed, 5 insertions(+), 13 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java index b13d971d589ac..a34d289caab74 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTopicApiTest.java @@ -18,37 +18,29 @@ */ package org.apache.pulsar.broker.admin; +import static java.nio.charset.StandardCharsets.UTF_8; +import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertTrue; +import java.util.List; import java.util.UUID; +import java.util.concurrent.TimeUnit; import lombok.Cleanup; - import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerConsumerBase; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.common.naming.TopicDomain; -import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.policies.data.PersistentTopicStats; import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.common.policies.data.stats.NonPersistentTopicStatsImpl; import org.apache.pulsar.common.policies.data.stats.TopicStatsImpl; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; -import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeClass; -import org.testng.annotations.BeforeMethod; import org.testng.annotations.DataProvider; import org.testng.annotations.Test; -import java.util.List; -import java.util.concurrent.TimeUnit; - -import static java.nio.charset.StandardCharsets.UTF_8; -import static org.testng.Assert.assertNotNull; -import static org.testng.Assert.assertNull; -import static org.testng.Assert.assertTrue; - @Test(groups = "broker-admin") public class AdminTopicApiTest extends ProducerConsumerBase { private static final Logger log = LoggerFactory.getLogger(AdminTopicApiTest.class);