diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index b632eb8850959..3886205d77603 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -49,6 +49,7 @@ import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.UriBuilder; import org.apache.bookkeeper.mledger.LedgerOffloader; +import org.apache.commons.collections4.ListUtils; import org.apache.commons.lang.mutable.MutableObject; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarServerException; @@ -62,6 +63,7 @@ import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.common.api.proto.CommandGetTopicsOfNamespace; import org.apache.pulsar.common.naming.NamedEntity; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceBundleFactory; @@ -162,6 +164,48 @@ protected void internalDeleteNamespace(AsyncResponse asyncResponse, boolean auth } } + protected CompletableFuture> internalGetListOfTopics(Policies policies, + CommandGetTopicsOfNamespace.Mode mode) { + switch (mode) { + case ALL: + return pulsar().getNamespaceService().getListOfPersistentTopics(namespaceName) + .thenCombine(internalGetNonPersistentTopics(policies), + (persistentTopics, nonPersistentTopics) -> + ListUtils.union(persistentTopics, nonPersistentTopics)); + case NON_PERSISTENT: + return internalGetNonPersistentTopics(policies); + case PERSISTENT: + default: + return pulsar().getNamespaceService().getListOfPersistentTopics(namespaceName); + } + } + + protected CompletableFuture> internalGetNonPersistentTopics(Policies policies) { + final List>> futures = Lists.newArrayList(); + final List boundaries = policies.bundles.getBoundaries(); + for (int i = 0; i < boundaries.size() - 1; i++) { + final String bundle = String.format("%s_%s", boundaries.get(i), boundaries.get(i + 1)); + try { + futures.add(pulsar().getAdminClient().topics() + .getListInBundleAsync(namespaceName.toString(), bundle)); + } catch (PulsarServerException e) { + throw new RestException(e); + } + } + return FutureUtil.waitForAll(futures) + .thenApply(__ -> { + final List topics = Lists.newArrayList(); + for (int i = 0; i < futures.size(); i++) { + List topicList = futures.get(i).join(); + if (topicList != null) { + topics.addAll(topicList); + } + } + return topics.stream().filter(name -> !TopicName.get(name).isPersistent()) + .collect(Collectors.toList()); + }); + } + @SuppressWarnings("deprecation") protected void internalDeleteNamespace(AsyncResponse asyncResponse, boolean authoritative) { validateTenantOperation(namespaceName.getTenant(), TenantOperation.DELETE_NAMESPACE); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java index c04773c222db9..e523b0a843ef7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java @@ -147,7 +147,7 @@ public void getTopics(@Suspended AsyncResponse response, validateNamespaceOperationAsync(NamespaceName.get(property, namespace), NamespaceOperation.GET_TOPICS) // Validate that namespace exists, throws 404 if it doesn't exist .thenCompose(__ -> getNamespacePoliciesAsync(namespaceName)) - .thenCompose(__ -> pulsar().getNamespaceService().getListOfTopics(namespaceName, mode)) + .thenCompose(policies -> internalGetListOfTopics(policies, mode)) .thenApply(topics -> filterSystemTopic(topics, includeSystemTopic)) .thenAccept(response::resume) .exceptionally(ex -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java index 89c4722555200..2325eb704a29f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java @@ -121,7 +121,7 @@ public void getTopics(@Suspended AsyncResponse response, validateNamespaceOperationAsync(NamespaceName.get(tenant, namespace), NamespaceOperation.GET_TOPICS) // Validate that namespace exists, throws 404 if it doesn't exist .thenCompose(__ -> getNamespacePoliciesAsync(namespaceName)) - .thenCompose(__ -> pulsar().getNamespaceService().getListOfTopics(namespaceName, mode)) + .thenCompose(policies -> internalGetListOfTopics(policies, mode)) .thenApply(topics -> filterSystemTopic(topics, includeSystemTopic)) .thenAccept(response::resume) .exceptionally(ex -> {