- 消息队列
- 流处理
- 后端
- 微服务
- 消息路由
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
Apache Pulsar 的命名空间主题列举(topic listing)逻辑在 PIP-452 中被重新设计:从硬编码扫描元数据存储,演进为支持客户端属性上下文(properties)、可插拔扩展(PulsarResourcesExtended)的灵活架构。本文将以 PIP-452 提案为主线,结合当前仓库中的协议定义、Broker 源码与端到端测试,完整讲解其设计动机、协议变更、配置方式、自定义实现方法及向后兼容与安全边界,帮助读者在 Pulsar 上落地多租户场景下"按属性筛选主题"的能力。
背景与动机:硬编码主题列举的局限
在 PIP-452 之前,Broker 中处理CommandGetTopicsOfNamespace的逻辑是硬编码的:直接扫描元数据存储(如 ZooKeeper)下命名空间节点的全部子节点。这在复杂多租户场景下暴露了两大问题:
- 没有客户端上下文(No Client Context):Broker 无法区分"是谁、因为什么目的在请求主题列表",也就无法基于客户端属性(这些属性可能对应或派生于主题属性)对结果进行过滤;
- 过滤效率低下(Inefficient Filtering):对于包含数百万主题的命名空间,Broker 必须先把完整主题列表加载进内存,再应用
topics_pattern正则做过滤。没有任何机制可以把过滤"下推"(push down)到数据源侧(例如带索引的数据库)。
PIP-452 的解决思路是双管齐下:让协议携带客户端属性,同时把主题列举逻辑改造成可插拔的扩展点。
目标范围:做什么与不做什么
In Scope(范围内)
- 协议:为
CommandGetTopicsOfNamespace增加properties字段,承载客户端侧上下文; - Broker:引入可插拔接口
PulsarResourcesExtended,允许自定义资源管理逻辑,先从 "Get Topics" 请求做起; - Client:更新 Java 客户端,在使用Regex 订阅时把 Consumer 属性转发给 lookup 服务;
- Admin API & CLI:REST API 与命令行支持在列举命名空间主题时传入
properties; - Configuration:新增 Broker 配置项,用于开关主题列表监听(topic list watching)。
Out of Scope(范围外)
- 不支持为自定义主题列举开启主题列表监听(topic list watching),可作为后续工作考虑;
- 若想使用该特性,必须关闭主题列表监听器,即在
broker.conf中设置enableBrokerTopicListWatcher = false。
这一点与测试代码相互印证:在 CustomizedPulsarResourcesExtendedTest.java 的setup()中,测试显式调用conf.setEnableBrokerTopicListWatcher(false)后才配置自定义扩展类。
总体设计:协议携带上下文 + Broker 可插拔资源接口
高层设计上,PIP-452 做了两处修改:
- Pulsar 协议新增
propertiesmap 字段; - Broker 侧新增可插拔接口
PulsarResourcesExtended,其默认实现保留原有行为(委托给NamespaceService)。
请求的处理链路为:Broker 连接处理器(connection handler)→NamespaceService→PulsarResourcesExtended。NamespaceService作为入口,把带属性的列举请求最终导向可插拔扩展,从而为"自定义主题列举策略"留出实现空间。
协议变更:PulsarApi.proto 新增 properties 字段
协议定义位于 PulsarApi.proto,源码中的CommandGetTopicsOfNamespace已完整包含 PIP-452 新增的字段:
message CommandGetTopicsOfNamespace { enum Mode { PERSISTENT = 0; NON_PERSISTENT = 1; ALL = 2; } required uint64 request_id = 1; required string namespace = 2; optional Mode mode = 3 [default = PERSISTENT]; optional string topics_pattern = 4; optional string topics_hash = 5; // Context properties from the client repeated KeyValue properties = 6; }字段说明:
| 字段 | 编号 | 类型 | 说明 |
|---|---|---|---|
request_id | 1 | required uint64 | 请求标识,用于关联响应 |
namespace | 2 | required string | 目标命名空间 |
mode | 3 | optional Mode | 列举模式,默认PERSISTENT,可选NON_PERSISTENT、ALL |
topics_pattern | 4 | optional string | 客户端下发的正则过滤模式(既有字段) |
topics_hash | 5 | optional string | 用于增量同步的主题哈希(既有字段) |
properties | 6 | repeated KeyValue | 新增字段:客户端属性上下文,用于自定义过滤 |
对应的响应消息CommandGetTopicsOfNamespaceResponse(PulsarApi.proto)保留了topics、filtered、topics_hash、changed等既有语义字段,Broker 最终过滤逻辑不受影响。
Broker 配置:两个关键配置项
PIP-452 在 Broker 配置中引入两个参数,其定义与默认值见 ServiceConfiguration.java:
# Enables watching topic add/remove events on broker side. It is separated from enableBrokerSideSubscriptionPatternEvaluation. # 是否在 Broker 侧监听主题新增/删除事件,用于订阅模式评估(默认 true) enableBrokerTopicListWatcher = true # Class name for the extended Pulsar resources. # 扩展资源实现类名,该类必须实现 org.apache.pulsar.broker.PulsarResourcesExtended # 默认实现为 DefaultPulsarResourcesExtended,保留原有行为 pulsarResourcesExtendedClassName = org.apache.pulsar.broker.DefaultPulsarResourcesExtended源码层面的细节值得注意:
pulsarResourcesExtendedClassName的默认值常量定义在 ServiceConfiguration.java,即"org.apache.pulsar.broker.DefaultPulsarResourcesExtended";enableBrokerTopicListWatcher默认值为true(ServiceConfiguration.java),且属于dynamic = false的静态配置,修改需要重启 Broker;- 配置项文档明确说明:旧的
enableBrokerSideSubscriptionPatternEvaluation不再控制主题列表监听,行为已迁移到enableBrokerTopicListWatcher(见 ServiceConfiguration.java)。默认配置文件 broker.conf 中仍保留enableBrokerSideSubscriptionPatternEvaluation=true。
配置注意事项:当前仓库的默认 broker.conf 尚未包含这两个新配置项(它们依靠ServiceConfiguration的默认值生效)。若需要自定义主题列举,应在broker.conf中显式添加pulsarResourcesExtendedClassName,并同时关闭主题列表监听:
enableBrokerTopicListWatcher = false pulsarResourcesExtendedClassName = com.example.MyCustomPulsarResourcesExtended可插拔接口:PulsarResourcesExtended 与默认实现
接口定义
接口源码位于 PulsarResourcesExtended.java,标注为@InterfaceStability.Evolving,共定义三个方法:
package org.apache.pulsar.broker; @InterfaceStability.Evolving public interface PulsarResourcesExtended { CompletableFuture<List<String>> listTopicOfNamespace(NamespaceName namespaceName, CommandGetTopicsOfNamespace.Mode mode, @Nullable Map<String, String> properties); void initialize(PulsarService pulsarService); void close(); }listTopicOfNamespace(...):核心方法,按命名空间 + 模式 + 属性过滤返回主题名列表;properties为 null 或空时不施加过滤;initialize(PulsarService):Broker 启动时调用,为自定义实现注入PulsarService及其依赖;close():Broker 关闭时调用,用于释放自定义实现持有的资源。
默认实现:委托 NamespaceService
DefaultPulsarResourcesExtended.java 是默认实现,其listTopicOfNamespace直接委托给pulsarService.getNamespaceService().getListOfTopics(namespaceName, mode),完全保留 PIP 之前的原生行为:
public class DefaultPulsarResourcesExtended implements PulsarResourcesExtended { @Override public CompletableFuture<List<String>> listTopicOfNamespace(NamespaceName namespaceName, CommandGetTopicsOfNamespace.Mode mode, Map<String, String> properties) { return pulsarService.getNamespaceService().getListOfTopics(namespaceName, mode); } // initialize() 中持有 pulsarService,close() 为空实现 }加载机制
扩展实例在 Broker 启动阶段创建。从 PulsarService.java 可以看到启动流程中调用loadPulsarResourcesExtended(),其实现(PulsarService.java)通过Reflections.createInstance(className, PulsarResourcesExtended.class, ...)按类名反射实例化并调用initialize(this)——因此自定义类必须具备无参构造,且pulsarResourcesExtendedClassName必须指向实现了该接口的类。
NamespaceService 改造:带属性过滤的主题列表入口
NamespaceService是请求处理与可插拔扩展之间的桥接层,源码位于 NamespaceService.java。
新增方法
public CompletableFuture<List<String>> getListOfTopicsByProperties(NamespaceName namespaceName, Mode mode, Map<String, String> properties) { return pulsar.getPulsarResourcesExtended().listTopicOfNamespace(namespaceName, mode, properties); } public CompletableFuture<List<String>> getListOfUserTopicsByProperties(NamespaceName namespaceName, Mode mode, Map<String, String> properties) { return getListOfUserTopicsInternal(cacheKeyWithProperties(namespaceName, mode, properties), () -> getListOfTopicsByProperties(namespaceName, mode, properties)); }getListOfTopicsByProperties:把请求下放给可插拔扩展,是自定义策略的钩子;getListOfUserTopicsByProperties:在扩展之上叠加既有逻辑——按"缓存键"合并并发请求,并对结果执行系统主题过滤。
请求合并与缓存键
getListOfUserTopicsInternal(NamespaceService.java)利用inProgressQueryUserTopics的computeIfAbsent合并同一键的并发查询,只让第一个线程真正发起底层查询,并通过thenApplyAsync(TopicList::filterSystemTopic, pulsar.getExecutor())剔除系统主题。
缓存键由cacheKeyWithProperties(NamespaceService.java)生成:先拼接mode + "://" + namespaceName,若 properties 非空,则按键排序后以|key=value追加——排序保证了等价属性集合产生相同缓存键,避免缓存碎片化。
连接处理器调用更新
Broker 侧处理CommandGetTopicsOfNamespace时,从原来的getListOfUserTopics改为getListOfUserTopicsByProperties(PIP-452 中internalHandleGetTopicsOfNamespace片段):
return getBrokerService().pulsar().getNamespaceService() .getListOfUserTopicsByProperties(namespaceName, mode, properties);同时该路径仍受maxTopicListInFlightLimiter(AsyncDualMemoryLimiter)的堆内存配额限流保护,说明带属性的列举与普通列举一样受 Broker 的 in-flight 主题列表内存限制约束。
客户端变更:Regex 订阅转发 Consumer 属性
LookupService 接口
Java 客户端内部LookupService接口更新签名,接受属性 map:
CompletableFuture<List<String>> getTopicsUnderNamespace( NamespaceName namespace, Mode mode, String topicsPattern, String topicsHash, Map<String, String> properties );PatternMultiTopicsConsumerImpl
正则订阅消费者实现PatternMultiTopicsConsumerImpl会从ConsumerConfigurationData中提取属性并传给LookupService:
ConsumerConfigurationData conf; Map<String, String> contextProperties = conf.getProperties(); lookup.getTopicsUnderNamespace( namespace, mode, topicsPattern.pattern(), topicsHash, contextProperties // Pass consumer properties here ).thenAccept(topics -> { // ... update subscriptions ... });这意味着:只要你在创建正则订阅 Consumer 时设置了.properties(...),这些属性就会随getTopicsUnderNamespace请求带到 Broker,供自定义PulsarResourcesExtended实现做属性匹配——"客户端属性对应或派生于主题属性"的假设由此闭环。
Admin API 与 CLI:按属性列举主题
PIP-452 为 REST API 和pulsar-admin增加了properties参数(注意 URL 中的值需要编码):
GET /admin/v2/persistent/{tenant}/{namespace}?properties=k1=v1,k2=v2pulsar-admin topics list <tenant>/<namespace> -p k1=v1 -p k2=v2测试代码 CustomizedPulsarResourcesExtendedTest.java 展示了 Admin API 侧的用法——通过ListTopicsOptions.builder().properties(propsFromClient).build()传入属性过滤条件:
List<String> list = admin.topics().getList("public/default", TopicDomain.persistent, ListTopicsOptions.builder().properties(propsFromClient).build());端到端验证:从自定义实现到测试用例
测试用例解读
CustomizedPulsarResourcesExtendedTest.java 完整演示了该特性的落地路径:
- 在
setup()中配置setEnableBrokerTopicListWatcher(false)和setPulsarResourcesExtendedClassName(CustomizedPulsarResourcesExtended.class.getName())(L44-L51); - 创建分区主题并给若干主题打上自定义属性(
setCustomProperties)或通过admin.topics().updateProperties(...)更新真实主题属性(L92-L96); - 通过 Admin API 传入
properties = {env:prod, region:us-west}列举主题,断言只返回带这些属性的主题(含分区-partition-0/1/2); - 用
.topicsPattern("persistent://public/default/.*").properties(propsFromClient)创建 Regex Consumer,断言其实际订阅的主题集合与属性过滤结果一致(L115-L124)。
该测试同时覆盖了Admin API 路径与Regex Consumer 路径,且在同一命名空间混有"无属性主题"(test-topic4的属性为env=test)以验证过滤确实生效。
自定义实现参考
测试配套的自定义实现 CustomizedPulsarResourcesExtended.java 是一个极佳的模板:它继承DefaultPulsarResourcesExtended,仅覆盖listTopicOfNamespace:
@Override public CompletableFuture<List<String>> listTopicOfNamespace(NamespaceName namespaceName, CommandGetTopicsOfNamespace.Mode mode, @Nullable Map<String, String> properties) { if (MapUtils.isEmpty(properties)) { return super.listTopicOfNamespace(namespaceName, mode, properties); } if (enabledTopicWatcher) { return CompletableFuture.failedFuture(new IllegalStateException( "Customized topic listing with properties is not supported when broker topic watcher is enabled.")); } List<String> list = queryTopicListByProperties(namespaceName.toString(), properties); // ... 返回过滤结果 }要点:
- 无属性请求直接走默认逻辑(委托原生
NamespaceService),只有携带属性时才走自定义查询; - 自定义查询时若开启了 topic watcher 会显式报错,与 PIP 的 Out of Scope 声明一致;
- 自定义实现内部可以用
Map<Namespace, Map<Topic, Map<PropertyKey, PropertyValue>>>这样的内存映射模拟"数据库索引",生产环境中完全可以把这里的查询替换为对带索引存储(如数据库)的访问,实现真正的过滤下推。
向后兼容与安全考量
向后兼容
- 协议层:
properties是 Protobuf 的 optional/repeated 字段,旧客户端不会发送该字段,旧 Broker 会忽略它,属于非破坏性变更; - 行为层:默认策略与现有行为完全一致(
DefaultPulsarResourcesExtended委托原生NamespaceService);且Broker 始终保留最终过滤逻辑——即使自定义策略返回了主题,Broker 仍会应用客户端请求的topics_pattern正则,确保自定义策略不可能返回违反客户端模式的主题。
安全考量
- 输入校验:
propertiesmap 是用户可控输入。PIP 明确要求自定义实现必须校验与清洗这些输入,尤其当它们被用于构造数据库查询时,需防范注入风险; - 授权边界:本 PIP 只控制主题的发现(列举),并不会绕过订阅/生产这些主题时的 Authorization Service——对主题的读写权限校验依然由既有的授权体系负责。
小结
PIP-452 以"协议携带属性 + Broker 可插拔扩展"的组合拳,把 Apache Pulsar 命名空间主题列举从硬编码扫描升级为可定制、可下推过滤的能力,同时通过默认实现与 Broker 最终过滤保证 100% 向后兼容。对开发者而言,落地路径清晰:关闭enableBrokerTopicListWatcher→ 实现并注册PulsarResourcesExtended→ 在 Regex Consumer 或 Admin API 中携带属性即可。相关源码与测试均可在仓库中查阅:协议定义见 PulsarApi.proto,接口与默认实现见 PulsarResourcesExtended.java 与 DefaultPulsarResourcesExtended.java,入口逻辑见 NamespaceService.java,端到端示例见 CustomizedPulsarResourcesExtendedTest.java。
- 消息队列
- 流处理
- 后端
- 微服务
- 消息路由
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 可插拔 Entry Filter 机制解析:基于 PIP-105 的 Dispatcher 消息过滤框架
Apache Pulsar 可插拔 Entry Filter 机制解析:基于 PIP 105 的 Dispatcher 消息过滤框架 导读 PIP 105(Su
消息队列流处理后端微服务消息路由Apache Pulsar 命名空间复制职责拆分:PIP-321 与 `allowed-clusters` 机制深度解析
Apache Pulsar 命名空间复制职责拆分:PIP 321 与 allowed clusters 机制深度解析 PIP 321(Split the res
消息队列后端Apache Pulsar PIP-191 深度解析:基于属性分组的批量消息 Entry Filter 过滤方案
Apache Pulsar PIP 191 深度解析:基于属性分组的批量消息 Entry Filter 过滤方案 导读 PIP 191 是 Apache Pul
消息队列流处理后端微服务消息路由
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考