Pulsar 架构深入
提出问题
Apache Pulsar 是近年来消息队列领域最受关注的新秀之一,但很多开发者只是听过它的名字,并不清楚它和 Kafka 到底有什么本质区别。面试中常被问到:Pulsar 的"计算与存储分离"是什么意思?为什么它能做到无 Rebalance、多租户原生支持?Pulsar 用 BookKeeper 做存储层,和 Kafka 的日志存储相比有什么优势?如果你的系统需要更高的读写隔离性、更低的消息延迟,或者希望冷数据自动下沉到 S3 降成本,Pulsar 的架构设计值得深入理解。
分析问题
计算与存储分离:Broker 无状态 + BookKeeper 持久化
Pulsar 最核心的架构决策是将 Broker(消息代理)与存储层彻底分离。Broker 不持有任何持久化数据,它只负责:接收客户端请求、路由到正确的 BookKeeper 节点、缓存热数据。这意味着 Broker 可以随时加入或退出而无需数据迁移——不像 Kafka 那样 Broker 挂了要重新对分区做 Leader 重新选举和副本同步。
存储层使用 Apache BookKeeper,一个专为 append-only 日志设计的分布式存储系统。BookKeeper 的核心概念是 Ledger(账本) 和 Entry(条目):一个 Topic 被切分成多个 Segment(Ledger),每个 Segment 包含若干 Entry。Entry 是 Pulsar 消息的最小存储单元。
// Pulsar 消息写入 BookKeeper 的简化示意
// 每个 Segment 是一个 Ledger,写入后不可变
LedgerHandle ledger = bkClient.createLedger(
LedgerHandle.ENSEMBLE_SIZE, // 3 副本
LedgerHandle.WRITE_QUORUM_SIZE, // 2 个确认
LedgerHandle.ACK_QUORUM_SIZE, // 2 个 ack 即可返回
BookKeeper.DigestType.CRC32,
"password".getBytes()
);
for (Message msg : messages) {
// Entry 写入后自动复制到 ensemble 中的副本
ledger.addEntry(msg.serialize());
}这段代码展示了 Pulsar 消息写入的配置模型:每个 Ledger 可以独立设置副本数(Ensemble)、写确认数(Write Quorum)和确认数(Ack Quorum),这是 Kafka 集群级别配置无法做到的细粒度控制。
分层存储:冷数据自动下沉降本
Pulsar 的分层存储(Tiered Storage)是另一个亮点。当 Segment 数据在 BookKeeper 中达到一定时间阈值或大小后,Pulsar 自动将其卸载到更廉价的存储后端,如 S3、GCS 或 HDFS。消费者读取历史消息时,Broker 从冷存储拉取数据,对客户端完全透明。
# pulsar-broker.conf 中分层存储配置示例
managedLedgerOffloadDriver=aws-s3
s3ManagedLedgerOffloadBucket=pulsar-tiered-storage
s3ManagedLedgerOffloadRegion=us-east-1
managedLedgerOffloadThresholdInBytes=1073741824 # 1GB 自动卸载
managedLedgerOffloadDeletionLagInMillis=604800000 # 卸载后保留7天才删除本地这个配置意味着:Segment 一旦超过 1GB 且已写入超过 7 天,Pulsar 会将其复制到 S3 并标记为可删除。相比 Kafka 需要手动管理磁盘空间或依赖 Tiered Storage 插件(KIP-405),Pulsar 的 Tiered Storage 是架构原生能力,不需要额外组件。
与 Kafka 的对比:关键差异
| 维度 | Pulsar | Kafka |
|---|---|---|
| 存储架构 | 计算与存储分离(Broker + BookKeeper) | 存储与计算耦合(Broker 即存储) |
| 扩容/缩容 | 无 Rebalance,Broker 无状态,即加即用 | 分区迁移,Rebalance 期间不可用 |
| 多租户 | 原生支持(tenant/namespace 层级隔离) | 靠 ACL + Quota 模拟 |
| Geo 复制 | 内置异步复制,支持跨集群 | MirrorMaker 或 Confluent Replicator |
| 订阅模式 | 独占/共享/灾备/Key_Shared 四种 | Consumer Group 一种(类似共享) |
| 消息确认 | 单条 ACK(BookKeeper 级别) | 基于 Offset 批量提交 |
| 冷存储 | 原生分层存储,自动卸载到 S3 | 需插件或手动迁移 |
四种订阅模式
Pulsar 的 Topic 支持四种订阅模式,这是它比 Kafka 灵活的重要体现:
- Exclusive(独占):一个 Topic 只有一个消费者,等价于 Kafka 的单个 Partition 消费。
- Shared(共享):多个消费者轮询消费消息,Kafka 的一个 Consumer Group 即此模式。
- Failover(灾备):一个主消费者,其余作为备份,主消费者挂掉后自动切换。
- Key_Shared(键共享):同一 key 的消息路由到固定消费者,保证有序性,解决了 Shared 模式下无序的问题。
// Key_Shared 订阅示意:按消息 key 路由到固定消费者
Consumer<byte[]> consumer = client.newConsumer()
.topic("persistent://public/default/orders")
.subscriptionType(SubscriptionType.Key_Shared)
.subscriptionName("order-processor")
.subscribe();总结
Pulsar 的计算与存储分离架构从根本上解决了 Kafka 在弹性伸缩、多租户隔离和存储成本方面的痛点。BookKeeper 作为存储层提供了细粒度的副本控制和分层存储能力,Broker 无状态化让扩容不再需要 Rebalance。但 Pulsar 并非没有代价:架构复杂度更高,部署运维需要同时管理 Broker 和 BookKeeper 两套集群,学习曲线也更陡。选择建议:如果需要多租户、高读写隔离、低成本冷存储,优先考虑 Pulsar;如果只做简单的高吞吐日志收集,Kafka 的成熟生态和低运维成本仍是稳妥选择。
参考
参考:Apache Pulsar 官方文档 - Architecture;BookKeeper 设计文档;《Streaming Systems》Pulsar 章节;Pulsar 与 Kafka 对比白皮书