☰
Apache Pulsar Admin 接口完全指南:pulsar-admin CLI、REST API 与 Java Admin API 实战
2026/9/25 3:34:44 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

导读

Apache Pulsar 的 Admin 接口(admin interface)是管理 Pulsar 实例中所有核心实体(tenant、namespace、topic 等)的统一入口,它包含pulsar-admin命令行工具、Broker 暴露的 REST API 以及 Java Admin API 三种交互方式。本文以 Pulsar 仓库中的官方文档为基础,结合 conf/client.conf、PulsarAdmin.java、PulsarAdminImpl.java 等源码与配置,完整讲解三种管理接口的用法、认证(Authentication)与 TLS 配置、Java 客户端构建方式,以及在 Kubernetes 上运行 Functions/Connectors 时的资源命名规则,帮助你从零搭建可用的 Pulsar 管理链路。

三种 Admin 交互方式总览

Pulsar 的 Admin 接口允许你管理实例中的全部重要实体,例如:

  • Tenant(租户):资源隔离与配额分配的逻辑单元;
  • Namespace(命名空间):topic 的组织与管理单元;
  • Topic(主题):消息生产和消费的基本单元。

官方文档明确指出,你可以通过以下三种方式与 Admin 接口交互:

  1. pulsar-adminCLI 工具:位于 Pulsar 安装目录的bin文件夹下,直接执行:

    bin/pulsar-admin

    在源码仓库中,bin/pulsar-admin-common.sh 定义了该工具的启动逻辑,其默认客户端配置指向$PULSAR_HOME/conf/client.conf(即仓库中的 conf/client.conf),日志配置使用conf/log4j2.yaml。

  2. HTTP REST API:由 Pulsar Broker 提供,面向REST的 HTTP 调用。需要注意:部分 RESTful API 可能会被重定向到负责该资源的 owner broker 进行服务,返回307 Temporary Redirect,因此 HTTP 调用方必须处理 307 重定向。如果使用curl,需要显式指定-L参数跟随重定向:

    curl -L http://localhost:8080/admin/v2/namespaces/public/default

    有趣的是,从源码看,Java Admin 客户端内部已经自动处理了这一情况:PulsarAdminImpl.java 在构建 JAX-RSClientConfig时设置了ClientProperties.FOLLOW_REDIRECTS, true,因此 Java 调用方无需自行处理 307。

  3. Java 客户端接口:面向 Java 开发者,通过PulsarAdmin对象以类型安全的方式调用全部管理能力。

核心原则:REST API 就是 Admin 接口本身。无论是pulsar-adminCLI 工具还是 Java 客户端,底层都封装并调用 Broker 的 REST API。如果你需要实现自己的 Admin 接口客户端,应当直接基于 REST API 开发。

开启认证后的 Admin 配置

如果 Pulsar 实例启用了认证,那么三种 Admin 接口(CLI、REST API、Java Admin API)都需要进行相应的专门配置。

pulsar-admin CLI:conf/client.conf 参数详解

pulsar-admin工具的配置默认位于conf/client.conf文件(仓库根目录下的 conf/client.conf 即是真实模板)。官方文档给出的核心参数如下:

参数说明默认值
webServiceUrl集群的 Web(REST)URLhttp://localhost:8080/
brokerServiceUrl集群的 Pulsar 二进制协议 URLpulsar://localhost:6650/
authPlugin认证插件类名空
authParams集群认证参数,逗号分隔的字符串空
useTls集群是否强制启用 TLSfalse
tlsAllowInsecureConnection是否接受客户端无法验证的 TLS 证书false
tlsTrustCertsFilePath受信任 TLS 证书文件路径空

对照仓库中的真实 conf/client.conf,除上述字段外还包含一批生产环境常用的 TLS/KeyStore 参数:

  • tlsEnableHostnameVerification=false:是否校验服务端主机名与证书 CN/SAN 一致;
  • useKeyStoreTls=false:是否使用 KeyStore 类型的 TLS 配置(false表示使用默认 PEM 类型配置);
  • tlsTrustStoreType=JKS:TrustStore 类型,支持JKS、PKCS12;
  • tlsTrustStorePath与tlsTrustStorePassword:TrustStore 路径与密码。

文件中还给出了 TLS 场景下的典型写法:

# TLS 时 Web 服务地址 webServiceUrl=https://localhost:8443/ # TLS 时二进制协议地址 brokerServiceUrl=pulsar+ssl://localhost:6651/ # TLS 认证插件示例 authPlugin=org.apache.pulsar.client.impl.auth.AuthenticationTls # 认证参数为 key:value 的逗号分隔列表,键取决于 authPlugin authParams=tlsCertFile:/path/to/client-cert.pem,tlsKeyFile:/path/to/client-key.pem

REST API:随文档发布

Broker 暴露的 REST API 详细参数、响应与示例由 REST 文档统一提供,本文不重复罗列,直接以{REST}端点文档为准。

Java Admin API:PulsarAdmin 构建实战

要使用 Java Admin API,需要实例化一个PulsarAdmin对象,并指定 Pulsar Broker 的 URL 与对应的PulsarAdminBuilder。下面是官方文档给出的localhost单 Broker 最小示例:

String url = "http://localhost:8080"; // 若启用了 Pulsar 安全认证,传入认证插件类的全限定名 String authPluginClassName = "com.org.MyAuthPluginClass"; // 若认证插件需要参数,传入 auth-param String authParams = "param1=value1"; boolean useTls = false; boolean tlsAllowInsecureConnection = false; String tlsTrustCertsFilePath = null; PulsarAdmin admin = PulsarAdmin.builder() .authentication(authPluginClassName, authParams) .serviceHttpUrl(url) .tlsTrustCertsFilePath(tlsTrustCertsFilePath) .allowTlsInsecureConnection(tlsAllowInsecureConnection) .build();

如果集群有多个 Broker,可以像连接 Pulsar 服务一样使用多主机地址(multi-host),例如同时指定三个 Broker:

String url = "http://localhost:8080,localhost:8081,localhost:8082"; // 认证插件类全限定名(启用安全认证时) String authPluginClassName = "com.org.MyAuthPluginClass"; // 认证插件参数 String authParams = "param1=value1"; boolean useTls = false; boolean tlsAllowInsecureConnection = false; String tlsTrustCertsFilePath = null; PulsarAdmin admin = PulsarAdmin.builder() .authentication(authPluginClassName, authParams) .serviceHttpUrl(url) .tlsTrustCertsFilePath(tlsTrustCertsFilePath) .allowTlsInsecureConnection(tlsAllowInsecureConnection) .build();
源码视角:PulsarAdminBuilder 提供了哪些能力

从源码看,PulsarAdmin.builder()静态方法位于 PulsarAdmin.java,返回一个PulsarAdminBuilder接口实例(实际实现为 PulsarAdminBuilderImpl.java)。Builder 接口 PulsarAdminBuilder.java 提供的方法族可以归为四类:

认证相关(三重重载)

  • authentication(String authPluginClassName, String authParamsString):传入插件类名与逗号分隔字符串参数;
  • authentication(String authPluginClassName, Map<String, String> authParams):传入插件类名与键值对 Map;
  • authentication(Authentication authentication):直接传入已构造好的Authentication实例。

底层实现通过AuthenticationFactory.create(...)创建认证提供者(见 PulsarAdminBuilderImpl.java),并在loadConf(...)时调用setAuthenticationFromPropsIfAvailable自动从配置项authPluginClassName/authParams中恢复认证。

TLS 相关

  • tlsTrustCertsFilePath(...):受信任证书文件路径;
  • allowTlsInsecureConnection(...):是否接受不可信证书,默认false;
  • enableTlsHostnameVerification(...):是否校验主机名(遵循 RFC 2818 3.1 Server Identity 校验);
  • useKeyStoreTls(...)、sslProvider(...)、tlsTrustStoreType(...)、tlsTrustStorePath(...)、tlsTrustStorePassword(...)、tlsCiphers(...)、tlsProtocols(...):KeyStore 模式与加密套件、协议配置。

超时与刷新

  • connectionTimeout(int, TimeUnit):建连超时;
  • readTimeout(int, TimeUnit):服务端响应读取超时;
  • requestTimeout(int, TimeUnit):请求整体超时;
  • autoCertRefreshTime(int, TimeUnit):TLS 认证场景下证书自动刷新间隔。

默认值定义在 PulsarAdminImpl.java 中:连接超时 60 秒、读取超时 60 秒、请求超时 300 秒、证书刷新间隔 300 秒。

通用能力

  • loadConf(Map<String, Object>):从配置 Map 批量载入参数(例如config.put("serviceHttpUrl", "http://localhost:8080"));
  • clone():复制 Builder,便于在同一份基础配置上派生多个客户端;
  • build():构建最终PulsarAdmin实例,URL 缺失时会抛出参数校验异常("Service URL needs to be specified")。
PulsarAdmin 能管理什么

PulsarAdmin对象是全部管理入口的门面(facade),其接口方法(见 PulsarAdmin.java)覆盖了 Pulsar 管理的方方面面:

  • 集群与节点:clusters()、brokers()、bookies()、brokerStats()、proxyStats();
  • 租户与命名空间:tenants()、namespaces()、resourcegroups()、resourceQuotas()(properties()自 2.0 起废弃,改用tenants());
  • 消息与数据:topics()、topicPolicies()、schemas()、lookups()、packages()、transactions();
  • 函数生态:functions()、sources()、sinks()、worker()。

此外,从 PulsarAdminImpl.java 的实现可以看到,PulsarAdmin底层基于 JAX-RSClientBuilder构建 HTTP 客户端,是否启用 TLS 由serviceUrl是否以https://开头自动判断(clientConfigData.getServiceUrl().startsWith("https://")),并默认注册 Jackson JSON 序列化与MultiPartFeature,这些细节都印证了"REST API 是 Admin 接口底层"的架构论断。

Kubernetes 环境下 Pulsar 资源命名规范

如果在 Kubernetes 上运行 Pulsar Functions 或 Connectors,无论使用哪种 Admin 接口,都需要遵循 Kubernetes 的命名规范来定义 Pulsar 资源名称。

Kubernetes 要求名称满足 RFC 1123 中 DNS 子域名的约束,而 Pulsar 本身支持比 Kubernetes 命名规范更多的合法字符(例如 Namespace 名称中允许包含冒号:)。Kubernetes 运行时会把这些 Pulsar 对象名翻译为符合 RFC 1123 的 Kubernetes 资源标签(resource labels),从而让 Functions/Connectors 得以在 Kubernetes 上运行。翻译规则如下:

  • 截断:名称截断至 63 个字符;
  • 替换为短横线(-):将所有非字母数字字符、下划线(_)、点(.)替换为-;
  • 首尾兜底:将开头和结尾的非字母数字字符替换为0。

提示

  • 如果翻译 Pulsar 对象名为 Kubernetes 资源标签时出现错误(例如对象名过长导致命名冲突),或希望自定义翻译规则,可参考 Kubernetes Runtime 的自定义配置(customize Kubernetes runtime)说明;
  • Kubernetes Runtime 的完整配置方式参见 Functions Runtime 的 configure Kubernetes runtime 章节。

一个直观的示例:若 Pulsar Namespace 名为public/default(含点号)且超过 63 字符的长名称,Kubernetes 运行时都会按上述规则转换为形如public-default-xxxxxxxxxx的合法标签。实践中最稳妥的做法是:从一开始就使用 63 字符以内的、仅含字母数字与短横线的资源名,避免触发翻译与潜在冲突。

小结

本指南完整覆盖了 Pulsar Admin 接口的三大入口:

  1. pulsar-adminCLI:开箱即用,配置集中在 conf/client.conf,适合日常命令行管理与脚本化运维;
  2. REST API:一切管理能力的真实载体,调用时务必处理 307 重定向(curl -L);
  3. Java Admin API:通过PulsarAdmin.builder()构建类型安全的管理客户端,PulsarAdminBuilder.java 提供了完整的认证、TLS、超时配置能力,并在多 Broker 场景下支持逗号分隔的多主机 URL。

在开启认证的环境中,请按需配置认证插件与 TLS 参数;在 Kubernetes 上运行函数与连接器时,请严格遵守 RFC 1123 兼容的资源命名规则。理解"REST API 即 Admin 接口"这一架构本质后,你无论是调用 CLI、写 Java 代码,还是自定义管理客户端,都能准确预判行为并排障。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载
上一篇:Data-Juicer与Ray集成:构建企业级数据处理平台的终极指南
下一篇:云存储集成:Instatic与AWS S3、Cloudflare R2配置指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询