
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本指南围绕 Pulsar 的集群Cluster生命周期管理展开覆盖集群的创建Provision、元数据初始化Initialize cluster metadata、配置查询与更新、删除、列表枚举以及 Peer Cluster 配置等核心操作。你将掌握pulsar-admin clusters命令、/admin/v2/clustersREST 端点以及PulsarAdminJava API 三套管理入口的完整用法并理解集群配置数据结构ClusterData与底层元数据存储机制可直接应用于 Pulsar 多集群实例的日常运维。集群在 Pulsar 架构中的角色在 Pulsar 中一个**集群Cluster**由以下三类组件构成一个或多个 PulsarBroker负责消息的生产、消费与路由一个或多个BookKeeper服务器即Bookie负责消息数据的持久化存储一个ZooKeeper集群负责提供配置与协调管理。集群与更高一层的Instance概念相对一个 Pulsar Instance 可以包含多个集群同一实例内的集群之间可以配置为Peer Clusters用于跨集群复制与联邦federation场景。集群的配置信息web service URL、broker service URL、TLS 相关配置、peer 集群列表等统一以ClusterData的形式保存在元数据存储metadata store通常是 ZooKeeper中。对集群的管理可以通过以下三种方式完成管理方式入口适用场景命令行工具pulsar-admin工具的clusters命令族日常运维、脚本化操作REST API/admin/v2/clusters端点集成到 HTTP 工具、监控系统或第三方平台Java APIPulsarAdmin对象上的clusters()方法Clusters接口Java 应用内嵌管理逻辑三种方式在功能上等价REST 端点的服务端实现位于 pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Clusters.java其业务逻辑继承自 ClustersBase而pulsar-adminCLI 与 Java 客户端最终都通过 HTTP 调用同一组 REST 端点例如 ClustersImpl.java 中adminClusters web.path(/admin/v2/clusters)的定位方式。集群配置数据结构ClusterData在深入操作之前先了解集群配置的载体 ——ClusterData。该接口定义于 pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/ClusterData.java核心字段如下字段说明serviceUrl集群的 Web 服务地址HTTP通常为http://host:8080供管理 API 与 HTTP 查询使用serviceUrlTls启用 TLS 时的 Web 服务地址https://host:8443brokerServiceUrlBroker 服务地址Pulsar 协议通常为pulsar://host:6650供客户端连接 Broker 使用brokerServiceUrlTls启用 TLS 时的 Broker 服务地址pulsarssl://host:6651proxyServiceUrl可选的代理服务地址配合proxyProtocol使用proxyProtocol代理协议类型见ProxyProtocol枚举peerClusterNames与该集群互为 peer 的集群名称列表LinkedHashSetStringauthenticationPlugin/authenticationParametersBroker 连接时使用的认证插件及其参数brokerClientTlsEnabled等 TLS 字段Broker 客户端到其他集群的 TLS 连接配置truststore 类型、路径、密码等listenerName监听器名称多监听器场景ClusterData采用 Builder 模式构建可通过ClusterData.builder().serviceUrl(...).build()生成具体实现类为ClusterDataImpl。此外从源码可见createCluster/updateCluster操作对名称有校验要求——服务端在创建时执行NamedEntity.checkName(cluster)集群名中不能包含/字符创建已存在的集群会返回 409Cluster already exists。创建集群Provision创建新集群需要通过管理接口完成该操作要求 superuser超级用户权限。使用 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创建成功后集群配置即被持久化到元数据存储。CLI 的create命令实现于 CmdClusters.java内部调用getAdmin().clusters().createCluster(cluster, clusterData)。使用 REST API 创建PUT /admin/v2/clusters/:cluster请求体为ClusterData的 JSON 表示例如{ serviceUrl: http://my-cluster.org.com:8080, brokerServiceUrl: pulsar://my-cluster.org.com:6650 }服务端对应实现createClusterClustersBase.java先校验 superuser 权限与策略只读状态再检查集群是否已存在存在则返回 409 CONFLICT通过NamedEntity.checkName校验名称合法性非法名称返回 412 PRECONDITION_FAILED最后写入元数据存储。使用 Java API 创建ClusterData clusterData ClusterData.builder() .serviceUrl(serviceUrl) .serviceUrlTls(serviceUrlTls) .brokerServiceUrl(brokerServiceUrl) .brokerServiceUrlTls(brokerServiceUrlTls) .build(); admin.clusters().createCluster(clusterName, clusterData);admin.clusters()返回Clusters接口定义于 pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Clusters.java其实现 ClustersImpl 通过asyncPutRequest(path, ...)将请求发送至PUT /admin/v2/clusters/{cluster}。初始化集群元数据Initialize cluster metadata为什么需要单独初始化创建集群后还需要初始化该集群的元数据metadata store。初始化时需一次性指定以下全部信息集群名称cluster name该集群的本地 ZooKeeper 连接串local ZooKeeper quorum整个 Instance 的配置存储连接串configuration store / global ZooKeeper集群的 Web 服务 URL供客户端与集群内 Broker 交互的 Broker 服务 URL。关键前提必须在启动任何属于该集群的 Broker之前完成元数据初始化。只能通过 CLI 初始化与其他管理功能不同集群元数据初始化无法通过 admin REST API 或 Java admin client 完成因为初始化过程需要直接与 ZooKeeper 通信。必须使用pulsarCLI 工具中的initialize-cluster-metadata命令。示例命令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 pulsarssl://pulsar.us-west.example.com:6651/只有在实例启用了 TLS 认证 时才需要--*-tls系列参数。命令背后的实现逻辑该命令的入口类是 PulsarClusterMetadataSetup.java其main方法按顺序完成一系列初始化动作校验并解析参数--metadata-store或--zookeeper与--configuration-metadata-store或--configuration-store/--global-zookeeper为必填项且新旧参数不能混用——--configuration-metadata-store会取代已废弃的--global-zookeeper与--configuration-store。初始化本地元数据存储与配置元数据存储连接MetadataStoreExtended默认会话超时 30 秒--zookeeper-session-timeout-ms默认 30000ms。格式化 BookKeeper ledger 元数据若/ledgers节点尚不存在且未通过--existing-bk-metadata-service-uri指定既有 BookKeeper 集群。初始化分布式日志DLog命名空间元数据以及 BookKeeper 流存储stream storage的容器元数据默认--initial-num-stream-storage-containers为 16。创建 Bookie 机架感知rack awareness根节点。通过ClusterData.builder()组装集群配置并写入元数据存储若集群不存在则创建。创建名为global的全局集群标记peer 集群与全局策略使用。创建public租户与system租户并将其allowedClusters加入当前集群详见createTenantIfAbsent。创建默认命名空间public/default与系统命名空间设置 replication 集群列表。创建事务协调器transaction coordinator的分区分配 topic分区数默认--initial-num-transaction-coordinators为 16。这也解释了为什么初始化必须在启动 Broker 前执行——Broker 启动时需要从元数据存储读取上述集群、租户、命名空间信息。获取集群配置Get configuration任何时候都可以获取现有集群的配置。使用 pulsar-adminclusters 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 }使用 REST APIGET /admin/v2/clusters/:cluster服务端实现位于 ClustersBase.getCluster要求 superuser 权限集群不存在时返回 404Cluster does not exist。使用 Java APIClusterData clusterData admin.clusters().getCluster(clusterName);更新集群配置Update已有集群的配置可以随时更新操作会整体覆盖ClusterData。使用 pulsar-adminclusters update子命令通过 flags 指定新配置值$ pulsar-admin clusters update cluster-1 \ --url http://my-cluster.org.com:4081 \ --broker-url pulsar://my-cluster.org.com:3350使用 REST APIPOST /admin/v2/clusters/:cluster请求体与 create 相同ClusterDataJSON。服务端 updateCluster 要求 superuser 权限集群不存在时返回 404。使用 Java APIClusterData clusterData ClusterData.builder() .serviceUrl(serviceUrl) .serviceUrlTls(serviceUrlTls) .brokerServiceUrl(brokerServiceUrl) .brokerServiceUrlTls(brokerServiceUrlTls) .build(); admin.clusters().updateCluster(clusterName, clusterData);从源码实现看ClustersImpl.updateClusterAsync 通过asyncPostRequest(path, ...)调用POST端点与createClusterPUT形成对应PUT 用于新建POST 用于更新。删除集群Delete集群可以从 Pulsar Instance 中删除。使用 pulsar-admin$ pulsar-admin clusters delete cluster-1使用 REST APIDELETE /admin/v2/clusters/:cluster使用 Java APIadmin.clusters().deleteCluster(clusterName);删除的前置校验删除并非无条件执行。服务端 deleteCluster 会做以下检查校验 superuser 权限与策略只读状态检查集群是否仍被租户使用isClusterUsed若还有命名空间或命名空间隔离策略namespace isolation policies绑定在该集群上返回 412 PRECONDITION_FAILEDCluster not empty若存在命名空间隔离策略但策略列表为空则先清理隔离策略数据删除关联的 failure domain故障域数据后再删除集群本身。此外CLI 的delete命令支持-a/--all标志CmdClusters.Delete指定后会先遍历删除该集群下所有租户、命名空间、分区 topic 与非分区 topic注意分区 topic 必须通过deletePartitionedTopic删除其 schema然后才调用deleteCluster。列出集群List使用 pulsar-admin$ pulsar-admin clusters list cluster-1 cluster-2使用 REST APIGET /admin/v2/clusters使用 Java APIListString clusters admin.clusters().getClusters();服务端 getClusters 返回所有集群名称的集合并会过滤掉名为global的全局集群标记该集群是initialize-cluster-metadata自动创建的内部标记不参与普通集群枚举。Java 客户端的getClustersAsync通过asyncGetRequest(adminClusters, ...)直接请求根路径GET /admin/v2/clustersClustersImpl.java。配置 Peer ClusterUpdate peer-cluster dataPeer Clusters 用于在同一 Pulsar Instance 内的多个集群之间建立对等关系是跨集群复制与联邦架构的基础配置。可以为指定集群配置其 peer 集群列表。使用 pulsar-admin$ pulsar-admin clusters update-peer-clusters cluster-1 --peer-clusters cluster-2注意命令格式为clusters update-peer-clusters对应子命令定义于 CmdClusters.java。使用 REST APIPOST /admin/v2/clusters/:cluster/peers请求体为 peer 集群名称的 JSON 数组例如[cluster-a, cluster-b]服务端 setPeerClusterNames 会逐一校验peer 列表中不能包含集群自身返回 412且每个 peer 集群必须真实存在不存在返回 412 Peer cluster does not exist校验通过后通过old.clone().peerClusterNames(peerClusterNames).build()保留原有配置并仅更新 peer 列表。使用 Java APIadmin.clusters().updatePeerClusterNames(clusterName, peerClusterList);其中peerClusterList为LinkedHashSetString保持顺序且去重。对应的读取操作为GET /admin/v2/clusters/:cluster/peersJava 侧为admin.clusters().getPeerClusterNames(clusterName)。实战建议与常见问题权限准备除查询类操作get/list同样要求 superuser 权限外集群管理操作一律要求 superuser 权限请提前通过认证配置如 TLS、token 等配置超级用户。操作顺序新集群的标准流程是initialize-cluster-metadata初始化元数据 → 启动 Broker → 通过任一管理入口create集群若初始化时已写入则跳过→ 按需update、配置 peer clusters → 业务发布。名称约束集群名不能包含/字符global是保留的全局集群标记不应作为普通集群名使用。TLS 场景使用 TLS 时务必同时配置serviceUrlTls与brokerServiceUrlTlspulsarssl://协议对应启用 TLS 的 Pulsar 协议端口默认 6651。集群名变更的代价集群名一旦用于命名空间 replication 配置删除重建会影响关联策略建议上线前规划好集群命名。延伸阅读集群与 Instance、Broker 等术语的关系参见 reference-terminology集群配置底层存储与元数据服务说明参见 concepts-architecture-overviewBroker 级别的管理操作参见 admin-api-brokers完整的pulsar-admin clusters命令参数create/update 的全部 flags如--proxy-url、TLS 相关 flags 等可直接运行pulsar-admin clusters --help查看或参考 CmdClusters.java 源码。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐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 Pulsa消息队列后端流处理Apache Pulsar Namespace 管理实战指南pulsar-admin、REST API 与 Java Admin API 三端详解Apache Pulsar Namespace 管理实战指南pulsar admin、REST API 与 Java Admin API 三端详解 本指南基于消息队列后端流处理上一篇FiftyOne自定义数据集构建完全指南标注规范与数据结构设计最佳实践下一篇SonarQube批量分析配置多项目管理与报告生成创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考