Apache Pulsar 存储与订阅模型深度实战:从 BookKeeper Ledger 条带化、Quorum 写入到分层存储的工程全解
在消息队列的选型讨论里,Pulsar 常被简化成一句"Kafka 的计算存储分离版本"。这个说法不算错,但它掩盖了真正有意思的部分:Pulsar 把"消息流"这件事拆成了两个彼此独立的问题——如何持久化一段只追加的字节序列,以及如何在其上叠加多个彼此隔离的消费进度。前者交给 Apache BookKeeper,后者由 Broker 侧的 ManagedCursor 承担。理解这条分界线,是判断 Pulsar 是否适合你的负载、以及线上出问题时该调哪个旋钮的前提。
一、为什么 Kafka 的分区模型会碰到天花板
Kafka 的 topic 分区本质上是一个绑定在特定 broker 上的本地日志目录。这份"绑定"带来了极高的吞吐——顺序写、page cache 接管、零拷贝 sendfile 一条龙——但也带来几个结构性的约束:
- 扩容必须搬数据。新增 broker 后要做分区重分配,跨节点复制数十 GB 的 segment 文件,期间占用大量磁盘与网络带宽。
- 分区数有实际上限。分区越多,leader 选举的元数据操作、文件句柄数、以及 follower 拉取的 fan-out 成本都以超线性方式增长。单集群数千分区尚可,数万分区就开始难受。
- 多消费者共享同一份顺序日志。Kafka 的 offset 是分区级别的单一游标,消费者组各自保存 offset,但底层的 retention 策略(时间/大小)无法感知"某条消息是否已被所有消费者消费完"。
Pulsar 的解法是把 broker 变成无状态的计算节点,把持久化下沉到 BookKeeper 集群。Broker 只负责协议处理、订阅游标推进和缓存,所有消息字节都写在 Bookie 上。于是一个 Pulsar 分区不再是"某台机器上的一个目录",而是一条逻辑上的Managed Ledger。
二、三层映射:Topic → Managed Ledger → Ledger
这是理解 Pulsar 存储的第一把钥匙:
Topic (persistent://tenant/ns/topic)
└── Partition-0 → ManagedLedger(逻辑日志,无限长)
├── Ledger 42 (fragment 0: entries 0 – 49999)
├── Ledger 43 (fragment 0: entries 0 – 49999)
│ (fragment 1: entries 50000– 51999) ← ensemble 变更
└── Ledger 44 (正在写入)
ManagedLedger 是逻辑概念,没有物理实体;它是若干个 BookKeeper Ledger 的有序链表。每个 Ledger 是一段有界的、只能追加的日志,达到阈值(managedLedgerMaxEntriesPerLedger,默认 50000 条)或超时(managedLedgerMinLedgerRolloverTimeMinutes)后就滚动(rollover)出新的 Ledger。
只有最后一个 Ledger 是可写的。历史 Ledger 一旦关闭就不可变——这个性质是后面所有删除、分层存储、游标回放逻辑的基础。
滚动之外还有一次强制切换:Bookie 宕机导致 ensemble 变更时,当前 Ledger 会被关闭,新开一个 fragment 用新的 ensemble 继续写。所以一个 Ledger 内部可能有多个 fragment,每个 fragment 记录着自己的 ensemble 列表。元数据里存的是 fragment 数组,读取时按 entry ID 定位到对应 fragment。
三、BookKeeper 写入:条带化 + Quorum
BookKeeper 的数据模型只有四个概念,但组合起来非常精巧:
| 概念 | 含义 |
|---|---|
| Entry | 一条记录,带单调递增的 entry ID |
| Ledger | entry 的有序序列,写完后 seal 变为不可变 |
| Fragment | Ledger 内使用同一套 ensemble 的一段连续 entry |
| Ensemble / Qw / Qa | 参与该 fragment 的 Bookie 集合 / 每条 entry 写几份 / 等几份确认即返回 |
关键在于 Qw ≤ ensemble size。默认配置 ensembleSize=3, writeQuorum=3, ackQuorum=2 表示:选 3 个 Bookie 组成 ensemble,每条 entry 都写满这 3 份,但只需等 2 个确认就向客户端返回。
而更常见的生产配置是 ensemble > writeQuorum 的条带化模式,例如 ensemble=5, Qw=3, Qa=2:
// ensemble=5, Qw=3, Qa=2 时的 entry 分布(轮转条带)
// entry 0 → bookie[0,1,2]
// entry 1 → bookie[1,2,3]
// entry 2 → bookie[2,3,4]
// entry 3 → bookie[3,4,0]
// entry 4 → bookie[4,0,1]
int ensembleSize = 5, writeQuorum = 3;
for (long entryId = 0; entryId < 10; entryId++) {
for (int i = 0; i < writeQuorum; i++) {
int idx = (int) ((entryId + i) % ensembleSize);
System.out.printf("entry %d -> bookie %d%n", entryId, idx);
}
}
这样配置的收益是吞吐横向扩展:写 N 条 entry 的 I/O 被摊到 5 台机器上,而不是 3 台。代价是读取要能从多个 Bookie 拼装,且单条 entry 的冗余度下降到 3 副本。这是典型的 RUM 权衡——用读放大换写吞吐和空间。
LAC:异步确认如何保证读一致性
Qa < Qw 意味着返回成功时还有副本没落盘。如果读者随便挑一个 Bookie 读,就可能读不到。BookKeeper 用 LAC(Last Add Confirmed) 解决这个问题:
- 写请求响应中携带该 Bookie 已连续确认的最大 entry ID。
- 客户端维护所有副本的 LAC,取"至少有 Qa 个副本都已确认"的那个值作为 ledger 级别的 LAC。
- 只有 LAC 以内的 entry 对外可见。
所以读路径是安全的:任意 entry ID ≤ LAC 的数据,必然已存在于至少 Qa 个副本上,读者可以从任意副本读到。LAC 随写请求捎带传播(piggyback),无需独立的确认往返。
// 客户端侧 LAC 计算:找出 quorum 内被覆盖的最大值
long computeLAC(long[] replicaLacs, int ackQuorum) {
long[] sorted = replicaLacs.clone();
Arrays.sort(sorted);
// 升序后,倒数第 ackQuorum 个即"至少 ackQuorum 个副本都达到"的水位
return sorted[sorted.length - ackQuorum];
}
Ledger 关闭时,最终 LAC 被写进元数据(ZooKeeper / etcd),此后该 Ledger 完全不可变。
四、游标:Pulsar 真正的护城河
Kafka 里消费者组的 offset 存在 __consumer_offsets;Pulsar 里每个订阅对应一个 ManagedCursor,它也是一个 BookKeeper Ledger(cursor ledger)。这意味着游标本身的推进是有持久化、可恢复的。
更重要的是:Broker 的 retention 策略可以基于"是否所有游标都已越过该位置"。当所有订阅都 ack 了某段消息后,对应的 Ledger 就可以被删除,而不需要像 Kafka 那样按时间一刀切。这对多消费者场景是质变——10 个消费者订阅同一 topic,Kafka 必须按最慢的那个设 retention,Pulsar 可以让每个订阅各自推进,落后的订阅只保留它需要的那部分数据。
四种订阅类型与 ack 语义
| 类型 | 并行度 | 顺序保证 | 典型场景 |
|---|---|---|---|
| Exclusive | 单消费者 | 全局有序 | 严格串行处理 |
| Failover | 单活跃 + 备 | 全局有序 | 主备高可用 |
| Shared | 多消费者 | 无顺序保证 | 高吞吐、可乱序 |
| Key_Shared | 多消费者 | 按 key 有序 | 需要并行又要 key 内有序 |
Key_Shared 是实践中价值最高也最容易踩坑的一个。它保证同一个 message key 始终路由到同一个消费者,从而在保持水平扩展的同时维护 per-key 顺序:
Consumer<OrderEvent> consumer = client.newConsumer(Schema.AVRO(OrderEvent.class))
.topic("persistent://prod/orders/events")
.subscriptionName("order-processor")
.subscriptionType(SubscriptionType.Key_Shared)
.keySharedPolicy(KeySharedPolicy.autoSplitHashRange()
.setAllowOutOfOrderDelivery(false)) // 关键:禁止乱序派发
.acknowledgmentGroupTime(100, TimeUnit.MILLISECONDS)
.negativeAckRedeliveryDelay(5, TimeUnit.SECONDS)
.deadLetterPolicy(DeadLetterPolicy.builder()
.maxRedeliverCount(3)
.deadLetterTopic("persistent://prod/orders/dlq")
.build())
.subscribe();
setAllowOutOfOrderDelivery(false) 是必须显式设置的。开启乱序派发后吞吐更高,但一旦某个 key 的消息处理失败触发重投递,该 key 后续消息可能已被处理——顺序语义就破了。
五、Bookie 的存储引擎:Journal + Entry Log + Index
单台 Bookie 内部是三层结构,理解它对性能调优至关重要:
写请求
├──► Journal(顺序 WAL,独立磁盘) ← 决定写延迟
└──► Memtable (write cache)
│ flush
▼
Entry Log(数据文件,多 ledger 混写)
│
Index(RocksDB / LedgerCache:ledgerId+entryId → 文件偏移)
Journal 是纯顺序追加的 WAL,只保证崩溃后可恢复。它必须放在独立的物理磁盘上——这是 BookKeeper 部署中最重要的一条铁律。journalSyncData=false(依赖 OS flush 而非每次 fsync)能大幅提升吞吐,但掉电会丢最后几条;对多数消息队列场景,配合副本冗余是可以接受的取舍。
Entry Log 是多个 ledger 的 entry 混写在一起的大文件。混写带来顺序 I/O,但也意味着单个 ledger 删除后空间不能立即回收,需要后台 Garbage Collection 线程做压实:把仍有存活 entry 的 entry log 重写到新文件,再删掉旧文件。这会产生明显的写放大,线上应配置 GC 限流与低峰窗口:
# bookkeeper.conf 关键项
journalDirectory=/mnt/nvme0/bk-journal # 独立 NVMe
ledgerDirectories=/mnt/ssd0,bk-data # 便宜的大容量盘
journalMaxSizeMB=2048
journalMaxBackups=2
journalSyncData=false
# GC 压实触发阈值与限流
minorCompactionThreshold=0.2
majorCompactionThreshold=0.5
minorCompactionInterval=3600
majorCompactionInterval=86400
compactionRateByEntries=10000 # 限制压实速度,避免打满 I/O
Index 默认用 RocksDB,dbStorage_rocksDB_writeBufferSizeMB 和 block cache 需要根据 ledger 数量调整。 ledger 数量多时,把 index 放到 SSD 甚至内存盘(ledgerStorageClass=Memory 用于小集群测试,生产不推荐)。
六、分层存储:让冷数据落到对象存储
Ledger 不可变这个设计在这里兑现了红利——它可以被整体搬到对象存储而无需任何合并或改写。
# broker.conf
managedLedgerOffloadDriver=S3
managedLedgerOffloadThresholdInBytes=1073741824 # 1 GB 后触发卸载
managedLedgerOffloadDeletionLagInMillis=3600000 # 卸载后本地保留 1h
managedLedgerOffloadMaxThreads=4
# S3 驱动配置(走 s3managed 或 jclouds)
s3ManagedLedgerOffloadRegion=ap-northeast-1
s3ManagedLedgerOffloadBucket=pulsar-tiered
触发条件是当前 topic 的累计数据超过阈值,卸载的是已关闭的历史 Ledger。卸载后读取冷数据时,Broker 从 S3 拉取 segment 到本地缓存再服务消费者——延迟从毫秒级升到百毫秒级,但成本降到 1/10 量级。
这里有个实践陷阱:managedLedgerOffloadThresholdInBytes 是命名空间级别累计的,不是"超过 1GB 就卸载最早的部分"。实际行为是当 managed ledger 的总字节数超阈值时才启动卸载流程。因此对于写入量大的 topic,卸载频率取决于 rollover 速度而非时间——需要配合 managedLedgerMaxEntriesPerLedger 一起调。
七、生产调优清单与踩坑
- Bookie 磁盘务必 journal / ledger 分离。混盘时 GC 压实和 journal 追加互相争抢磁头,p99 写延迟会飙到几百毫秒。
ensembleSize不要盲目调大。ensemble 越大,写入 fan-out 的连接数和元数据体积越大。中小集群 3-3-2 或 5-3-2 足够;跨 AZ 部署可考虑ensemble=6, Qw=4, Qa=2配合机架感知。- 未确认消息会撑爆内存。Shared / Key_Shared 订阅下,Broker 需要追踪每条已投递未 ack 的消息。设置
maxUnackedMessagesPerConsumer(默认 50000)并监控pulsar_subscription_unacked_messages,超阈值说明消费者处理跟不上或 ack 逻辑有 bug。 - Key_Shared 的 stuck consumer。某个消费者卡住时,默认会把它的 key range 标记 pending(
allowOutOfOrderDelivery=false时),导致那些 key 完全停摆。必须配死信队列 + 监控重投递次数。 - Bundle 分裂与负载不均。命名空间被切成若干 bundle(
defaultNumberOfNamespaceBundles,默认 4),bundle 是负载均衡的最小单位。热点 topic 集中在一个 bundle 时会造成单 broker 过载,需要主动调大 bundle 数或手动 unload。 - 延迟消息走内存索引。DelayedDeliveryTracker 默认在 Broker 内存中维护时间索引,Broker 重启后要从 cursor 重建。百万级延迟消息会显著拖慢 topic 加载,长延迟场景应改用外部调度。
结语
Pulsar 的设计哲学可以概括成一句话:把"顺序日志"和"消费进度"解耦,各自用最合适的结构实现。BookKeeper 的 ledger/fragment/quorum 模型提供了可横向扩展且写入后不可变的持久化层,ManagedCursor 则在之上叠加了多维、可独立推进的消费视图。这带来了 Kafka 难以做到的多订阅 retention、秒级分区扩容和冷热分层,代价是架构复杂度更高、故障域更多(多了 Bookie 和元数据存储两个组件),排障链路更长。
如果你的场景是"数千分区 + 多消费者组 + 需要长期保留历史数据",Pulsar 这套设计的收益非常实在。如果场景是"几百分区 + 极致吞吐 + 团队已熟悉 Kafka 运维",那么它的复杂度未必划算。技术选型从来没有银弹,只有对代价的清醒认知。

发表评论 取消回复