- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
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 接口交互:
pulsar-adminCLI 工具:位于 Pulsar 安装目录的bin文件夹下,直接执行:bin/pulsar-admin在源码仓库中,bin/pulsar-admin-common.sh 定义了该工具的启动逻辑,其默认客户端配置指向
$PULSAR_HOME/conf/client.conf(即仓库中的 conf/client.conf),日志配置使用conf/log4j2.yaml。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-RS
ClientConfig时设置了ClientProperties.FOLLOW_REDIRECTS, true,因此 Java 调用方无需自行处理 307。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)URL | http://localhost:8080/ |
brokerServiceUrl | 集群的 Pulsar 二进制协议 URL | pulsar://localhost:6650/ |
authPlugin | 认证插件类名 | 空 |
authParams | 集群认证参数,逗号分隔的字符串 | 空 |
useTls | 集群是否强制启用 TLS | false |
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.pemREST 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 接口的三大入口:
pulsar-adminCLI:开箱即用,配置集中在 conf/client.conf,适合日常命令行管理与脚本化运维;- REST API:一切管理能力的真实载体,调用时务必处理 307 重定向(
curl -L); - 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
相关推荐
Apache Pulsar Admin 接口完全指南:pulsar-admin CLI、REST API 与 Java Admin API
Apache Pulsar Admin 接口完全指南:pulsar admin CLI、REST API 与 Java Admin API 导读 Apache
消息队列后端流处理Apache Pulsar Admin 接口完全指南:pulsar-admin CLI、REST API 与 Java 客户端
Apache Pulsar Admin 接口完全指南:pulsar admin CLI、REST API 与 Java 客户端 Apache Pulsar 的
消息队列后端流处理Apache Pulsar Admin 接口详解:pulsar-admin CLI、REST API 与 Java Admin API 的体系化使用指南
Apache Pulsar Admin 接口详解:pulsar admin CLI、REST API 与 Java Admin API 的体系化使用指南 Apa
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考