- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
Apache Pulsar 从设计之初就是一套多租户(multi-tenant)的消息系统,租户(Tenant)与命名空间(Namespace)共同构成了其资源隔离与策略治理的基石。本文基于 Pulsar 官方概念文档与仓库源码,系统讲解 Pulsar 的 Topic URL 三级结构、租户与命名空间的划分与管理方式,并深入剖析用于支撑 Topic 级策略的命名空间变更事件(Namespace Change Events)机制及其__change_events系统主题的底层实现。读完本文,你将掌握 Pulsar 多租户模型的核心概念、管理 API 的实际用法,以及系统主题在策略传播与缓存刷新中的完整工作链路。
多租户设计:从底层构建的隔离模型
Pulsar 不是通过后期改造来支持多租户,而是在创建之初就将其作为核心能力设计。多租户意味着一个 Pulsar 实例(instance)可以同时服务于多个相互隔离的业务方,且这种隔离可以跨越多个集群:
- 租户(Tenant)可以横跨多个集群,即同一个租户的资源可以分布在不同的物理集群上;
- 每个租户可以单独配置认证与授权方案(authentication and authorization),互不干扰;
- 租户是管理存储配额、消息 TTL 和隔离策略(isolation policies)的行政单元,这些策略以租户为边界进行设定与计量。
因此,多租户并非简单的"多用户"概念,而是一套从命名、鉴权、配额到策略传播的完整治理体系。仓库中pulsar-client-admin-api模块下的 Tenants.java 与 Namespaces.java 接口,即是对这套治理模型的管理面抽象。
Topic URL 的三级结构:tenant/namespace/topic
Pulsar 多租户特性最直观的体现是 Topic 的命名结构。所有 Topic 的 URL 都遵循以下格式:
persistent://tenant/namespace/topicURL 由三段组成,从语义上看:
| 段 | 含义 | 说明 |
|---|---|---|
persistent:// | 域(domain) | 标识持久化主题;非持久化主题使用non-persistent:// |
tenant | 租户 | 多租户体系中最基本的分类单元 |
namespace | 命名空间 | 租户内部的细分命名空间 |
topic | 主题 | 实际承载消息的主题名 |
值得强调的是:租户是比命名空间和 Topic 名称更基本的分类维度。Topic 的归属关系是"租户 → 命名空间 → Topic"逐级嵌套的,任何 Topic 都必须归属于某个租户,而租户的鉴权、配额等策略会通过命名空间逐级向下生效。在源码中,这一层级结构由 TopicName.java 与 NamespaceName.java 等命名解析类共同实现。
Tenant:租户
租户是 Pulsar 多租户体系中的核心管理单元。针对实例中的每个租户,管理员可以为其配置:
- 授权方案(authorization scheme):每个租户可以绑定独立的认证/授权策略,例如基于角色的访问控制,实现租户间管理权限的彻底隔离;
- 集群集合(set of clusters):指定租户配置生效的集群范围,使租户可以横跨多个集群进行消息读写与复制。
租户级配置由TenantInfo承载。从管理实现看,TenantsImpl.java 将租户管理 REST 端点定义为web.path("/admin/v2/tenants"),其中:
createTenantAsync(String tenant, TenantInfo config)通过PUT/admin/v2/tenants/{tenant}创建租户;updateTenantAsync通过 POST 更新租户配置(如调整授权方案与集群集合);deleteTenant支持普通删除与force强制删除;getTenantInfoAsync用于查询租户详情。
同时,在 PulsarAdminImpl.java 中,this.tenants = new TenantsImpl(root, auth, readTimeoutMs)表明所有pulsar-admin租户命令最终都会经由这条 REST 通道落地到 Broker。
Namespace:命名空间
租户与命名空间是支撑 Pulsar 多租户的两大关键概念,二者的定位分工如下:
- 租户是资源配额与容量分配的单位:Pulsar 会为指定租户进行容量规划,为其分配适当的存储、带宽等资源;
- 命名空间是租户内部的行政细分单位:命名空间内设置的所有配置策略,会统一作用于该命名空间下创建的所有 Topic。
命名空间的自管理能力
一个租户可以包含多个命名空间。多租户体系允许租户通过**自管理(self-administration)**方式使用REST API与pulsar-adminCLI 工具自行创建命名空间。例如,一个承载了多个不同应用的租户,可以为每个应用创建独立的命名空间,实现策略上的相互隔离。
同一命名空间下的 Topic 名称形如:
persistent://tenant/app1/topic-1 persistent://tenant/app1/topic-2 persistent://tenant/app1/topic-3从管理面接口 Namespaces.java 可以看出,命名空间支持丰富的策略配置能力,例如:
createNamespace(String namespace)、createNamespace(String namespace, int numBundles)创建命名空间,并可指定 bundle 数量;setNamespaceMessageTTL(String namespace, int ttlInSeconds)设置消息 TTL;setNamespaceReplicationClusters(String namespace, Set<String> clusterIds)配置跨集群复制;grantPermissionOnNamespace/revokePermissionsOnNamespace管理命名空间权限;setAutoTopicCreation、setDeduplicationStatus等控制自动建主题、去重等行为。
这些方法体现了"命名空间即策略作用域"的设计:一旦策略在命名空间级别设定,其下所有 Topic 自动继承。
命名空间的全局唯一性
命名空间作为 Topic URL 的中间段,在同一 Pulsar 实例内以tenant/namespace组合唯一标识,由 NamespaceName.java 进行解析与校验。这也保证了下文所述的命名空间变更事件可以精确地归属到具体的租户与命名空间。
Namespace Change Events 与 Topic 级策略
Pulsar 是多租户的事件流系统,管理员可以在不同层级(租户级、命名空间级)设置策略。但保留策略(retention policy)、存储配额策略(storage quota policy)等传统上只能设定在命名空间级别。在实际使用中,用户经常需要对单个 Topic 设置策略。为此,Pulsar 提出了命名空间变更事件(Namespace Change Events)机制,以高效的方式支持 Topic 级策略。
为什么选择事件驱动方案
将"策略变更"以事件的形式写入系统主题而非直接操作元数据服务,主要带来三方面收益:
- 避免给 ZooKeeper 增加负载:策略变更不再频繁读写 ZooKeeper,从而减轻元数据服务的压力;
- 利用 Pulsar 本身作为事件日志:策略缓存(policy cache)的传播借助 Pulsar 自身的消息流转发能力,天然具备水平扩展能力;
- 可审计、可查询:可以使用 Pulsar SQL 查询命名空间变更记录,对系统进行审计。
__change_events系统主题
每个命名空间都有一个名为__change_events的系统主题(system topic),用于存放该命名空间下的变更事件。在源码中,该主题名被定义为常量:
// pulsar-common/src/main/java/org/apache/pulsar/common/events/EventsTopicNames.java public static final String NAMESPACE_EVENTS_LOCAL_NAME = "__change_events";主题的实际构造在 NamespaceEventsSystemTopicFactory.java 中完成:
public TopicPoliciesSystemTopicClient createTopicPoliciesSystemTopicClient(NamespaceName namespaceName) { TopicName topicName = TopicName.get("persistent", namespaceName, EventsTopicNames.NAMESPACE_EVENTS_LOCAL_NAME); return new TopicPoliciesSystemTopicClient(client, topicName); }也就是说,__change_events是一个持久化系统主题,其完整路径形如persistent://tenant/namespace/__change_events。它同时承载两种职责:一方面作为策略变更的写入目标,另一方面作为各 Broker 订阅读取的变更日志。
四步工作流程
下图描述了使用命名空间变更事件实现 Topic 级策略的完整链路:
Pulsar Admin 客户端 │ ① 调用 Admin RESTful API 更新 Topic 级策略 ▼ Admin RESTful API │ ② 收到请求的 Broker 将 Topic 策略变更事件发布到 │ 对应命名空间的 __change_events 主题 ▼ __change_events 主题(持久化事件日志) │ ③ 拥有该命名空间 bundle 的每个 Broker 订阅 │ __change_events 主题,接收变更事件并应用到策略缓存 ▼ 策略缓存(policy cache) │ ④ 缓存更新完成后,Broker 向 Pulsar Admin 客户端返回响应 ▼ Pulsar Admin 客户端具体步骤可拆解为:
- Pulsar Admin 客户端调用Admin RESTful API提交 Topic 级策略更新请求;
- 收到该 HTTP 请求的 Broker将 Topic 策略变更事件(含
UPDATE/DELETE等动作类型)发布到该命名空间对应的__change_events主题; - 每个拥有该命名空间 bundle 的 Broker订阅
__change_events主题,持续接收命名空间的变更事件,并将其应用到本地的策略缓存(policy cache); - 策略缓存更新完成后,Broker 向 Pulsar Admin 客户端返回响应,客户端感知到策略已生效。
源码级实现印证
上述工作流程在 SystemTopicBasedTopicPoliciesService.java 中得到了完整实现,该服务即"基于系统主题的 Topic 策略服务",其核心结构包括:
policiesCache与globalPoliciesCache:本地方针缓存与全局方针缓存(ConcurrentHashMap),按 Topic 维度缓存最新策略;readerCaches:以命名空间为键缓存的系统主题 Reader,每个 Reader 对应一个__change_events订阅;listeners:Topic 策略监听器注册表,用于在事件到达时通知相关组件刷新策略。
写入侧:updateTopicPoliciesAsync/deleteTopicPoliciesAsync最终汇聚到sendTopicPolicyEvent,该方法通过namespaceEventsSystemTopicFactory.createTopicPoliciesSystemTopicClient(...)创建系统主题客户端,再newWriterAsync()打开 Writer,将携带ActionType.UPDATE/DELETE与PulsarEvent的事件写入__change_events:
return systemTopicClient.newWriterAsync() .thenCompose(writer -> { PulsarEvent event = getPulsarEvent(topicName, actionType, policies); CompletableFuture<MessageId> writeFuture = ActionType.DELETE.equals(actionType) ? writer.deleteAsync(event) : writer.writeAsync(event); ... });订阅侧:start()方法注册了NamespaceBundleOwnershipListener,当 Broker 加载(onLoad)或卸载(unLoad)某个命名空间 bundle 时,相应地addOwnedNamespaceBundleAsync/removeOwnedNamespaceBundleAsync。加载 bundle 后,服务会为该命名空间创建__change_events的 Reader,通过initPolicesCache回放历史事件初始化策略缓存,再由readMorePolicies循环读取后续事件,并在refreshTopicPoliciesCache中按INSERT/UPDATE/DELETE动作类型更新缓存、经notifyListener通知监听器。
值得注意的是,事件读取采用了带退避(Backoff)的重试机制:createSystemTopicClientWithRetry中使用Backoff(1s, 3s, 10s)配合RetryUtil.retryAsynchronously创建 Reader,以增强系统主题订阅的健壮性。
全局策略与本地策略:getPulsarEvent中有一个关键分支——当策略不是全局策略(!policies.isGlobalPolicies())时,会将replicateTo设置为仅包含本地集群,注释明确写道"本地策略无需复制到远端集群";而refreshTopicPoliciesCache则依据是否携带replicateTo分别写入globalPoliciesCache或policiesCache。这保证了全局策略可跨集群传播、本地策略仅作用于本集群,是 Topic 级策略在多租户、多集群场景下正确生效的基础。
与租户、命名空间的关系
整个事件链路中的每一条TopicPoliciesEvent都携带完整的归属信息。从 TopicPoliciesEvent.java 的数据模型看,事件中记录了domain、tenant、namespace、topic与policies字段;消费侧在refreshTopicPoliciesCache中通过TopicName.get(event.getDomain(), event.getTenant(), event.getNamespace(), event.getTopic())重建完整 Topic 名。这再次印证了多租户的根基:任何一条策略变更都能精确回溯到某个租户下某个命名空间中的某个 Topic。
小结
Apache Pulsar 的多租户模型以persistent://tenant/namespace/topic的三级 URL 为骨架,以租户为资源配额与鉴权的基本单元、以命名空间为策略治理的作用域,实现了"一个实例、多租户隔离、策略可下钻到单 Topic"的完整能力。而命名空间变更事件机制则巧妙地将策略传播从 ZooKeeper 中解放出来,借助每个命名空间下的__change_events系统主题,以"事件日志 + 订阅消费 + 本地缓存"的方式高效、可扩展、可审计地驱动 Topic 级策略的生效。
对于开发者而言,理解这套机制意味着:
- 在设计多租户部署时,可以按"租户 → 命名空间"的两级维度规划资源、鉴权与策略边界;
- 在排查策略未生效的问题时,可以沿着"Admin REST →
__change_events写入 → Broker 订阅消费 → 策略缓存刷新"这条链路逐层定位,而 SystemTopicBasedTopicPoliciesService.java 与 EventsTopicNames.java 是深入源码的最佳切入点; - 在需要审计策略变更时,可借助 Pulsar SQL 查询
__change_events中的事件记录,还原每一次 Topic 级策略调整的全过程。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
深入解析 Apache Pulsar 多租户架构:Tenant、Namespace 与命名空间变更事件机制
深入解析 Apache Pulsar 多租户架构:Tenant、Namespace 与命名空间变更事件机制 Apache Pulsar 从设计之初就定位为原生多
消息队列后端流处理Apache Pulsar 多租户架构:Tenant、Namespace 与 Namespace 变更事件(__change_events)机制解析
Apache Pulsar 多租户架构:Tenant、Namespace 与 Namespace 变更事件(__change_events)机制解析 多租户(M
消息队列后端流处理Apache Pulsar 多租户机制深度解析:Tenant、Namespace 与 Namespace 变更事件
Apache Pulsar 多租户机制深度解析:Tenant、Namespace 与 Namespace 变更事件 Apache Pulsar 从底层设计之初就
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考