Apache Pulsar PIP-452 深度解析:基于属性过滤的可插拔命名空间主题列表机制 消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache 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下命名空间节点的全部子节点。这在复杂多租户场景下暴露了两大问题没有客户端上下文No Client ContextBroker 无法区分是谁、因为什么目的在请求主题列表也就无法基于客户端属性这些属性可能对应或派生于主题属性对结果进行过滤过滤效率低下Inefficient Filtering对于包含数百万主题的命名空间Broker 必须先把完整主题列表加载进内存再应用topics_pattern正则做过滤。没有任何机制可以把过滤下推push down到数据源侧例如带索引的数据库。PIP-452 的解决思路是双管齐下让协议携带客户端属性同时把主题列举逻辑改造成可插拔的扩展点。目标范围做什么与不做什么In Scope范围内协议为CommandGetTopicsOfNamespace增加properties字段承载客户端侧上下文Broker引入可插拔接口PulsarResourcesExtended允许自定义资源管理逻辑先从 Get Topics 请求做起Client更新 Java 客户端在使用Regex 订阅时把 Consumer 属性转发给 lookup 服务Admin API CLIREST API 与命令行支持在列举命名空间主题时传入propertiesConfiguration新增 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_id1required uint64请求标识用于关联响应namespace2required string目标命名空间mode3optional Mode列举模式默认PERSISTENT可选NON_PERSISTENT、ALLtopics_pattern4optional string客户端下发的正则过滤模式既有字段topics_hash5optional string用于增量同步的主题哈希既有字段properties6repeated KeyValue新增字段客户端属性上下文用于自定义过滤对应的响应消息CommandGetTopicsOfNamespaceResponsePulsarApi.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.DefaultPulsarResourcesExtendedenableBrokerTopicListWatcher默认值为trueServiceConfiguration.java且属于dynamic false的静态配置修改需要重启 Broker配置项文档明确说明旧的enableBrokerSideSubscriptionPatternEvaluation不再控制主题列表监听行为已迁移到enableBrokerTopicListWatcher见 ServiceConfiguration.java。默认配置文件 broker.conf 中仍保留enableBrokerSideSubscriptionPatternEvaluationtrue。配置注意事项当前仓库的默认 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 { CompletableFutureListString listTopicOfNamespace(NamespaceName namespaceName, CommandGetTopicsOfNamespace.Mode mode, Nullable MapString, String properties); void initialize(PulsarService pulsarService); void close(); }listTopicOfNamespace(...)核心方法按命名空间 模式 属性过滤返回主题名列表properties为 null 或空时不施加过滤initialize(PulsarService)Broker 启动时调用为自定义实现注入PulsarService及其依赖close()Broker 关闭时调用用于释放自定义实现持有的资源。默认实现委托 NamespaceServiceDefaultPulsarResourcesExtended.java 是默认实现其listTopicOfNamespace直接委托给pulsarService.getNamespaceService().getListOfTopics(namespaceName, mode)完全保留 PIP 之前的原生行为public class DefaultPulsarResourcesExtended implements PulsarResourcesExtended { Override public CompletableFutureListString listTopicOfNamespace(NamespaceName namespaceName, CommandGetTopicsOfNamespace.Mode mode, MapString, String properties) { return pulsarService.getNamespaceService().getListOfTopics(namespaceName, mode); } // initialize() 中持有 pulsarServiceclose() 为空实现 }加载机制扩展实例在 Broker 启动阶段创建。从 PulsarService.java 可以看到启动流程中调用loadPulsarResourcesExtended()其实现PulsarService.java通过Reflections.createInstance(className, PulsarResourcesExtended.class, ...)按类名反射实例化并调用initialize(this)——因此自定义类必须具备无参构造且pulsarResourcesExtendedClassName必须指向实现了该接口的类。NamespaceService 改造带属性过滤的主题列表入口NamespaceService是请求处理与可插拔扩展之间的桥接层源码位于 NamespaceService.java。新增方法public CompletableFutureListString getListOfTopicsByProperties(NamespaceName namespaceName, Mode mode, MapString, String properties) { return pulsar.getPulsarResourcesExtended().listTopicOfNamespace(namespaceName, mode, properties); } public CompletableFutureListString getListOfUserTopicsByProperties(NamespaceName namespaceName, Mode mode, MapString, String properties) { return getListOfUserTopicsInternal(cacheKeyWithProperties(namespaceName, mode, properties), () - getListOfTopicsByProperties(namespaceName, mode, properties)); }getListOfTopicsByProperties把请求下放给可插拔扩展是自定义策略的钩子getListOfUserTopicsByProperties在扩展之上叠加既有逻辑——按缓存键合并并发请求并对结果执行系统主题过滤。请求合并与缓存键getListOfUserTopicsInternalNamespaceService.java利用inProgressQueryUserTopics的computeIfAbsent合并同一键的并发查询只让第一个线程真正发起底层查询并通过thenApplyAsync(TopicList::filterSystemTopic, pulsar.getExecutor())剔除系统主题。缓存键由cacheKeyWithPropertiesNamespaceService.java生成先拼接mode :// namespaceName若 properties 非空则按键排序后以|keyvalue追加——排序保证了等价属性集合产生相同缓存键避免缓存碎片化。连接处理器调用更新Broker 侧处理CommandGetTopicsOfNamespace时从原来的getListOfUserTopics改为getListOfUserTopicsByPropertiesPIP-452 中internalHandleGetTopicsOfNamespace片段return getBrokerService().pulsar().getNamespaceService() .getListOfUserTopicsByProperties(namespaceName, mode, properties);同时该路径仍受maxTopicListInFlightLimiterAsyncDualMemoryLimiter的堆内存配额限流保护说明带属性的列举与普通列举一样受 Broker 的 in-flight 主题列表内存限制约束。客户端变更Regex 订阅转发 Consumer 属性LookupService 接口Java 客户端内部LookupService接口更新签名接受属性 mapCompletableFutureListString getTopicsUnderNamespace( NamespaceName namespace, Mode mode, String topicsPattern, String topicsHash, MapString, String properties );PatternMultiTopicsConsumerImpl正则订阅消费者实现PatternMultiTopicsConsumerImpl会从ConsumerConfigurationData中提取属性并传给LookupServiceConsumerConfigurationData conf; MapString, 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}?propertiesk1v1,k2v2pulsar-admin topics list tenant/namespace -p k1v1 -p k2v2测试代码 CustomizedPulsarResourcesExtendedTest.java 展示了 Admin API 侧的用法——通过ListTopicsOptions.builder().properties(propsFromClient).build()传入属性过滤条件ListString 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的属性为envtest以验证过滤确实生效。自定义实现参考测试配套的自定义实现 CustomizedPulsarResourcesExtended.java 是一个极佳的模板它继承DefaultPulsarResourcesExtended仅覆盖listTopicOfNamespaceOverride public CompletableFutureListString listTopicOfNamespace(NamespaceName namespaceName, CommandGetTopicsOfNamespace.Mode mode, Nullable MapString, 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.)); } ListString list queryTopicListByProperties(namespaceName.toString(), properties); // ... 返回过滤结果 }要点无属性请求直接走默认逻辑委托原生NamespaceService只有携带属性时才走自定义查询自定义查询时若开启了 topic watcher 会显式报错与 PIP 的 Out of Scope 声明一致自定义实现内部可以用MapNamespace, MapTopic, MapPropertyKey, 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。赞分享消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载相关推荐Apache Pulsar 可插拔 Entry Filter 机制解析基于 PIP-105 的 Dispatcher 消息过滤框架Apache Pulsar 可插拔 Entry Filter 机制解析基于 PIP 105 的 Dispatcher 消息过滤框架 导读 PIP 105Su消息队列流处理后端微服务消息路由Apache Pulsar 命名空间复制职责拆分PIP-321 与 allowed-clusters 机制深度解析Apache Pulsar 命名空间复制职责拆分PIP 321 与 allowed clusters 机制深度解析 PIP 321Split the res消息队列后端Apache Pulsar PIP-191 深度解析基于属性分组的批量消息 Entry Filter 过滤方案Apache Pulsar PIP 191 深度解析基于属性分组的批量消息 Entry Filter 过滤方案 导读 PIP 191 是 Apache Pul消息队列流处理后端微服务消息路由上一篇淘宝淘金币自动化脚本每天节省20分钟的终极解决方案下一篇LMDeploy 部署 CogVLM / CogVLM2 多模态模型模型准备、离线推理与 PyTorch 引擎实现剖析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考