diff --git a/service-pulsar-query/src/main/java/com/lanyuanxiaoyao/service/pulsar/PulsarQueryApplication.java b/service-pulsar-query/src/main/java/com/lanyuanxiaoyao/service/pulsar/PulsarQueryApplication.java index e7730b8..43bd703 100644 --- a/service-pulsar-query/src/main/java/com/lanyuanxiaoyao/service/pulsar/PulsarQueryApplication.java +++ b/service-pulsar-query/src/main/java/com/lanyuanxiaoyao/service/pulsar/PulsarQueryApplication.java @@ -1,5 +1,6 @@ package com.lanyuanxiaoyao.service.pulsar; +import cn.hutool.core.util.IdUtil; import cn.hutool.core.util.ObjectUtil; import cn.hutool.core.util.StrUtil; import com.eshore.odcp.hudi.connector.exception.PulsarInfoNotFoundException; @@ -14,8 +15,9 @@ import java.util.List; import java.util.Map; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; -import org.apache.pulsar.client.api.MessageId; -import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.*; +import org.apache.pulsar.client.impl.schema.StringSchema; +import org.apache.pulsar.common.policies.data.ConsumerStats; import org.apache.pulsar.common.policies.data.*; import org.eclipse.collections.api.factory.Lists; import org.eclipse.collections.api.list.ImmutableList; @@ -279,23 +281,20 @@ public class PulsarQueryApplication { @Cacheable(value = "exists-topic", sync = true) @GetMapping("exists_topic") - public Boolean existsTopic(@RequestParam("url") String url, @RequestParam("topic") String topic) throws PulsarClientException, PulsarAdminException { - String name = name(url); - logger.info("Detected name: {}", name); - PulsarInfo info = getInfo(name); - try (PulsarAdmin admin = PulsarAdmin.builder() - .serviceHttpUrl(adminUrl(info)) + public Boolean existsTopic(@RequestParam("url") String url, @RequestParam("topic") String topic) { + try (PulsarClient client = PulsarClient.builder() + .serviceUrl(url) .build()) { - List tenants = admin.tenants().getTenants(); - for (String tenant : tenants) { - List namespaces = admin.namespaces().getNamespaces(tenant); - for (String namespace : namespaces) { - List topics = admin.topics().getList(namespace); - if (topics.contains(topic)) { - return true; - } - } + try (Consumer consumer = client.newConsumer(new StringSchema()) + .topic(topic) + .subscriptionMode(SubscriptionMode.NonDurable) + .subscriptionType(SubscriptionType.Exclusive) + .subscriptionName("Pulsar_Exists_Topic_Detect_" + IdUtil.nanoId(8)) + .subscribe()) { + return consumer.isConnected(); } + } catch (PulsarClientException e) { + logger.warn("Pulsar {} not contain topic {}", url, topic); } return false; }