Apache Pulsar核心解析:存储计算分离与MessageId设计实践
2026/9/7 19:26:27 网站建设 项目流程

算算日子,COSCon‘25 同场活动 Pulsar Developer Day 就剩三天了。今年几个技术群里早就在传议程,消息中间件方向的专场活动能做到这种热度,确实不多见。借着这个由头,我打算把 Pulsar 这些年的一些核心设计、实践心得,以及大家在群里最常问的几个问题(比如“messageId 为什么长得那么奇怪”)一次性捋清楚。不管你是准备去现场,还是打算云围观,这篇内容应该都能让你对消息中间件这件事有个更立体的认识。

1. 先给还没入坑的朋友:消息中间件为什么是必修课

1.1 从单体应用到分布式,消息队列解决的三件事

如果你刚接触分布式系统,可能还不太理解,为什么大家都在聊消息中间件。我用最朴素的话来解释:当你的系统从一个进程变成几十个服务互相调用时,总会有一些请求是一瞬间涌进来的,也有一些操作根本不需要让用户一直等着。消息中间件做的事,本质上就是三件:削峰填谷、服务解耦、异步处理

拿电商下单来举例。用户点下“提交订单”那一刻,后端要处理库存扣减、优惠券核销、积分变动、发送短信通知、生成物流单……如果所有这些逻辑都同步跑完再返回“下单成功”,高峰期接口耗时可能直接飙到几秒甚至超时。实际生产环境里,大家会把短信通知、积分累计这类非核心链路扔进消息队列,订单接口只管落库和发消息,剩下的事情由下游服务异步消费处理。这既保证了用户体验,也避免了大促时数据库被瞬间打垮。

从单体到微服务的演进过程中,消息队列几乎成了标配。但“标配”不等于“随便选一个就行”,选型和架构设计直接决定了你未来三到五年的运维体感。

1.2 选型号之前,先搞清楚你的场景是哪种

目前主流消息中间件各有一批忠实用户:RabbitMQ 以灵活的路由和轻量部署著称,适合内部系统间的事件通知;Kafka 凭借超高吞吐和成熟的生态成为日志采集、数据管道的事实标准;RocketMQ 在电商和金融场景里口碑不错,事务消息做得很成熟;而 Apache Pulsar,这几年凭借“存储计算分离”和“多租户”两大王牌,在云原生环境里增长非常快。

选型时我一般会问自己四个问题:吞吐量要求是多少量级?是否要求消息可以回放、重跑?是否需要跨地域复制?扩缩容时能不能做到不停机?如果团队规模小、业务简单,RabbitMQ 或者直接用云厂商的托管队列就够了;如果每天要处理几十亿条数据,Kafka 或 Pulsar 才是真正的选项。Pulsar 的差异化优势在于:它的 Broker 不存储数据,只负责读写调度,扩容时不需要搬数据;而 Kafka 的 Broker 既管读写又管存储,扩容往往伴随着 rebalance 和数据迁移,痛点比较明显。

2. Pulsar 这几年凭什么站稳脚跟

2.1 存储计算分离:把“分层”发挥到极致

Pulsar 最核心的设计就是存储和计算分离。它底层依赖 Apache BookKeeper 作为分布式日志存储服务,Broker 层不落盘数据,所有的消息数据都交给 BookKeeper 集群管理。打个比方,Kafka 像一家“前店后厂”的餐馆,做菜和上菜是同一批人,翻台率受限于后厨;Pulsar 则像中央厨房模式,前厅只管接单上菜,菜品由统一的中央厨房配送,哪个分店忙了,多开几家前厅就行,不需要重新备菜。

存储计算分离带来的第一个红利是扩容优雅。Kafka 集群如果磁盘空间吃紧,通常需要新增节点,然后把 partition 副本迁移过去,期间可能影响线上流量。Pulsar 只需要给 BookKeeper 集群加节点,数据会自动做 rebalance,Broker 层毫不知情。

第二个红利是读写分离更彻底。生产者和消费者的流量可以单独调度,Broker 节点只维持连接和计算任务,状态全部放在 BookKeeper。某个 Broker 挂了,客户端会自动重连到其他 Broker,session 恢复成本极低。

第三个红利是存储成本下降。Pulsar 从早期版本就支持分层存储(Tiered Storage),老数据可以自动从 BookKeeper 卸载到 S3 或 HDFS 这类廉价对象存储。Kafka 在这块的能力也在补强,但 Pulsar 生来就是这么设计的,配合上更顺滑。

2.2 多租户和云原生基因不是噱头

很多中间件在传统 IDC 时代根本不需要考虑多租户,每个部门一套集群就完事了。但到了云原生阶段,资源要共享,成本要分摊,权限要隔离,多租户就成了刚需。Pulsar 的Tenant → Namespace → Topic三级模型,天生就是为了共享集群而设计的。不同的部门可以共享同一套 Pulsar 集群,但数据完全隔离,权限可以细粒度控制,配额可以分别管理。

我记得有次分享会上,一位维护过万级 Topic 集群的哥们说,他最喜欢 Pulsar 的一点是“Topic 是轻量级的”。Kafka 里 Topic 多了以后,分区数膨胀会拖垮 Broker 的整体性能;Pulsar 因为存储和计算分离,单集群承载百万级 Topic 都不是奇怪事,这在做多业务接入时特别香。你不需要每接入一个新业务就申请一个新集群,开个 Namespace 就够了。

2.3 和 Kafka 的对比:没有银弹,只有适不适合

我在这几年实际落地中,两种中间件都深度用过。Kafka 在日志管道、大数据生态集成方面依然无可替代,计算引擎和 Kafka 的集成深度远超 Pulsar。如果你整个技术栈都围绕 Flink、Spark、ClickHouse 转,Kafka 可能是心智负担最低的选择。

但如果你在乎的是云原生部署、跨地域容灾、多租户隔离、以及未来可能出现的突发扩容需求,Pulsar 的架构优势就会逐渐体现出来。尤其是跨地域复制,Pulsar 原生支持多集群的异步复制配置,而 Kafka 的 MirrorMaker 一直有种“外挂工具”的糙劲儿。正如没有银弹一样,选型时列个表格,把你的业务场景、团队熟悉度、未来规划填进去,该选谁答案自然浮现。

3. 热词解析:Pulsar 的 MessageId 为什么会是这个样子

3.1 先来破解那一串字符的秘密

兄弟群里有人甩了个问题:“为什么我拿到一个 messageId 长这样messageId|28077:20854:-1,这到底是个啥?” 我第一次看到这个格式的时候也愣了一下,因为习惯 Kafka 的 offset 是单调递增的整数,一眼能看懂。Pulsar 的 MessageId 却是一串“ledgerId:entryId:partitionIndex”的结构。

要理解这个格式,得从 Pulsar 的存储结构说起。我在上面提到,Pulsar 用 BookKeeper 来存数据,而 BookKeeper 最基础的存储单元是Ledger。Ledger 是一段追加写入的日志文件,里面包含一条条 Entry。对 Pulsar 来说,一个 Topic 的消息会顺序写入一系列 Ledger 中,每条消息在 Ledger 里对应一个 Entry。所以:

  • 28077Ledger ID,消息属于哪一个 Ledger;
  • 20854Entry ID,消息在 Ledger 里的物理偏移序号;
  • -10Partition Index,表示这条消息来自哪个分区,-1 通常用于非分区 Topic 或管理场景。

所以messageId|28077:20854:-1的意思是:这条消息位于第 28077 个 Ledger 的第 20854 条 Entry 上,这条消息所在的 Topic 是一个非分区主题。用文件系统来类比的话,Ledger 就像一本书,Entry 就是书里的页码,两者定位才能精确找到一条消息。

3.2 Ledger 机制:读懂了它就读懂了 Pulsar 一半

Ledger 是 Pulsar 存储的一个核心抽象。每个 Ledger 都有几个特性:追加写入不可变性按 Entry 序号随机读取。一个 Ledger 写满一定大小或者存活超过一定时间后,就会自动关闭,然后新建一个 Ledger 继续写。这个过程叫 Ledger Rollover。

为什么要设计成小段小段的 Ledger,而不是一个大文件写到底?这里面的门道很深。第一个好处是便于容量管理和恢复:如果某个 Bookie 节点挂了,只需要对它负责的那部分 Ledger 做数据恢复,而不是整库重建。第二个好处是便于实现高效的 TTL 删除:消息过期后,整段 Ledger 直接删掉就行,不需要像 Kafka 那样做日志紧凑和分段清理。第三个好处是并行度更好:多个 Ledger 可以分布在不同的 Bookie 上,写入流程天然可以并行扩展。

这里还要提一个隐藏机制——Ensemble。Pulsar 写入时不是把所有副本都写在同一批 Bookie 上,而是会为每个 Ledger 动态选一组 Bookie(Ensemble),数据以条带方式分布写入。这个设计保证了当某个 Bookie 出问题时,数据依然是可读的,而且不影响其他 Ledger 的写入。生产环境里最常见的配置是 E=3, W=3, A=2,意思是 3 个 Bookie 组成一个 Ensemble,需要写入全部 3 份才算成功,但允许 1 个节点故障不影响读取。这套逻辑源自 BookKeeper 的 Quorum 机制,理解了你就能解释为什么 Pulsar 能在高可用和数据一致性之间维持不错的平衡。

3.3 为什么不能直接用一个自增整数当 offset

很多人会问:Kafka 的 offset 是分区内单调递增的整数,直观又简单,Pulsar 为什么要搞这么复杂?这个问题的答案,还是要回到存储计算分离上。

Kafka 的 offset 是分段日志里的位置索引,它跟本地磁盘上的文件偏移强相关,所以天然只能由某个 Broker 自己管理,消费者恢复位点时必须找对 Broker。而 Pulsar 的消息位点由 LedgerId + EntryId 唯一定位,Broker 层不存储实际数据,消费者无论连接哪个 Broker,都能通过这个 ID 去 BookKeeper 里读取。这种设计让 Pulsar 的客户端连接可以无缝漂移,Broker 故障时消费者几乎无感知。

此外,Pulsar 的 MessageId 还包含 Batch 概念。生产端开启批量后,多条消息会打包进同一个 Entry,但对外暴露的 MessageId 依然能区分出 Batch 内部的单条消息。比如(ledgerId, entryId, partitionIndex, batchIndex)四个维度,比 Kafka 的“offset + batch 内相对位置”要更清晰。刚接触的人可能会觉得格式复杂,但当你开始做消息回溯、精确消费到某一条消息时,这种设计带来的准确性是简单整数 offset 做不到的。

3.4 从 MessageId 衍生出去:Cursur 与消息确认机制

理解了 MessageId 之后,“游标(Cursor)”的概念就比较好懂了。Pulsar 的每个订阅都有一个游标,游标里存储的是这个消费者组当前消费到的 MessageId。消费成功后,客户端会发送 ACK 给 Broker,Broker 把游标前移。这个游标信息默认保存在 BookKeeper 里,叫作Cursor Ledger,相当于把消费进度也做成了高可用存储。

实际排查问题的时候,游标和 MessageId 的配合非常有用。比如消费者组堆积了,你可以用pulsar-admin topics stats命令查看msgBacklog,那是当前游标到最大 MessageId 之间的消息条数。也可以用peek-messages命令指定 MessageId 来查看某条历史消息的内容,这在定位“某条消息到底有没有被消费到”的场景下特别好用。这些能力在 Kafka 里实现起来很别扭,因为消息位点放在消费者端,服务端对堆积状态的管理要弱不少。

4. 生产环境绕不开的:订阅模式、分层存储与运维要点

4.1 三种订阅模式怎么选

Pulsar 的订阅模型是它的又一大卖点。它原生支持三种消费模式,很多人刚接触时容易混淆,我在这里用一个生活场景来拆解:

  • 独占订阅(Exclusive):一个 Topic 同一时刻只能有一个消费者,适合强顺序场景,比如把订单状态流转消息按顺序处理,不能并发。这种模式最简单,但扩展性最差。
  • 共享订阅(Shared):多个消费者共同消费一个 Topic 的消息,消息按 round-robin 或 pending-ack 状态分配给不同消费者,吞吐量最高,但消息顺序性无法保证。适合大多数数据处理场景,比如消息量很大、处理逻辑之间没有严格依赖关系。
  • 灾备订阅(Failover):多个消费者中只有一个活跃消费者接收消息,其他消费者作为备用,活跃消费者挂掉后自动切换。相当于带高可用的独占模式。

还有一种是 Key_Shared 订阅,介于 Shared 和 Exclusive 之间:相同 key 的消息只发给同一个消费者,比如把同一个用户 ID 的所有订单事件固定发给一个消费者处理,这样既能并行消费,又能保证单用户的局部有序。生产环境里我用 Key_Shared 解决过一个典型问题:用户积分变动必须按时间顺序处理,但不同用户之间可以并行,用 Exclusive 太浪费,用 Shared 会乱序,Key_Shared 刚好完美解决。

4.2 分层存储:省钱大法

Pulsar 的分层存储(Tiered Storage)是个非常实用的功能,但宣传得还不够。默认情况下,数据只存在 BookKeeper 里,而 BookKeeper 一般用的是 SAS 盘或者 SSD。数据量大了以后,存储成本非常扎眼。Pulsar 可以配置自动卸载策略,把超过一定时间或者一定容量的旧数据搬到 S3、阿里云 OSS、腾讯云 COS 这类对象存储里,读的时候如果 BookKeeper 里没有,会自动从对象存储拉回来。

这个机制在“消息回溯”场景里特别有用。很多团队要求消息至少保留 7 天甚至 30 天,以便排查问题或者重新跑数。如果不做分层存储,你就得为这 30 天的数据预留大量高性能磁盘。开了分层存储之后,热数据在 BookKeeper,冷数据在对象存储,成本能降一个数量级。我见过有团队把 Pulsar 保留了 90 天的数据,存储成本跟 Kafka 保留 3 天差不多,这就是分层存储的魔力。

4.3 运维踩坑:几个值得记录的真实案例

Pulsar 的整体运维体验比传统消息中间件要好,但也不是没有坑。我把自己落地过程中遇到比较多的问题整理成一张速查表:

现象可能原因排查方法
生产端发送延迟突增BookKeeper 写入延迟变高,磁盘 I/O 忙bookkeeper shell ledgercheck检查磁盘状态;观察 bookie 的 Journal 和 EntryLog 写入延迟
消费端 backlog 持续增加消费者处理慢;某个消费者挂掉没有重连pulsar-admin topics stats看 backlog;检查客户端日志连接状态
Topic 写入失败 “No such ledger”Ledger 被自动删除,但客户端还在写(时间窗口问题)调整保留策略,确保写入过程中 Ledger 不会被清理
客户端连接不断重连Broker 和 Bookie 之间的网络抖动;认证过期检查认证 token 有效期,查看 broker 日志里的连接错误
分区数量无法减少设计阶段分区数不合理先评估消息量和消费并发度,一般先小后大,确实不够再扩容分区

大家最容易踩的坑其实是“分区数设太多”。Pulsar 对 topic 数量的容忍度极高,但分区太多会导致 BookKeeper 的 Ledger 数量爆炸,增加 ZooKeeper 和 Bookie 的元数据压力。我一般建议:初期按消费者并行度来定分区数,后期根据实际流量再逐个加,不要未雨绸缪地建一堆分区。

4.4 客户端参数调优:三个我建议必调的配置

Pulsar 客户端很多参数都有默认值,但生产环境里默认值往往不是最优解。

第一个是生产端的 Batching。默认 batching 是开启的,但如果你追求极低延迟,比如控制在 5ms 以内,需要把batchingMaxPublishDelay调小或直接关闭 Batching。反过来,如果要追求高吞吐,可以把batchingMaxMessages调到 1000 以上,batchingMaxBytes调到 128KB 以上。这个取舍跟 TCP 的 Nagle 算法有点像,具体按你的业务对延迟的敏感度来。

第二个是消费端的 Ack 超时。很多人设了ackTimeout,但设得太短(比如 1 秒),消费逻辑稍微慢一点就会触发重投,导致消息重复处理的概率上升。如果你下游不要求精确一次,建议把 ackTimeout 调大甚至设为 0(不超时),配合消息重试队列去处理失败场景。注意:ackTimeout 和negativeAckRedeliveryDelay是两套独立的重试机制,别混淆。

第三个是接收队列大小receiverQueueSize默认 1000,对处理很快的逻辑来说,这个值偏小,容易让消费者频繁陷入等待。我通常在 CPU 密集型处理时设为 100~200,IO 密集型或调用下游接口场景设到 1000 或更高,避免消费者线程空转。

5. 活动背面的技术信号:为什么开发者日值得你专门跑一趟

5.1 从议程看行业趋势

说实话,国内专门围绕 Pulsar 做线下开发者日的机会并不多,一年到头可能也就几场。去年我在一次类似的活动中听到了好几个非常有启发的议题:有团队分享了怎么把 Pulsar 用在金融级交易系统里,保证消息不丢不重;也有人讲了怎么用 Pulsar Functions 做轻量级流处理,替代一部分 Flink 任务;还有人在分享 Pulsar + Lakehouse 的实践,把消息中间件和数据湖打通的路径已经非常清晰。

从这些议题能明显感觉到,Pulsar 已经不是早期那个“技术超前但生态不足”的项目了。它开始向金融、运营商、车联网这些传统行业渗透,在云原生、数据集成方面也出现了一大批成熟的同场产品。今年 COSCon’25 同场活动定名为 Pulsar Developer Day,来的 speaker 和 sponsor 大多是实际落地的一线工程师,不是空谈架构的纸上谈兵,这种内容密度比主流技术大会的高层 keynote 要实在得多。

5.2 带着问题去,收获才更大

参加这类活动的正确姿势,不是坐在台下听 PPT,而是带着自己的实际问题去。我每次去都会准备一张问题清单,比如:

  • 我们在生产环境遇到 Bookie 扩容后写入抖动,官方有没有推荐的 rebalance 策略?
  • 在 Pulsar 里实现“延迟消息”和“定时消息”,最佳实践是什么?
  • Pulsar 和 Flink 集成的 checkpoint 一致性如何处理?
  • 在跨地域复制场景下,消息延迟的监控指标怎么设阈值?

这些问题在现场很容易通过 speaker 的分享和 Q&A 环节找到答案,甚至可以当面加微信交流后续细节。线上文档写得再多,也不如和核心维护者、一线实践者面对面聊十分钟收获大。今年活动倒计时只剩三天,票务信息在 COSCon 官网可以看到,对消息中间件生态感兴趣的朋友,我觉得完全可以抽出半天时间过去转转。

5.3 如果你没法到场,可以这样保持同步

对于没法到场的朋友,我的建议是:重点留意活动后的 PPT 和视频回放,一般 COSCon 系列的产出质量高,官方渠道都会公开。同时把活动中提到的实践项目、仓库链接都收藏起来,动手跑一遍,比收藏一堆“技术文章”有用得多。另一个做法是加入 Pulsar 中文社区或邮件列表,很多讨论从活动当天会一直延续到线上,你在群里提问经常能得到活跃贡献者的直接回复。

6. 写在最后:我的个人体会

我自己从 Kafka 转向深入使用 Pulsar,大概经历了一个从“这玩意怎么这么复杂”到“原来这样设计是合理的”再到“回不去了”的过程。MessageId 的复杂结构、Ledger 的切分机制、游标的持久化,初看都是增加理解成本的东西,但用久了会发现,这些设计全都是为了支撑同一个目标:让存储和计算解耦,让扩缩容不再痛苦,让大数据量下的消息系统保持稳定

马上就是 Pulsar Developer Day 的日子了,这种主题的交流机会且行且珍惜。不管你是消息中间件的老兵,还是刚刚入行的新人,去听听一线踩坑的人怎么说,远比自己在电脑前瞎琢磨效率高。准备去现场的朋友,我建议提前列出你当前系统里三个最痛的问题,直奔对应议题的讲师,别不好意思问。三天后见。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询