news 2026/10/9 1:27:42

Apache Pulsar PIP-452 深度解析:基于属性过滤的可插拔命名空间主题列表机制

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar PIP-452 深度解析:基于属性过滤的可插拔命名空间主题列表机制
  • 消息队列
  • 流处理
  • 后端
  • 微服务
  • 消息路由

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pu/pulsar
点击查看免费下载

Apache Pulsar 的命名空间主题列举(topic listing)逻辑在 PIP-452 中被重新设计:从硬编码扫描元数据存储,演进为支持客户端属性上下文(properties)、可插拔扩展(PulsarResourcesExtended)的灵活架构。本文将以 PIP-452 提案为主线,结合当前仓库中的协议定义、Broker 源码与端到端测试,完整讲解其设计动机、协议变更、配置方式、自定义实现方法及向后兼容与安全边界,帮助读者在 Pulsar 上落地多租户场景下"按属性筛选主题"的能力。

背景与动机:硬编码主题列举的局限

在 PIP-452 之前,Broker 中处理CommandGetTopicsOfNamespace的逻辑是硬编码的:直接扫描元数据存储(如 ZooKeeper)下命名空间节点的全部子节点。这在复杂多租户场景下暴露了两大问题:

  1. 没有客户端上下文(No Client Context):Broker 无法区分"是谁、因为什么目的在请求主题列表",也就无法基于客户端属性(这些属性可能对应或派生于主题属性)对结果进行过滤;
  2. 过滤效率低下(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 做了两处修改:

  1. Pulsar 协议新增propertiesmap 字段;
  2. 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_id1required uint64请求标识,用于关联响应
namespace2required string目标命名空间
mode3optional Mode列举模式,默认PERSISTENT,可选NON_PERSISTENT、ALL
topics_pattern4optional string客户端下发的正则过滤模式(既有字段)
topics_hash5optional string用于增量同步的主题哈希(既有字段)
properties6repeated 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=v2
pulsar-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 完整演示了该特性的落地路径:

  1. 在setup()中配置setEnableBrokerTopicListWatcher(false)和setPulsarResourcesExtendedClassName(CustomizedPulsarResourcesExtended.class.getName())(L44-L51);
  2. 创建分区主题并给若干主题打上自定义属性(setCustomProperties)或通过admin.topics().updateProperties(...)更新真实主题属性(L92-L96);
  3. 通过 Admin API 传入properties = {env:prod, region:us-west}列举主题,断言只返回带这些属性的主题(含分区-partition-0/1/2);
  4. 用.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

项目地址:https://gitcode.com/gh_mirrors/pu/pulsar
点击查看免费下载
上一篇:淘宝淘金币自动化脚本:每天节省20分钟的终极解决方案
下一篇:LMDeploy 部署 CogVLM / CogVLM2 多模态模型:模型准备、离线推理与 PyTorch 引擎实现剖析

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/9 1:27:33

8个可验证的ChatGPT写作指令策略系统

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/9 1:26:33

题解:洛谷 AT_abc451_a [ABC451A] illegal(废)

本文分享的必刷题目是从蓝桥云课、洛谷、AcWing等知名刷题平台精心挑选而来,并结合各平台提供的算法标签和难度等级进行了系统分类。题目涵盖了从基础到进阶的多种算法和数据结构,旨在为不同阶段的编程学习者提供一条清晰、平稳的学习提升路径。 欢迎大家订阅我的专栏:算法…

作者头像 李华
网站建设 2026/10/9 1:26:09

SVM鸢尾花分类实战:基于sklearn的机器学习作业源码解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/9 1:24:29

校园局域网VLAN间路由与NAT实战配置指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/9 1:23:43

PhysX仿真揭秘:fetchResults背后的等待与交付

先想象一个场景: 游戏主线程把“这一帧的物理计算”交给施工队,然后去处理其他工作。等需要使用物理结果时,再回来验收。 在 PhysX 中: scene->simulate(dt); // 开工 // 主线程做其他可以安全并行的事情 scene->fetchResults(true); // 等待并完成验收但这…

作者头像 李华