- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本篇技术指南以 Apache Pulsar 官方文档 admin-api-clusters.md 为核心骨架,系统讲解 Pulsar 集群(Cluster)的完整生命周期管理:从创建集群(Provision)、初始化集群元数据,到获取配置、更新、删除、列出集群以及配置对等集群(peer-clusters)。读完本文,你将掌握pulsar-adminCLI、/admin/v2/clustersREST API 与 Java Admin API 三种管理方式的具体用法、核心参数含义,以及底层 broker 源码中这些操作的真实实现路径。
Pulsar 集群的构成与管理入口
一个 Pulsar 集群(Cluster)由三部分组成:
- 一个或多个 Pulsarbroker,负责消息路由与读写服务;
- 一个或多个BookKeeper服务器(又称bookie),负责消息的持久化存储;
- 一个ZooKeeper集群,负责提供配置管理与协调管理。
在 Pulsar 的层级体系中,一个instance(实例)可以包含多个 cluster,因此"集群"是命名空间(namespace)、租户(tenant)之上的重要逻辑边界,也是多租户与跨地域复制(geo-replication)的基础单元。相关术语可参考 reference-terminology.md。
集群可以通过以下三种途径进行管理:
| 管理方式 | 入口 | 说明 |
|---|---|---|
| 命令行工具 | pulsar-admin clusters子命令 | 最常用,见 reference-pulsar-admin.md |
| REST API | /admin/v2/clusters端点 | 供脚本、HTTP 客户端调用 |
| Java API | PulsarAdmin对象上的clusters()方法 | 供 Java 应用内嵌调用 |
从源码看,REST 端点由 v2/Clusters.java 暴露,声明为@Path("/clusters"),真正的业务逻辑全部实现在其父类 impl/ClustersBase.java 中;Java Admin API 客户端侧的实现位于 internal/ClustersImpl.java。三条路径最终都指向同一套服务端逻辑。
Provision:创建(预置)新集群
使用 pulsar-admin 创建集群
通过pulsar-admin的clusters create子命令即可预置一个新集群:
$ pulsar-admin clusters create cluster-1 \ --url http://my-cluster.org.com:8080 \ --broker-url pulsar://my-cluster.org.com:6650注意:该操作需要superuser 超级用户权限。
clusters create的完整可用参数(来源:reference-pulsar-admin.md 中clusters create小节):
| Flag | 描述 |
|---|---|
--broker-url | broker 服务(Pulsar 二进制协议)的 URL,默认端口 6650 |
--broker-url-secure | 用于安全连接(TLS)的 broker 服务 URL |
--url | 集群的 web service URL,默认端口 8080 |
--url-secure | 用于安全连接(TLS)的 web service URL |
使用 REST API 创建集群
创建集群对应的 REST 端点为:
PUT /admin/v2/clusters/:cluster其中:cluster为要创建的集群名称。请求体为ClusterData的 JSON 表示。服务端 ClustersBase.java 中createCluster的 Swagger 示例给出了典型 JSON 结构:
{ "serviceUrl": "http://pulsar.example.com:8080", "brokerServiceUrl": "pulsar://pulsar.example.com:6651" }该端点可能的响应码包括:204(创建成功)、403(无超级用户权限)、409(集群已存在)、412(集群名称不合法)、500(内部错误)。同时,服务端还要求集群名不能包含/字符。
使用 Java API 创建集群
ClusterData clusterData = new ClusterData( serviceUrl, serviceUrlTls, brokerServiceUrl, brokerServiceUrlTls ); admin.clusters().createCluster(clusterName, clusterData);ClusterData是集群配置的数据载体,其接口定义位于 pulsar-client-admin-api 的 ClusterData.java。该接口包含的字段远不止上述四个,还支持proxyServiceUrl、proxyProtocol、peerClusterNames、authenticationPlugin、authenticationParameters,以及一组 broker 客户端侧的 TLS 配置(brokerClientTlsEnabled、tlsAllowInsecureConnection、brokerClientTlsTrustStore、brokerClientTrustCertsFilePath、listenerName等),并通过ClusterData.builder()构建器模式创建实例。这意味着集群配置天然支持后续扩展为 TLS、代理与鉴权场景。
Initialize cluster metadata:启动 broker 前的关键一步
创建集群记录之后,还需要初始化该集群的元数据(metadata store 中的初始节点)。初始化时需要明确指定以下所有信息:
- 集群名称(
--cluster) - 该集群的本地 ZooKeeper连接串(
--zookeeper/--metadata-store) - 整个 instance 的配置存储(configuration store)连接串(
--configuration-store/--configuration-metadata-store) - 集群的web service URL(
--web-service-url) - 用于与集群内 broker 交互的broker service URL(
--broker-service-url)
必须在启动任何属于该集群的 broker 之前完成元数据初始化,否则 broker 将无法找到集群的元数据节点而启动失败。
初始化命令使用pulsarCLI 工具(而非pulsar-admin),完整示例如下:
bin/pulsar initialize-cluster-metadata \ --cluster us-west \ --zookeeper zk1.us-west.example.com:2181 \ --configuration-store zk1.us-west.example.com:2184 \ --web-service-url http://pulsar.us-west.example.com:8080/ \ --web-service-url-tls https://pulsar.us-west.example.com:8443/ \ --broker-service-url pulsar://pulsar.us-west.example.com:6650/ \ --broker-service-url-tls pulsar+ssl://pulsar.us-west.example.com:6651/其中--*-tls系列参数仅在实例启用了 TLS 认证时才需要使用,可参考 security-tls-authentication.md。
为什么不能用 REST API 或 Java API 做元数据初始化
与 Pulsar 绝大多数管理功能不同,集群元数据初始化无法通过 admin REST API 或 admin Java 客户端完成。原因在于元数据初始化需要直接与 ZooKeeper(元数据存储)通信、创建海量初始节点,而这两类 API 都经由 broker 的 HTTP 管理端口间接转发。因此必须使用pulsarCLI 的initialize-cluster-metadata命令。
参数明细与源码实现
initialize-cluster-metadata的完整参数(来源:reference-cli-tools.md 中initialize-cluster-metadata小节):
| Flag | 描述 | 默认值 |
|---|---|---|
-c,--cluster | 集群名称 | 必填 |
-uw,--web-service-url | 新集群的 web service URL | 必填 |
-tw,--web-service-url-tls | 启用 TLS 加密的 web service URL | |
-ub,--broker-service-url | 新集群的 broker service URL | |
-tb,--broker-service-url-tls | 启用 TLS 加密的 broker service URL | |
-zk,--zookeeper | 本地 ZooKeeper quorum 连接串 | |
-cs,--configuration-store | 配置存储 quorum 连接串 | |
--existing-bk-metadata-service-uri | 要复用的已有 BookKeeper 集群的元数据服务 URI | |
--initial-num-stream-storage-containers | BookKeeper stream storage 的存储容器数量 | 16 |
--initial-num-transaction-coordinators | 集群中分配的事务协调器数量 | 16 |
--zookeeper-session-timeout-ms | 本地 ZooKeeper 会话超时(毫秒) | 30000 |
-h,--help | 显示帮助信息 | false |
该命令的底层实现位于 PulsarClusterMetadataSetup.java(对应类PulsarClusterMetadataSetup)。从当前仓库源码看,其参数体系已随版本演进有所变化:
--metadata-store(如zk:my-zk:2181)取代了旧参数--zookeeper,后者被标记为hidden;--configuration-metadata-store取代了--configuration-store与--global-zookeeper,后者同样被隐藏标记为弃用;- 若同时传入新旧参数,程序会报错并提示新参数"supersedes the deprecated argument";
--bookkeeper-metadata-service-uri是--existing-bk-metadata-service-uri的兼容别名,已被标记@Deprecated。
命令执行时会同时连接本地元数据存储(local store)与配置元数据存储(configuration store),逐项创建/clusters/<name>、命名空间 bundle、系统租户与命名空间等元数据节点,并可选地初始化 Dlog 命名空间与 BookKeeper stream storage 容器。因此这是一次性的引导操作,同一集群只需(也只能)执行一次。
Get configuration:获取集群配置
已存在的集群可以在任意时刻查询其配置信息。
pulsar-admin
使用clusters get子命令并指定集群名:
$ pulsar-admin clusters get cluster-1 { "serviceUrl": "http://my-cluster.org.com:8080/", "serviceUrlTls": null, "brokerServiceUrl": "pulsar://my-cluster.org.com:6650/", "brokerServiceUrlTls": null "peerClusterNames": null }返回的 JSON 即ClusterData的序列化结果(注:原文档示例中brokerServiceUrlTls行后缺少一个逗号,实际输出以服务端为准)。
REST API
GET /admin/v2/clusters/:cluster成功时返回200与集群配置数据(ClusterDataImpl),404表示集群不存在。
Java API
admin.clusters().getCluster(clusterName);Update:更新集群配置
集群配置创建后可以随时更新。
pulsar-admin
使用clusters update子命令,通过 flag 指定新的配置值:
$ pulsar-admin clusters update cluster-1 \ --url http://my-cluster.org.com:4081 \ --broker-url pulsar://my-cluster.org.com:3350REST API
POST /admin/v2/clusters/:cluster服务端updateCluster方法会先校验集群存在(404不存在),再写入新配置;403表示无权限或策略只读。
Java API
ClusterData clusterData = new ClusterData( serviceUrl, serviceUrlTls, brokerServiceUrl, brokerServiceUrlTls ); admin.clusters().updateCluster(clusterName, clusterData);Delete:删除集群
集群可以从所在的 Pulsar instance 中删除。
pulsar-admin
$ pulsar-admin clusters delete cluster-1REST API
DELETE /admin/v2/clusters/:clusterJava API
admin.clusters().deleteCluster(clusterName);实战提醒:服务端
deleteCluster的 Swagger 定义中列出了412响应码,含义为 "Cluster is not empty"(集群非空)。也就是说,当集群下仍存在命名空间、租户等资源时,删除会被拒绝。正确流程是先清理该集群下的命名空间与租户,再删除集群本身。
List:列出集群
可以获取当前 Pulsar instance 中全部集群的列表。
pulsar-admin
$ pulsar-admin clusters list cluster-1 cluster-2REST API
GET /admin/v2/clusters返回集群名的集合;服务端实现对应 ClustersBase.java 中的getClusters方法。
Java API
admin.clusters().getClusters();Update peer-cluster data:配置对等集群
**对等集群(peer clusters)**是 Pulsar 多集群 / 跨地域复制场景中的关键概念:将多个集群互相声明为 peer 后,这些集群可以互为对等方,用于拓扑感知与故障转移场景。
pulsar-admin
$ pulsar-admin clusters update-peer-clusters cluster-1 --peer-clusters cluster-2(原文档示例写作pulsar-admin update-peer-clusters cluster-1 --peer-clusters cluster-2,按 reference-pulsar-admin.md 中clusters子命令列表,正确形式应为pulsar-admin clusters update-peer-clusters ...,运行时以pulsar-admin clusters --help输出为准。)
REST API
POST /admin/v2/clusters/:cluster/peers请求体为 peer 集群名列表,例如:
[ "cluster-a", "cluster-b" ]Java API
admin.clusters().updatePeerClusterNames(clusterName, peerClusterList);源码中的校验逻辑
服务端setPeerClusterNames方法(见 ClustersBase.java)在写入前会执行严格校验:
- 先校验被配置的集群本身是否存在;
- 遍历传入的
peerClusterNames,逐一检查每个 peer 集群是否真实存在,若不存在则直接返回412 Peer cluster doesn't exist; - 校验通过后,将
peerClusterNames写入集群配置(ClusterData.peerClusterNames),并记录操作日志。
也就是说,--peer-clusters中出现的每一个名字都必须是在当前 instance 中已创建成功的集群,否则整个操作会被拒绝。这与 ClusterData.java 中peerClusterNames字段(类型为LinkedHashSet<String>,保证顺序且去重)一一对应。
总结与延伸阅读
集群管理是搭建 Pulsar 生产环境的第一个环节。标准操作序列是:
- 用
bin/pulsar initialize-cluster-metadata初始化集群元数据(一次性、必须早于 broker 启动); - 用
pulsar-admin clusters create(或 REST/Java API)登记集群配置; - 用
clusters get / update / list维护配置,用clusters delete回收资源; - 多集群场景下用
clusters update-peer-clusters建立对等关系。
如需继续深入,可参阅同目录下的相关文档:
- reference-pulsar-admin.md:
pulsar-admin全部子命令与 flag 参考; - reference-cli-tools.md:
pulsarCLI(含initialize-cluster-metadata)参考; - reference-configuration.md:broker 等组件的完整配置项;
- admin-api-brokers.md:集群内 broker 的运维管理;
- concepts-architecture-overview.md:理解 metadata store、broker、bookie 在架构中的角色。
源码层面的关键入口:REST 端点 v2/Clusters.java、实现逻辑 impl/ClustersBase.java、数据模型 ClusterData.java,以及元数据初始化工具 PulsarClusterMetadataSetup.java。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 集群管理完全指南:pulsar-admin、REST API 与 Java Admin API 实战
Apache Pulsar 集群管理完全指南:pulsar admin、REST API 与 Java Admin API 实战 本指南以 admin api
消息队列后端流处理Apache Pulsar 集群管理实战:pulsar-admin、REST API 与 Java Admin API 全面指南
Apache Pulsar 集群管理实战:pulsar admin、REST API 与 Java Admin API 全面指南 本文基于 Apache Pul
消息队列后端流处理Apache Pulsar 集群管理完全指南:基于 pulsar-admin、REST API 与 Java Admin API 的 Clusters 资源操作
Apache Pulsar 集群管理完全指南:基于 pulsar admin、REST API 与 Java Admin API 的 Clusters 资源操作
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考