消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载非持久化主题Non-persistent Topics是 Apache Pulsar 中一种消息数据永不落盘、只驻留内存的主题类型适用于对消息丢失容忍度高、追求更低端到端延迟的实时流式场景。本文以 Pulsar 仓库中的官方 cookbook 为骨架结合源码实现完整讲解非持久化主题的命名规则、Broker 启用配置、pulsar-client生产消息、pulsar-admin管理命令以及底层消息分发机制读完即可在独立集群或 standalone 模式下正确使用并评估非持久化主题。什么是非持久化主题默认情况下Pulsar 会将所有未被确认的消息持久化存储在多个 BookKeeper bookie存储节点上因此持久化主题上的数据可以在 Broker 重启和订阅者故障切换后幸存。而非持久化主题的消息数据从不持久化到磁盘只存在于内存之中。这意味着如果杀掉一个 Pulsar broker或某个订阅者与该主题断开连接该主题上所有传输中的消息都会丢失客户端可能观察到消息丢失由于省略了 BookKeeper 持久化写入环节非持久化主题在部分场景下消息投递会略快于持久化主题但同时也会失去 Pulsar 持久化存储带来的可靠性保障。非持久化主题的完整概念说明见 概念与架构 文档。主题命名规则非持久化主题的名称与持久化主题结构相同只是以non-persistent作为主题类型标识non-persistent://tenant/namespace/topic对照示例持久化主题persistent://public/default/my-topic非持久化主题non-persistent://public/default/my-topic从源码看主题类型判断正是通过名称中的域名部分完成的NonPersistentTopic.isPersistent()直接返回false见 NonPersistentTopic.javaBroker 据此将主题路由到非持久化主题的创建与分发路径。启用非持久化主题使用非持久化主题前必须确保 Broker 配置中已启用它。Broker 通过enableNonPersistentTopics参数控制配置项说明默认值enableNonPersistentTopics是否允许 Broker 加载非持久化主题trueenablePersistentTopics是否允许 Broker 加载持久化主题true这两个开关的定义位于 ServiceConfiguration.java默认均为true。因此开箱即用一般无需额外操作即可使用非持久化主题。集群broker模式配置在集群模式下编辑 broker.conf第 443-453 行附近# Max concurrent non-persistent message can be processed per connection maxConcurrentNonPersistentMessagePerConnection1000 # Number of worker threads to serve non-persistent topic numWorkerThreadsForNonPersistentTopic # Enable broker to load persistent topics enablePersistentTopicstrue # Enable broker to load non-persistent topics enableNonPersistentTopicstruestandalone 模式配置如果以 standalone 模式运行同样的参数位于 standalone.conf第 276-286 行# Max concurrent non-persistent message can be processed per connection maxConcurrentNonPersistentMessagePerConnection1000 # Number of worker threads to serve non-persistent topic numWorkerThreadsForNonPersistentTopic8 # Enable broker to load persistent topics enablePersistentTopicstrue # Enable broker to load non-persistent topics enableNonPersistentTopicstrue完整参数说明见 reference-configuration.md。只启用非持久化主题如果希望某个 Broker仅承载非持久化主题可以关闭持久化主题、保持非持久化主题开启enablePersistentTopicsfalse enableNonPersistentTopicstrueBroker 在创建主题时会检查这一开关。从 BrokerService.createNonPersistentTopic 的源码可以看到当isEnableNonPersistentTopics()为false时创建请求会直接以NotAllowedException(Broker is not unable to load non-persistent topic)失败返回从而保证配置与实际行为一致。相关性能配置项除了启用开关还有两个与非持久化主题吞吐相关的配置同样定义于 ServiceConfiguration.javamaxConcurrentNonPersistentMessagePerConnection每条连接上最多可同时处理的非持久化消息数默认1000用于限制单连接上的并发消息处理量防止内存压力过大numWorkerThreadsForNonPersistentTopic服务非持久化主题的工作线程数默认取Runtime.getRuntime().availableProcessors()即 CPU 核数。standalone 配置中显式给出了8集群 broker.conf 中留空则由代码按 CPU 核数自动决定。这些线程通过 BrokerService 内部numThreads(...)构造专门用于非持久化主题的消息服务。使用非持久化主题使用非持久化主题与持久化主题几乎完全一样唯一的区别是主题名称中的类型标识必须为non-persistent。下面这条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--num-produce 1仅生产 1 条消息--messages指定消息内容。该命令的完整选项说明见 reference-cli-tools.md。使用 Pulsar 客户端Producer / Consumer使用 Java、C、Python 等 Pulsar 客户端连接非持久化主题时无需对客户端做任何修改只需在创建 producer / consumer 时传入以non-persistent开头、格式正确的主题名称即可其余 API 用法与持久化主题一致。值得注意的一点非持久化主题同样支持 Pulsar 的三种订阅类型——独占exclusive、共享shared和故障转移failover见 concepts-messaging.md 中关于非持久化主题的说明。使用 pulsar-admin 管理非持久化主题非持久化主题的管理操作通过pulsar-admin non-persistent命令行接口完成主要包括创建分区的非持久化主题create partitioned non-persistent topic获取主题统计信息stats列出某个命名空间下的非持久化主题list以及查看 internal stats、unload 主题等操作。从管理端源码v2/NonPersistentTopics.java可以看到REST 管理端点统一挂载在/admin/v2/non-persistent路径下提供以下能力HTTP 方法与路径功能GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/partitions获取分区元数据GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/stats获取主题统计信息GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/internalStats获取内部统计信息PUT /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/partitions创建分区的非持久化主题GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/partitioned-stats获取分区主题统计信息GET /admin/v2/non-persistent/{tenant}/{namespace}列出命名空间下的非持久化主题GET /admin/v2/non-persistent/{tenant}/{namespace}/{bundle}按 bundle 列出非持久化主题PUT /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/unload卸载主题DELETE /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/truncate清空主题另外pulsar-admin topics子命令同样适用于非持久化主题——只要在主题名称前显式声明类型即可例如pulsar-admin topics create-partitioned-topic non-persistent://tenant/namespace/topic参见 reference-pulsar-admin.md。从管理员视角更完整的操作指引见 Non-persistent topics 管理指南。底层实现原理源码剖析消息发布不写 ledger直接内存分发非持久化主题的核心实现在 NonPersistentTopic.java 的publishMessage方法第 177-210 行。其流程是校验消息大小超过最大值则拒绝立即以(ledgerId0, entryId0)完成发布回调——没有任何磁盘写入对每个订阅NonPersistentSubscription将消息以retainedDuplicate()复制一份构造内存Entry直接交给订阅的 dispatcher 分发如果存在跨集群复制器NonPersistentReplicator同样通过内存复制转发。也就是说Broker 收到消息后不经过 ManagedLedger、不写 BookKeeper而是直接向当前在线的订阅者推送。这正是仅存内存、延迟更低这一特性的来源代价是任何在线消费者缺席都会造成消息丢失。没有消费者时消息被丢弃对于共享订阅场景分发逻辑在 NonPersistentDispatcherMultipleConsumers.java 的sendMessages方法第 185-204 行当所有消费者都没有可用 permittotalAvailablePermits 0时消息不会排队等待而是直接记录到msgDrop消息丢弃率统计并释放。这意味着非持久化主题没有积压队列——消费者不在线或消费速度跟不上消息就会立即被丢弃。重投与确认机制的简化由于消息不在任何地方保存非持久化主题的重投redelivery机制也被禁用该 dispatcher 使用RedeliveryTrackerDisabled见上文源码第 58-66 行并且订阅数据只保存在内存中incrementTopicEpoch等元数据操作也仅停留在内存层面。不支持的操作从源码可以明确看到非持久化主题对以下能力明确抛出UnsupportedOperationException或以 no-op 方式跳过NonPersistentTopic.java能力非持久化主题上的行为消息去重message deduplicationcheckMessageDeduplicationInfo为 no-op消息过期message expirycheckMessageExpiry为 no-op积压配额backlog quotagetBacklogQuota抛 UnsupportedOperationException获取最新位置 / 最新消息 IDgetLastPosition、getLastMessageId抛 UnsupportedOperationException事务消息publishTxnMessage、endTxn抛 UnsupportedOperationException截断主题truncate抛 NotAllowedException压缩读取readCompacted订阅时报readCompacted only valid on persistent topics静态日志 / 游标元数据无持久化元数据全部在内存这些限制与非持久化的设计目标一致数据不落盘自然无法支持任何依赖历史数据或持久元数据的机制。因此在选型时需要确认业务是否依赖上述能力。适用场景与使用建议适合实时指标流、传感器数据、日志转发等允许丢失部分数据的场景追求极致低延迟、且消费端始终在线或有快速降级方案的管道。不适合需要消息可靠投递、支持消费者离线补拉历史消息、依赖去重/事务/过期清理的业务——这些场景应继续使用持久化主题。使用前请务必确认业务是否可以承受 Broker 重启、消费者断连带来的消息丢失消费者的消费速率是否能跟上生产速率避免因无积压缓冲导致大面积丢弃Broker 内存资源是否充足——非持久化主题的所有消息都驻留在 Broker 堆内/堆外内存中。关于非持久化主题的更多背景与订阅类型细节可继续阅读 概念与架构 与 reference-configuration.md 两篇文档。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 非持久化主题Non-persistent Topics完整指南概念、配置、CLI 与源码原理Apache Pulsar 非持久化主题Non persistent Topics完整指南概念、配置、CLI 与源码原理 非持久化主题non persi消息队列后端流处理Apache Pulsar 非持久化消息Non-persistent Topics完全指南内存级 Topic 的配置、管理与源码解析Apache Pulsar 非持久化消息Non persistent Topics完全指南内存级 Topic 的配置、管理与源码解析 本篇技术指南围绕 A消息队列后端流处理Apache Pulsar 主题Topics管理完全指南Admin API、pulsar-admin 与分区路由实战Apache Pulsar 主题Topics管理完全指南Admin API、pulsar admin 与分区路由实战 本篇技术指南以 Apache Pul消息队列后端流处理上一篇netscan轻量级网络探测的高效解决方案下一篇解决跨软件模型传输难题的GoB让3D资产互导效率提升5倍创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考