- 消息队列
- 流处理
- 后端
- 微服务
- 消息路由
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
PIP-92(Topic policy across multiple clusters)是 Apache Pulsar 社区针对地理复制(geo-replication)场景提出的一项关键设计:将 Topic 级别的策略区分为"全局策略(Global Topic Policy)"与"本地策略(Local Topic Policy)",前者随复制链路同步到整个复制组的所有集群,后者只作用于当前集群。本文以仓库中的 pip-92.md 为骨架,结合 TopicPolicies.java、TopicPoliciesService.java、CmdTopicPolicies.java 等源码实现,完整讲解该机制的设计动机、数据模型、策略生效优先级、命令行用法、底层存储与复制原理,并给出可复现的验证路径。
一、背景与动机:为什么 Topic 策略需要"全局/本地"之分
在 Pulsar 的地理复制架构中,一个 Topic 可以在多个集群之间双向复制消息(replicationClusters配置了跨集群复制)。此时,管理员为这个 Topic 设置的策略(retention、backlog quota、dispatch rate、message TTL 等)会面临一个天然的二义性问题:
- 某些策略希望作用于整个复制组。例如"该 Topic 的消息在所有集群都保留 3 天"——如果只在其中一个集群设置,其他集群的存储行为就会不一致;
- 某些策略只希望作用于单个集群。例如"A 集群限制该 Topic 的最大生产者数为 10,B 集群不限制",这是因为各集群的容量、租户配额、资源水位各不相同。
在 PIP-92 之前,Topic 策略只有一份,通过系统 Topic(System Topic)在集群间同步,无法表达上述差异。PIP-92 的核心提案因此非常直接:为 TopicPolicies 数据结构增加一个isGlobal布尔标志,让 Replicator(复制器)只把带全局标志的策略复制到其他集群,从而形成全局策略与本地策略两套并存的数据。
二、核心数据模型:TopicPolicies 与 isGlobal 标志
PIP-92 给出的最简设计是在TopicPolicies上增加一个字段:
public class TopicPolicies { boolean isGlobal = false; }这一设计在当前仓库中已经落地,见 pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/TopicPolicies.java:
@Builder.Default private Boolean isGlobal = false;该字段用@Builder.Default标注,默认值为false,即默认所有策略都是本地策略——这与 PIP-92 的兼容性声明完全一致:"现有策略中isGlobal默认即为 false,因此本方案不引入任何兼容性问题"。
TopicPolicies是 Topic 级别所有可配置策略的聚合载体,除isGlobal外还包括(均可在同一文件中核对):
| 策略字段 | 作用 |
|---|---|
backLogQuotaMap | 每个订阅级别的积压配额 |
retentionPolicies | 消息保留时长与大小限制 |
deduplicationEnabled | 是否开启消息去重 |
messageTTLInSeconds | 消息 TTL(秒) |
maxProducerPerTopic/maxConsumerPerTopic | 单 Topic 生产者/消费者上限 |
maxConsumersPerSubscription | 单订阅最大消费者数 |
maxUnackedMessagesOnConsumer/OnSubscription | 未确认消息上限 |
delayedDeliveryTickTimeMillis/delayedDeliveryEnabled | 延迟投递相关 |
dispatchRate/subscriptionDispatchRate/replicatorDispatchRate | 各类投递限流 |
compactionThreshold | 压缩阈值 |
publishRate/subscribeRate | 生产/订阅限流 |
offloadPolicies | 分层存储卸载策略 |
inactiveTopicPolicies | 不活跃 Topic 的删除/保留策略 |
schemaCompatibilityStrategy | Schema 兼容性策略 |
subscriptionPolicies | 订阅级别策略(按订阅名细分) |
源码中还提供了一个便捷方法isGlobalPolicies()(TopicPolicies.java),返回isGlobal != null && isGlobal,用于在 Broker 侧统一判断"这份策略是否是全局策略"。
三、策略生效优先级:四层覆盖规则
PIP-92 明确规定,当同时存在多份策略时,Topic 最终生效的策略按照以下优先级从高到低取值:
1. 本地集群 Topic 策略(Local cluster topic policy) 2. 全局 Topic 策略(Global topic policy) 3. 命名空间策略(Namespace policy) 4. Broker 默认配置(Broker default configuration)也就是说:本地策略优先于全局策略,全局策略优先于命名空间策略,命名空间策略优先于 Broker 配置。本地策略只影响当前集群,全局策略则通过复制影响复制组内所有集群。
这一优先级在测试中被严格验证。仓库测试 pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java 中的testPriorityOfGlobalPolicies与testPriorityOfGlobalPolicies2分别验证了两个方向:
- 先设置本地策略、再设置全局策略:断言"Topic 策略(本地)的优先级高于全局策略",即本地值仍然生效;
- 先设置全局策略、再设置本地策略:断言本地策略生效;
- 随后删除本地策略(通过
isGlobal=false的移除操作),断言"全局策略重新生效"; - 反向删除全局策略后,断言本地策略重新生效。
此外testGlobalPolicyStillAffectsAfterUnloading验证了 Topic 卸载(unload)重载后,本地与全局策略依然同时生效;testRetentionGlobalPolicyAffects验证了全局 retention 策略的实际作用。这些用例共同构成了 PIP-92 "Test Plan" 中"Priority of Topic Policies matches our setting"一项的实现证据。
四、命令行实操:--global 选项的完整用法
PIP-92 要求每一个 Topic 策略管理 API 都要增加--global选项,覆盖 Broker REST API、Admin SDK 与命令行(CMD)三个层面。
4.1 命令行示例
PIP-92 给出的核心命令示例为设置全局 retention:
bin/pulsar-admin topics set-retention -s 1G -t 1d --global my-topic各参数含义:
| 参数 | 含义 |
|---|---|
-s 1G | 每个分区最大保留存储大小为 1 GiB |
-t 1d | 每条消息最大保留时长为 1 天 |
--global | 本次设置的是全局 Topic 策略,将随复制同步到其他集群 |
my-topic | 目标 Topic(可带persistent://tenant/namespace/topic全名) |
PIP-92 明确说明:如果不加--global,行为与之前完全一致,只更新本地策略——这正是isGlobal默认false的语义在 CLI 层的体现。
4.2 当前仓库中的选项实现
该--global选项在当前仓库的命令行实现中随处可见,例如 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopicPolicies.java 中,几乎每个子命令都声明了:
@Option(names = { "--global", "-g" }, description = "Whether to set this policy globally. " + "If set to true, broker returned global topic policies") private boolean isGlobal = false;注意这里同时注册了短选项-g,即以下两条命令等价:
bin/pulsar-admin topics set-message-ttl --global -t 3600 my-topic bin/pulsar-admin topics set-message-ttl -g -t 3600 my-topic从源码结构可以确认,该选项覆盖的策略命令至少包括:set/get/remove-retention、set/get/remove-message-ttl、set/get/remove-max-consumers、set/get/remove-max-producers、set/get/remove-max-unacked-messages-on-consumer、set/get/remove-max-unacked-messages-on-subscription、set/get/remove-max-consumers-per-subscription、set/get/remove-subscription-expiration-time、set/get/remove-subscription-types-enabled、set/get/remove-delayed-delivery-policy、set/get/remove-entry-filters等,与 PIP-92 "Every topic policy API adds the--globaloption" 的要求一一对应。
4.3 Admin SDK 与 REST API 层面
在管理 API 层面,isGlobal同样贯穿始终:
- Admin SDK:例如 pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicPoliciesImpl.java 中,
setRetention、setMaxProducers等方法均带有isGlobal参数(或admin.topicPolicies(true)这种按全局开关的调用形式,见测试代码admin.topicPolicies(true).setDispatchRate(...)); - REST API:测试 TopicPoliciesTest.java 中通过 HTTP 客户端直接携带查询参数
isGlobal=false/isGlobal=true发起策略设置请求,验证了 REST 层对该参数的透传。
4.4 读取时的区分
与写入对应,读取策略时也按全局/本地区分。在 Broker 侧的 TopicPoliciesService.java 中定义了读取枚举:
enum GetType { GLOBAL_ONLY, // only get the global policies LOCAL_ONLY, // only get the local policies }配合getTopicPoliciesAsync(TopicName topicName, GetType type)使用:GLOBAL_ONLY只返回全局策略,LOCAL_ONLY只返回本地策略。这也印证了 PIP-92 "以前内存中只缓存一份本地策略,现在变成两份:一份 Local、一份 Global" 的设计——两份策略在读写两侧都是独立管理的。
五、底层实现:System Topic 存储、复制与全局/本地区分
5.1 策略的存储介质
PIP-92 指出"Topic 策略存储在 System Topic 中,可以直接利用 Replicator 把数据复制到其他集群"。当前仓库的实现正是如此:策略变更以PulsarEvent事件的形式写入系统 Topic(__change_events),由 SystemTopicBasedTopicPoliciesService.java 负责读写、缓存与监听。
5.2 全局策略的区分标记
PIP-92 最初设想通过在 Replicator 接口上新增setFilterFunction(Function<Message, Boolean>)过滤器(返回false则过滤掉该消息),用一个"只放行isGlobal=true数据"的函数来实现选择性复制。需要说明的是:在当前仓库代码中,setFilterFunction尚未出现(搜索setFilterFunction仅命中 pip-92.md 本身),即该 API 属于 PIP 提案内容,后续实现演进为另一种更轻量的机制:
当前实现通过消息 key 前缀区分全局与本地策略。在 TopicPoliciesService.java 中可以看到:
String GLOBAL_POLICIES_MSG_KEY_PREFIX = "__G__";配套的三个静态方法构成了区分与还原逻辑:
// 写入时:全局策略的 key 加上 "__G__" 前缀,本地策略保持原 key static String wrapEventKey(String originalKey, boolean isGlobalPolicies) { if (!isGlobalPolicies) { return originalKey; } return GLOBAL_POLICIES_MSG_KEY_PREFIX + originalKey; } // 读取时:通过 key 前缀判断事件是否为全局策略 static boolean isGlobalPolicy(Message<PulsarEvent> msg) { return msg.getKey().startsWith(GLOBAL_POLICIES_MSG_KEY_PREFIX); } // 还原:去掉前缀得到真实 Topic 名 static TopicName unwrapEventKey(String originalKey) { ... }结合 SystemTopicBasedTopicPoliciesService.java 的实现可以看到:写入时调用getEventKey(event, isGlobalPolicy)计算带前缀的 key,系统 Topic 消费者侧则用isGlobalPolicy(msg)判断后分别更新globalPoliciesCache或本地策略缓存——该服务持有独立的globalPoliciesCache(final Map<TopicName, TopicPolicies> globalPoliciesCache),与本地缓存分开维护。由于系统 Topic 本身具备跨集群复制能力,带__G__前缀的全局策略事件会随之同步到其他集群,从效果上等价于 PIP-92 中"只复制全局策略"的目标。
5.3 缓存的读写一致性
值得补充的是,PIP-92 提出的"内存中缓存两份策略"在其后的演进中进一步解决了并发一致性隐患:相关设计在 pip-428.md 中做了系统性完善——通过"克隆 + 函数式更新"(updateTopicPoliciesAsync(topicName, isGlobalPolicy, skipUpdateWhenTopicPolicyDoesntExist, policyUpdater),其中policyUpdater接收一个可安全修改的克隆副本)、按Pair<TopicName, Boolean(isGlobal)>键控的顺序化更新队列(sequencer)以及"读己之写"(read-your-writes)保证,解决了多策略快速连续更新时的互相覆盖与丢失问题。当前 TopicPoliciesService.java 中的接口签名正是该演进后的版本。
六、策略在 Topic 上的应用:本地优先的合并逻辑
全局与本地两份策略最终要合并到 Topic 的实际运行参数上,合并原则即第三节的四层优先级。这一逻辑体现在 Topic 加载时的应用过程,见 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java:
boolean isGlobalPolicies = data.isGlobalPolicies(); ... topicPolicies.getRetentionPolicies().updateTopicValue(data.getRetentionPolicies(), isGlobalPolicies); topicPolicies.getMaxProducers().updateTopicValue(normalizeValue(data.getMaxProducers()), isGlobalPolicies); ...updateTopicValue(value, isGlobalPolicies)这类方法接收一个isGlobalPolicies标志,在内部决定"本地值覆盖全局值 / 全局值作为兜底"的合并策略——当本地值已设置时保持本地值,本地值未设置时才采用全局值,由此实现"Local > Global"的优先级;而TopicPolicies本身未设置时,继续回落到命名空间策略与 Broker 默认配置。
另外在 AbstractTopic.java 中可以看到,Topic 加载时会并行获取两份策略:
CompletableFuture<Optional<TopicPolicies>> globalPoliciesFuture = topicPoliciesService.getTopicPoliciesAsync(topicName, TopicPoliciesService.GetType.GLOBAL_ONLY); CompletableFuture<Optional<TopicPolicies>> localPoliciesFuture = topicPoliciesService.getTopicPoliciesAsync(topicName, TopicPoliciesService.GetType.LOCAL_ONLY); return globalPoliciesFuture.thenCombine(localPoliciesFuture, (global, local) -> { ... });即全局与本地策略分别异步读取后合并应用,从代码结构上印证了 PIP-92 的双缓存、双读取设计。
七、删除策略时的行为差异
PIP-92 特别提出了一个问题:"删除一个集群会怎样?"并给出了明确答案:
删除集群只会删除本地 Topic 策略,不会影响其他集群——因为只有全局策略才会被复制到其他集群。
其推理链条是:跨集群复制的对象仅是isGlobal=true的全局策略;本地策略留在本集群的系统 Topic 中,不参与复制。因此无论删除单个集群的本地策略,还是删除整个集群,都不会波及其他集群的策略状态。当前实现中同样可以看到相关语义:TopicPoliciesService的deleteTopicPoliciesAsync提供了keepGlobalPoliciesAfterDeleting开关,并在 SystemTopicBasedTopicPoliciesService.java 的删除逻辑中区分"跳过删除全局策略(对应 PIP-422 的语义)"与"删除全局策略"两种路径,进一步细化了这一行为。
八、兼容性与测试计划
8.1 兼容性
PIP-92 声明该方案不引入任何兼容性问题,理由是isGlobal字段在既有策略中默认即为false,旧数据无需迁移即可继续按本地策略解释。从 TopicPolicies.java 的@Builder.Default private Boolean isGlobal = false;可以确认该默认值语义,且所有既有getTopicPolicies(isGlobal)调用路径都保持默认false的行为(如 CmdTopicPolicies.java 中private boolean isGlobal = false;的默认值)。
8.2 测试计划对照
PIP-92 列出的 Test Plan 在当前仓库均有对应实现:
| PIP-92 测试计划项 | 仓库对应证据 |
|---|---|
| Existing TopicPolicies 不受影响 | TopicPoliciesUpdateTest.java、TopicPoliciesDisableTest.java 等既有用例持续运行 |
| 只复制全局策略、Filter 逻辑按预期工作 | SystemTopicBasedTopicPoliciesService.java 中基于__G__前缀的isGlobalPolicy(msg)判定 +globalPoliciesCache独立缓存;TopicPoliciesTest.java 的testGlobalPolicyStillAffectsAfterUnloading等用例覆盖复制后生效场景 |
| 策略优先级符合预期 | TopicPoliciesTest.java 的testPriorityOfGlobalPolicies/testPriorityOfGlobalPolicies2/testRetentionGlobalPolicyAffects |
九、总结与实战要点
PIP-92 为 Pulsar 的地理复制场景补齐了"策略作用域"这一维度,其设计思想可以归纳为三条主线:
- 数据层:
TopicPolicies.isGlobal(默认false)用最小代价区分策略作用域,一份数据结构、两种语义; - 复制层:全局策略通过系统 Topic 复制同步到整个复制组,本地策略留在本集群;当前实现以
__G__消息 key 前缀区分二者(PIP-92 原提案中的 ReplicatorsetFilterFunction方案可作为历史背景理解); - 应用层:严格遵循"本地策略 > 全局策略 > 命名空间策略 > Broker 默认配置"的优先级,本地优先、全局兜底。
实战中最容易混淆的三点需要特别留意:
- 不加
--global(或-g)时,所有策略操作都只影响当前集群,与旧版行为完全一致; - 同一 Topic 的本地策略与全局策略是两份独立数据,可并存、可分别删除,读取时也要明确区分(如 REST 的
isGlobal查询参数、服务的GetType.GLOBAL_ONLY/LOCAL_ONLY); - 删除集群或本地策略不会影响其他集群,因为复制链路中流通的只有全局策略。
如需深入源码继续研究,推荐按以下路径阅读:数据模型 TopicPolicies.java → 服务接口与 key 前缀定义 TopicPoliciesService.java → 系统 Topic 实现与双缓存 SystemTopicBasedTopicPoliciesService.java → 命令行选项 CmdTopicPolicies.java → 优先级验证测试 TopicPoliciesTest.java。
- 消息队列
- 流处理
- 后端
- 微服务
- 消息路由
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar PIP-422 深度解析:全局 Topic 级 replicated clusters 策略与 Topic 级策略删除 API
Apache Pulsar PIP 422 深度解析:全局 Topic 级 replicated clusters 策略与 Topic 级策略删除 API 导读
消息队列流处理后端微服务消息路由AscendNPU-IR架构深度解析:HFusion、HIVM、HACC自研方言与分层解耦设计完全指南
AscendNPU IR架构深度解析:HFusion、HIVM、HACC自研方言与分层解耦设计完全指南 AscendNPU IR 是基于 MLIR(Multi
消息队列流处理后端微服务消息路由Apache Pulsar PIP-284 深度解析:使用 TableView 重构 Topic Policies 底层实现
Apache Pulsar PIP 284 深度解析:使用 TableView 重构 Topic Policies 底层实现 导读 PIP 284(Migrat
消息队列流处理后端微服务消息路由
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考