Apache Pulsar 非持久化消息(Non-persistent Topics)完全指南:内存级 Topic 的配置、管理与源码解析
【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar
本篇技术指南围绕 Apache Pulsar 的非持久化主题(Non-persistent Topics)展开,系统讲解其核心概念、启用与配置方式、命令行管理与客户端使用实战,并结合当前仓库源码(Broker 服务实现、Admin REST 接口与 CLI 工具)深入剖析其底层行为。读完本文,你将掌握如何通过non-persistent://主题名使用纯内存消息投递、如何通过enableNonPersistentTopics等参数调优 Broker,以及非持久化主题在消息丢失、性能与适用场景上的取舍。
非持久化主题是什么
默认情况下,Pulsar 会将所有未确认的消息持久化存储到多个 BookKeeper Bookie 节点上。持久化主题(persistent topics)上的消息数据因此可以在 Broker 重启、订阅者故障切换等场景下存活。
与之相对,Pulsar 也支持非持久化主题(non-persistent topics):这类主题上的消息从不落盘,只存在于内存中。一旦 Broker 宕机或订阅者断开连接,该非持久化主题上所有"在途"消息都会随之丢失,客户端可能因此观察到消息丢失。
非持久化主题的完整名称形如(注意名称中的non-persistent类型标识):
non-persistent://tenant/namespace/topic在 Pulsar 的 Topic 命名体系中,persistent/non-persistent是标识主题类型的两个前缀,默认类型为持久化——如果你不指定类型前缀,主题即视为持久化主题。关于两种主题类型的更高层概念介绍,可参阅 Concepts and Architecture 文档中的 Non-persistent topics 章节。
非持久化主题的适用场景与代价
非持久化主题的价值在于降低消息投递路径上的持久化开销。由于无需写入 BookKeeper、无需等待落盘确认,它更适合对延迟敏感、可以容忍消息丢失的场景,例如:
- 实时监控与指标流(丢失少量采样点可接受);
- 日志流汇聚与过滤管道;
- 高速遥测数据、传感器数据的中转。
需要清醒认识的是其代价:可靠性让位于性能。Broker 重启、订阅者断连都会导致在途消息丢失;同时消息仅存在于内存,Broker 内存压力也会随之上升。因此在选择主题类型时,应结合业务对数据完整性的要求做出取舍,而不是盲目追求低延迟。
启用非持久化主题
要在 Broker 上启用非持久化主题,需要将enableNonPersistentTopics参数设置为true。该参数默认即为true,因此通常情况下无需任何额外操作即可使用非持久化消息。
在集群部署中,该参数位于 conf/broker.conf:
# Enable broker to load non-persistent topics enableNonPersistentTopics=true在 standalone 单机模式下,同样的配置参数位于 conf/standalone.conf:
# Enable broker to load non-persistent topics enableNonPersistentTopics=true如果你希望某个 Broker只提供非持久化主题服务,可以将enablePersistentTopics设为false、同时保持enableNonPersistentTopics=true(见 conf/broker.conf 中的注释:# Enable broker to load persistent topics):
enablePersistentTopics=false enableNonPersistentTopics=true从源码看开关如何生效
enableNonPersistentTopics不仅仅是一个"开关"语义的配置项,它还参与了 Broker 的负载均衡决策与主题创建拦截:
- 在负载均衡层面,
SimpleLoadManagerImpl与ModularLoadManagerImpl都会在负载报告中上报nonPersistentTopicsEnabled字段(见 pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/SimpleLoadManagerImpl.java),负载均衡器据此判断某个 Broker 是否有能力承载非持久化主题,避免把非持久化主题分配到未启用该能力的节点上; - 在主题创建层面,BrokerService.createNonPersistentTopic() 会先检查
isEnableNonPersistentTopics():若为false,直接返回NotAllowedException("Broker is not unable to load non-persistent topic"),拒绝加载非持久化主题;随后通过NonPersistentTopic实例完成主题初始化、命名空间归属检查(checkTopicNsOwnership)与复制检查(checkReplication)。
其他相关 Broker 配置项
除总开关外,Broker 还提供若干与非持久化主题运行时行为相关的参数,这些参数同样同时存在于 conf/broker.conf 与 conf/standalone.conf 中:
| 配置项 | 默认值 | 说明 |
|---|---|---|
enablePersistentTopics | true | 是否允许 Broker 加载持久化主题,可与enableNonPersistentTopics组合实现"仅非持久化"或"仅持久化"部署 |
enableNonPersistentTopics | true | 是否允许 Broker 加载非持久化主题 |
maxConcurrentNonPersistentMessagePerConnection | 1000 | 每条连接上可同时处理的非持久化消息数上限,用于限制单连接的并发消息处理压力(见 conf/broker.conf) |
numWorkerThreadsForNonPersistentTopic | 8(standalone)/ 空(broker,表示由系统决定) | 服务非持久化主题的工作线程数,直接影响并发吞吐(见 conf/standalone.conf 与 conf/broker.conf) |
其中maxConcurrentNonPersistentMessagePerConnection与numWorkerThreadsForNonPersistentTopic是调优高吞吐非持久化场景的关键旋钮:前者防止单连接消息处理过载,后者决定处理线程池规模,两者共同约束了内存中消息流转的并发度。
使用非持久化主题
使用非持久化主题非常简单:无需修改任何 Broker 配置以外的设置,只需在交互时通过主题名加以区分即可。
例如,下面的pulsar-client produce命令会在 standalone 集群中向一个非持久化主题生产一条消息:
$ bin/pulsar-client produce non-persistent://public/default/example-np-topic \ --num-produce 1 \ --messages "This message will be stored only in memory"该命令与普通持久化主题的唯一区别就是主题名中的non-persistent://前缀。生产端无需感知底层存储差异,Pulsar 客户端协议层会依据主题名路由到对应的 Topic 实现。
从管理视角出发的非持久化主题更完整指南,可参阅 admin-api-topics 文档中的 Non-persistent topics 部分。
通过 CLI 管理非持久化主题
非持久化主题可以通过pulsar-admin non-persistent命令族进行管理。在仓库中,该命令族由 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNonPersistentTopics.java 实现(对应 REST 服务端为 pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java)。
常用操作包括:
- 创建分区的非持久化主题:对应
create-partitioned-topic子命令,服务端实现为PUT /non-persistent/{tenant}/{namespace}/{topic}/partitions(见 NonPersistentTopics.java); - 查看主题统计信息:
stats子命令,对应GET /non-persistent/{tenant}/{namespace}/{topic}/stats,返回NonPersistentTopicStats类型的统计结果(见 NonPersistentTopics.java); - 列出命名空间下的非持久化主题:
list子命令,对应GET /non-persistent/{tenant}/{namespace}(见 NonPersistentTopics.java)。
示例命令形式:
# 列出 public/default 命名空间下的非持久化主题 $ bin/pulsar-admin non-persistent list public/default # 查看某个非持久化主题的统计信息 $ bin/pulsar-admin non-persistent stats non-persistent://public/default/example-np-topic # 创建分区的非持久化主题(例如 4 个分区) $ bin/pulsar-admin non-persistent create-partitioned-topic non-persistent://public/default/example-np-topic -p 4需要注意,CmdNonPersistentTopics在仓库中被标记为hidden = true,属于管理后台的隐藏命令族;其核心职责是把non-persistent://tenant/namespace/topic形式的主题名参数做合法性校验(见 CliCommand.validateNonPersistentTopic),再转发给 Admin API 客户端。
与 Pulsar 客户端配合使用
使用非持久化消息时,你的 Pulsar 客户端代码几乎不需要做任何改动——唯一需要保证的是使用正确的主题名,即以non-persistent作为主题类型前缀,例如:
non-persistent://my-tenant/my-namespace/my-topic无论是 Java、C++、Python 还是 Go 客户端,生产者与消费者的创建方式与持久化主题完全一致,主题类型由名称自动推导,客户端 API 无需感知底层是否落盘。这意味着你可以把现有应用中的某个主题从persistent切换为non-persistent(或反之),仅通过修改主题名即可完成,业务代码零改动。
源码视角:非持久化主题的内部实现
为了更深入理解非持久化主题的行为,可以阅读其核心实现类 NonPersistentTopic.java(位于pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/目录)。从源码结构可以推断出以下关键设计:
- 内存驻留:与持久化主题依赖 ManagedLedger + BookKeeper 不同,
NonPersistentTopic将消息存放在内存数据结构中,不经过磁盘写入路径,因此吞吐与延迟表现更优; - 无持久化订阅游标:非持久化主题不支持基于磁盘的游标恢复,订阅者需要重新连接并接受消息回退窗口之外的丢失;
- Broker 生命周期绑定:消息生命周期与 Broker 进程强绑定,进程退出即数据消失,这也解释了为何官方文档明确提示"killing a broker 或 disconnecting a subscriber 会导致在途消息全部丢失"。
此外,非持久化主题同样支持分区(partitioned)形态——Admin 层通过/non-persistent/{tenant}/{namespace}/{topic}/partitions创建分区,分区后的统计信息可通过partitioned-stats接口聚合获取(见 NonPersistentTopics.java)。
小结与最佳实践
| 维度 | 非持久化主题 | 持久化主题 |
|---|---|---|
| 消息存储 | 仅内存,不落盘 | BookKeeper 多副本持久化 |
| 消息可靠性 | Broker 重启 / 订阅者断连即丢失 | 可跨 Broker 重启与订阅者故障切换存活 |
| 延迟与吞吐 | 省去落盘开销,通常更优 | 受持久化路径影响 |
| 主题名前缀 | non-persistent:// | persistent://(默认) |
| 适用场景 | 实时指标、日志管道等可容忍丢失的场景 | 消息关键、需要严格保证投递的业务 |
实践建议:
- 确认开关:生产环境部署前,确认 conf/broker.conf 中
enableNonPersistentTopics=true;standalone 环境检查 conf/standalone.conf; - 明确命名:所有客户端与 CLI 操作都通过
non-persistent://tenant/namespace/topic名称区分主题类型,切勿混用前缀; - 评估丢失容忍度:在引入非持久化主题前,明确业务对"Broker 故障丢消息"的容忍边界,必要时用
maxConcurrentNonPersistentMessagePerConnection与numWorkerThreadsForNonPersistentTopic控制内存中的并发消息量; - 结合监控:使用
pulsar-admin non-persistent stats观察非持久化主题的实时统计,及时掌握消息流量与连接状态。
【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考