PIP-54 深度解读:Apache Pulsar 批量消息按索引级确认(Batch Index Level Acknowledgement)的实现原理与配置指南

📅 发布时间:2026/10/9 4:41:30
PIP-54 深度解读:Apache Pulsar 批量消息按索引级确认(Batch Index Level Acknowledgement)的实现原理与配置指南
消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载导读本文以 Apache Pulsar 官方 PIP-54Proposal for the Improvement of Pulsar为骨架深入剖析批量消息索引级确认acknowledgement at batch index level这一核心能力它解决了 Producer 批量发送batching场景下Broker 无法获知批内单个消息的 ack 状态、从而向消费者重复投递已确认消息的问题。读完本文你将掌握该特性的设计动机、客户端与 Broker 的协作机制、线协议与 managed ledger 持久化格式的字段变更以及acknowledgmentAtBatchIndexLevelEnabled等关键配置的用法与底层实现并能在自己的 Pulsar 集群中正确启用与调优该能力。背景与动机为什么需要按批内索引确认消息在 Pulsar 中Producer 为了提高吞吐默认会将多条消息打包成一个batch批量消息并以(ledgerId, entryId)二元组作为一条 entry 写入 BookKeeper。对于 Broker 端的 managed cursor 而言它维护已确认消息的状态主要依赖两类信息mark delete position标记删除位置表示该位置之前的消息都已被确认individual delete messages单条删除记录用于记录空洞式的部分确认。这两类记录的对象粒度都是 batch 消息的(ledgerId, entryId)而batch 内部单个消息的索引batch index并不可见。这就带来一个实际问题当某个 batch 中只有部分消息被消费者确认时Broker 无法区分批内哪些索引已 ack、哪些未 ack于是它会把整个 batch 中尚未确认甚至已确认的消息重新投递给消费者造成重复消费。PIP-54 的核心目标正是在 Broker 侧维护每个 batch 消息内部各索引的 ack 状态从而避免把已确认的批内消息再次派发给用户。总体设计思路客户端与 Broker 的协作PIP-54 提出的方案需要客户端与 Broker 协同配合整体流程可以概括为下行携带已 ack 索引、上行回报 ack 索引、Broker 维护 BitSet、全 ack 后整批删除Broker 分发时携带 ack 信息当 Broker 向消费者派发 batch 消息时会同时携带该 batch 中已被 ack 的索引集合客户端过滤客户端收到后将其中已被 ack 的 batch index 过滤掉只把尚未确认的消息交给上层应用客户端回报 ack客户端需要把 batch 索引级别的 ack 信息上报给 BrokerBroker 据此维护 batch index 的 ack 状态Broker 内存维护 持久化managed cursor 在内存中使用BitSet记录每个 batch 消息中已被 ack 的索引该 BitSet 可以持久化到 ledger 与 metastore从而避免 Broker 崩溃时状态丢失整批删除当一个 batch 消息的所有索引都被 ack 后cursor 便会删除整个 batch 消息。这一设计对 ack 状态的前进方式做了重要约定Broker 收到 batch index ack 请求后把已 ack 的索引写入 BitSet分发时从 managed cursor 读取该 batch 的 index ack 状态并随消息下发。关键注意点consumer permits可用配额的补偿计算PIP 明确强调了一个极易踩坑的细节由于客户端会过滤掉已 ack 的 batch indexBroker 必须按已 ack 的索引个数相应增加该消费者的可用 permitsavailable permits。因为 Broker 的流控flow control按消息个数发放 permits客户端过滤后实际处理的消息数少于 Broker 发放数。如果 Broker 不补偿这部分的 permits长时间累积后可用 permits 会被已 ack 的索引耗尽Broker 将停止向该消费者派发消息造成订阅停滞。因此 permits 补偿计算是实现该特性时不可忽视的一环。线协议Wire Protocol变更新增ack_set字段为了在客户端与 Broker 之间传递批内已 ack 索引集合PIP-54 对 Pulsar 的二进制线协议做了两处扩展均在PulsarApi.proto中落地1.CommandMessageBroker → 客户端派发消息的命令新增字段message CommandMessage { // ... 原有字段 ... repeated int64 ack_set 4; }ack_set以稀疏区间interval set编码的方式描述批内哪些索引已被 ack从而以很小的体积传输大批量的 ack 索引。2.MessageIdData消息 ID 数据结构新增字段message MessageIdData { // ... 原有字段 ... repeated int64 ack_set 5; }MessageIdData内嵌于CommandMessage见 PulsarApi.proto其中CommandMessage的required MessageIdData message_id 2与repeated int64 ack_set 4相邻定义MessageIdData中的repeated int64 ack_set 5见 PulsarApi.proto两个字段配合使每条派发的 batch 消息都能携带其内部各索引的 ack 状态。Managed Ledger 数据格式变更批量删除索引的持久化为避免 Broker 崩溃后 batch index 的 ack 状态丢失PIP-54 设计了专门的持久化结构。PIP 中给出的三种数据格式如下1. 新增BatchedEntryDeletionIndexInfo记录某个 batch entry 的删除索引集合message BatchedEntryDeletionIndexInfo { required NestedPositionInfo position 1; repeated int64 deleteSet 2; }其中position指向具体的 batch entry 位置ledgerId/entryIddeleteSet记录该批内已删除的索引集合。2.ManagedCursorInfo新增字段batchedEntryDeletionIndexInfo 7message ManagedCursorInfo { // Store which index in the batch message has been deleted repeated BatchedEntryDeletionIndexInfo batchedEntryDeletionIndexInfo 7; }cursor 的持久化元数据中保存一份批级删除索引的列表。3.PositionInfo新增字段batchedEntryDeletionIndexInfo 5message PositionInfo { // Store which index in the batch message has been deleted repeated BatchedEntryDeletionIndexInfo batchedEntryDeletionIndexInfo 5; }这样一来批内索引的 ack 状态既能在内存中以 BitSet 高效维护又能随 cursor 状态写入 ledger 与 metastore实现崩溃恢复。源码级印证ManagedCursorImpl 中的 BitSet 实现PIP 提出的设计在当前仓库中已经完整落地。核心实现位于 ManagedCursorImpl.java其中使用ConcurrentSkipListMapPosition, BitSet batchDeletedIndexes维护每个 batch entry 位置对应的已删除索引 BitSet见 ManagedCursorImpl.java与 PIP 中内存中使用 BitSet 维护 batch index ack 状态的设计完全对应在 ack 上报时向该 BitSet 中累加已 ack 索引并保证只前进不回退源码中明确注释为了防止记录在 batchDeletedIndexes 中的 batch index 回滚仅当上报的 batch index 更大时才更新 batchDeletedIndexes见 ManagedCursorImpl.java当某个 batch 的所有索引都被删除后从batchDeletedIndexes中移除该位置见 ManagedCursorImpl.java对应 PIP 中所有索引 ack 后删除整个 batch 消息的语义相关的恢复逻辑会依据持久化的 BitSet 重建内存状态见 ManagedCursorImpl.java验证了 PIPBitSet 可持久化以避免 Broker 崩溃的设计。同时Broker 模块的测试类 BatchMessageTest.java 对该特性的行为默认开启状态下批量消息的 ack/派发语义进行了覆盖验证。配置说明acknowledgmentAtBatchIndexLevelEnabled与相关参数PIP-54 在broker.conf中新增了开关acknowledgmentAtBatchIndexLevelEnabled。该配置在源码与配置文件中均已落地在 ServiceConfiguration.java 中定义默认值为true即默认开启其FieldContext描述为Whether to enable the acknowledge of batch local index在 conf/broker.conf 中写为acknowledgmentAtBatchIndexLevelEnabledtrue可直接按需修改conf/standalone.conf 中同样包含该特性相关配置项便于单机模式验证。除了总开关仓库还提供了配套的持久化容量控制参数managedLedgerMaxBatchDeletedIndexToPersist默认值200000见 ServiceConfiguration.java含义每个订阅中部分确认的 batch 消息允许被持久化 batch 删除索引的最大数量当该限制被超过时剩余 batch 的删除索引只在内存中跟踪若 Broker 重启或发生负载均衡topic 迁移这些内存中的索引会被清空并随消息重投递redelivery到消费者时重新计算相关参数说明同时出现在 conf/broker.conf 与 conf/standalone.conf 的注释中。实操建议对于批内消息 ack 比例高、重投递频繁或 Broker 重启场景敏感的生产环境建议根据订阅数与 backlog 规模合理调大managedLedgerMaxBatchDeletedIndexToPersist并注意它与managedLedgerMaxUnackedRangesToPersist、managedLedgerPersistIndividualAckAsLongArray等相邻配置在存储占用上的关联持久化体积与 backlog 跨度及 ack 空洞数量相关。总结PIP-54 通过客户端过滤 Broker 维护 BitSet 持久化到 ledger/metastore permits 补偿的组合方案把确认粒度从批量消息(ledgerId, entryId)细化到了批内单条索引从根上避免了批量发送场景下已确认消息被重复派发的问题。该设计在今天仍是 Pulsar 高效批量消费、Exactly-Once 语义与事务消息相关后续 PIP 亦复用ack_set表示法如 PIP-84、PIP-107的基础设施之一。对于使用者在生产环境中落地需要记住三点该特性默认开启、无需额外改造客户端即可获益调优时重点关注managedLedgerMaxBatchDeletedIndexToPersist以平衡持久化开销与故障恢复精度并理解 permits 补偿逻辑避免在流控上出现误判。赞分享消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载相关推荐Apache Pulsar 批量接收消息Batch Receive实战指南从 PIP-38 设计到源码实现Apache Pulsar 批量接收消息Batch Receive实战指南从 PIP 38 设计到源码实现 批量接收Batch Receive是 Ap消息队列流处理后端微服务消息路由Apache Pulsar PIP-330 解析getMessagesById 按消息 ID 获取批量消息中的所有消息Apache Pulsar PIP 330 解析getMessagesById 按消息 ID 获取批量消息中的所有消息 导读 PIP 330Proposal消息队列后端Apache Pulsar PIP-374 深度解读借助 ConsumerInterceptor.onArrival() 实现 receiverQueue 消息可见性Apache Pulsar PIP 374 深度解读借助 ConsumerInterceptor.onArrival 实现 receiverQueue 消息可消息队列流处理后端微服务消息路由上一篇TypeScript 3.0 破坏性变更深度解读unknown 保留关键字与 strictNullChecks 关闭时的交叉类型简化下一篇WSABuilds 完整指南在 Windows 10/11 上从零跑通带谷歌商店和 Root 的安卓环境创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考