Skip to content

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 消息的最小存储单元。

java
// 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 从冷存储拉取数据,对客户端完全透明。

yaml
# 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 的对比:关键差异

维度PulsarKafka
存储架构计算与存储分离(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 模式下无序的问题。
java
// 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 官方文档 - ArchitectureBookKeeper 设计文档;《Streaming Systems》Pulsar 章节;Pulsar 与 Kafka 对比白皮书

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。