一篇关于消息中间件选型的文章,核心重点应该落在“为什么选”和“怎么落地”上,而不是堆砌官方文档式的功能清单。Pulsar这几年在技术圈讨论度很高,但真正把它用在生产环境并踩过坑的人,其实比用Kafka的人少一个数量级。这篇文章打算把Pulsar从架构原理到选型决策、再到实际部署使用中的关键细节串一遍,重点说清楚它和Kafka那套经典架构的本质差异,以及在什么场景下该选它、什么场景下别硬上。如果你是正在做技术选型或者准备把Pulsar引入团队的人,这篇应该能帮你省掉不少调研时间。
1. 为什么聊到消息中间件时,Pulsar值得单独写一篇
市面上消息中间件不少,Kafka、RabbitMQ、RocketMQ各有各的地盘。过去几年大部分团队的技术选型基本绕不开Kafka,毕竟生态成熟、资料多、踩坑的先行者也多。但Kafka在处理一些特定场景时,总是让人有种“能用但不太痛快”的感觉。比如存储和计算耦合在一起导致扩容要搬数据、分区数上去之后Broker的负担明显加重、多租户隔离做起来费劲。这些痛点其实一直存在,只是很多团队用了一堆workaround硬扛过去,习惯了也就不觉得是问题。
Pulsar之所以值得单独写一篇来说,是因为它从架构层面换了条路。它把存储和计算彻底拆开了,Broker只负责消息的路由、调度和缓存,真正的数据持久化交给底层的Apache BookKeeper去管。这一拆带来的连锁反应是:扩容不再需要搬数据,Broker可以做到无状态化,分区数对Broker的压力大幅下降,基于BookKeeper的Segment存储还能支持比Kafka大得多的消息积压能力。这一整套设计和经典Kafka是完全不同的思路。
标题里写了“系列六”,说明前面应该已经聊过Kafka、RabbitMQ、RocketMQ这些选型了。如果把消息中间件比作一个工具箱,Kafka像一把大号扳手,结实耐用但调起来费劲;RabbitMQ像一套精密螺丝刀组,细分场景好用但大吞吐场景吃力;那Pulsar就更像一套模块化的电钻,能换各种头去应对不同活路,前提是你愿意先花时间搞明白它的构造。
这篇文章会从架构原理讲到部署实操,再到生产环境常见问题,最后回到选型决策本身。适合刚接触Pulsar的技术负责人、后端开发,也适合正在为某个具体业务场景挑MQ的架构师。哪怕你对Pulsar完全没概念,只要知道消息队列是用来解耦、削峰、异步的,就能跟着这篇文章把它搞明白。
2. 架构决策的底层逻辑:Pulsar和Kafka的分水岭在哪
2.1 从Kafka的痛点反推Pulsar的设计目标
先说Kafka让我最头疼的几个问题。第一是扩容,Kafka的分区是有状态的,数据落在特定Broker上,想加节点就得做分区重平衡,这个过程既要搬家又要限速,分区多了之后重平衡还容易把集群搞出问题。第二是分区数没法无限扩展,分区越多,每个Broker要维护的元数据、文件句柄、内存开销就越大,到一定程度整个集群的性能会明显下降。第三是积压能力,如果消费者挂了或者下游处理变慢,Kafka的消息会一直攒在分区日志里,日志一长,清理和读取都会变慢。
Pulsar从根上换了思路。它把“消息来到了哪台机器”和“消息存到了哪里”拆成了两回事。Broker只管接收请求、维护订阅状态、把消息发给消费者;消息真正的副本和持久化由BookKeeper这个专门的存储层负责。Broker本身不存用户数据,所以它可以是无状态的——想加Broker就加,不需要大规模迁移存量数据。
这两套架构的差别用一句话概括就是:Kafka把所有事绑在一起,Pulsar把各层拆开各管各的。听起来好像是Kafka设计得不好,其实不是,Kafka是很多年前的架构,当时分布式存储的成熟度不高,把存储写在Broker本地反而是最优解。Pulsar敢于拆分,是因为BookKeeper这套专门的分布式日志存储系统足够成熟,才让这个思路变成了现实。
2.2 存储与计算分离带来哪些连锁收益
存储与计算分离不是一个空概念,它带来的是实打实的好处。第一,Broker扩容彻底简化了,不用搬数据,把新节点加进来让它分担流量就行。第二,存算分离之后,Broker之间不需要做数据副本同步,同一个Topic的分区副本由BookKeeper负责,Broker省掉了大量网络和磁盘IO开销。第三,在BookKeeper的Segment存储模型下,一个消息积压很久也不会拖垮性能,系统可以一直往新的Segment里追加写入,老Segment归档处理。
多租户能力也是这套架构的额外红利。Pulsar的Topic名自带层级结构,比如persistent://finance-prod/orders/payment-events,天然分成了tenant(租户)、namespace(命名空间)、topic(主题)三级。在做隔离和权限控制时,这种结构比Kafka基于ACL的实现直观得多。团队A和团队B哪怕共用同一套集群,资源配额能分别限制,互不干扰。
当然,这套架构也有需要接受的代价。最明显的就是链路变长了,一条消息要经过客户端->Broker->BookKeeper多跳,写延迟理论上比Kafka直接写本地磁盘要高一点点。好消息是实际部署中Pulsar在大多数场景下延迟也能稳定落在个位数毫秒到十几毫秒,对绝大多数业务足够了。对我来说,这些代价换回的是运维上的省心和扩展上的自由度,值得。
3. Pulsar核心概念与消息流转全过程
3.1 Broker、BookKeeper、元数据服务三段架构盘点
Pulsar集群从逻辑上由三部分组成。第一部分是Broker层,它负责处理客户端的生产消费请求、管理订阅游标(Cursor)、处理消息的TTL和积压策略,是无状态的服务节点。第二部分是BookKeeper存储层,由一组Bookie节点构成,负责消息数据的分布式持久化,每个消息会被写入多个Bookie形成多副本。第三部分是元数据服务,通常用ZooKeeper或Etcd来存Topic、Broker、Bookie这些元信息和全局配置。
打个比方,Broker是餐厅的前厅,负责接单上菜;Bookie是后厨,负责把菜做好保存住;元数据服务是餐厅的收银排班系统,记录每张桌子对应哪个服务员。客人(客户端)只跟前厅打交道,不知道也不用关心后厨是怎么运作的。前厅不够了,多开几个门面就行,后厨是独立的,不会因为前台扩容导致厨具不够用。
消息数据进到BookKeeper之后,是以Segment为基本单位存储的。一条消息进来,会被追加到当前Ledger的Segment中,写满了一个Segment就换下一个。一个Topic的数据由一串连续的Segment组成,这些Segment分布在不同Bookie上,哪台Bookie挂了,它的Segment会被其他Bookie上的副本顶上。
3.2 一条消息从生产到消费的完整链路
一条Pulsar消息从产生到被业务方消费,中间经过的路径是:生产者客户端把消息发给某个Topic对应的Broker;Broker接到消息后,把消息写入BookKeeper的当前Ledger;BookKeeper完成多副本写入后返回确认;Broker再把确认回给生产者。消费者这边,消费者客户端向Broker发起订阅请求并获取消息;Broker从缓存或存储层把消息读出来返回给消费者;消费者处理完后发ack确认,Broker更新订阅游标。
这个流程里有几个容易被忽略但很重要的细节。第一个是BookKeeper写入确认,它需要等待至少quorum数量的副本写入成功才会返回成功,这个quorum默认是多数派,比如3副本就至少2个成功,和多数派共识是一个道理。第二个是消费游标的推进,Pulsar的游标信息不会跟着用户数据一起存,而是作为BookKeeper里的特殊Ledger来保存,这样即使消费者全部下线,游标也不会丢,重连之后能从最后确认的位置继续读。
第三点是Pulsar支持三种订阅类型。独占订阅(Exclusive)是一条消息同一时间只能被一个消费者消费;共享订阅(Shared)让消息在多个消费者之间按轮询方式分发,谁有空谁处理;故障转移订阅(Failover)则是主消费者优先,主消费者挂了自动切换。选哪种订阅类型取决于业务场景,不能用错了,比如你需要的明明是共享订阅,结果建Topic时用了默认的独占订阅,消费者一多就会出现消息没人处理的现象。
3.3 消息积压为什么在Pulsar里没那么可怕
“积压”这个词在Kafka里是有点敏感的。Kafka消费者的消费速度如果跟不上生产速度,消息会一直堆在分区日志里,日志文件越长,磁盘和内存的压力越大,对后续读写性能的影响越明显。运维一般不希望大家把Kafka当积压缓冲池来用,而是希望消费者的速度尽量跟上生产速度。
Pulsar因为存储和计算分离,积压场景的处理轻松得多。消息写入的都是Bookie上的Segment,Segment会按大小滚动,系统可以方便地管理这些数据块。即使消息在Topic里堆了几亿条,对Broker来说,它的内存里只有一部分热数据缓存,剩下都在存储层放着,不会因为积压量大就拖垮Broker的性能。实际上Pulsar的宣传口径里明确提到过,一个Topic可以存储TB甚至PB级数据,同时还能维持正常的读写性能。
当然这不代表你可以无限放纵积压。积压太大会让游标位置离消息尾部越来越远,消费时需要在Bookie上读很老的Segment,磁盘顺序读的性能也会下降。但从架构角度看,Pulsar确实是目前开源MQ里积压能力最强的一个,这也是它适合做统一消息平台的原因之一。
4. Pulsar的关键机制与生产配置实操经验
4.1 消息确认与重试、死信设计的坑
消费确认机制是MQ使用中第一道关卡。Pulsar的ack分单条确认和累积确认——在独占和故障转移订阅模式下,客户端可以批量确认,ack一个消息就代表它前面的消息也都确认了;共享订阅模式则必须逐条确认,因为消息是分散给不同消费者处理的,没法用游标一次性推进。这个差异如果没弄清楚,写共享订阅的消费者时会发现明明没处理完的消息被标记成已确认了。
除了正常ack,还有负向确认nack。消费者拿到消息后如果发现处理失败,可以nack这条消息,让它重新进入待投递队列。但要小心,nack控制不好会导致消息无限循环重投。我自己在生产环境的做法是:正常业务异常直接nack,并且配一个最大的nack次数;超过次数就投递到死信主题(DLQ),由专门的任务去分析和人工处理。Pulsar支持在Topic上配置死信策略,比如maxDeliverCount设为3,超过3次后自动转入死信Topic,这个配置在命名空间级别可以统一设置,非常方便。
还有一个容易被忽视的点是消息的消费超时。消费者拿走了消息但长时间不确认,Pulsar会有个ackTimeout的概念,超时后会把消息重新投递给其他消费者。这个值设置得如果太小,比如设成5秒,而你的业务处理一条消息正常情况下就要8秒,那消息会被反复投递,产生大量重复消费。设置太大又会让故障恢复变慢,消息要等很久才会被重新投递。经验值是先压测出正常处理时长的P99,然后在这个基础上乘1.5到2倍作为ackTimeout。
4.2 消息幂等:前端点两次真的会变成两条消息吗
“前端点两次算是发两条消息吗”这个热词背后,其实就是消息中间件里老生常谈的重复消息问题。答案是:如果你不做任何防护,确实可能变成两条消息被发到MQ里,然后再被消费两次。这里的重复可能来自两个层面:生产者重复发送和消费者重复消费。
先说生产者侧。MQ的at-least-once语义决定了客户端网络超时重试时,服务端可能已经写入了消息但返回响应丢了,客户端重新发送就会产生重复。Pulsar生产者的send超时重试机制就可能导致这种问题。业界通用的解法是做生产者幂等,给每条业务消息生成一个全局唯一的消息ID(比如UUID或者基于业务唯一键生成),服务端通过去重来保证同一ID只落库一次。Pulsar的Batch消息里本身就能携带消息Key,你可以用消息Key做去重依据。
再说消费者侧。即使生产者只发了一条,消费者也可能因为ack超时、网络抖动、重平衡等原因收到同一条消息两次。所以消费端的幂等是必做的,不管用哪个MQ都一样。常见方案有:依赖数据库唯一索引做插入约束;利用Redis的setnx做处理状态标记;或者把消息里的业务唯一键作为主键做幂等写。不要把幂等的希望寄托在MQ配置上,所有MQ都做不到精确一次(exactly-once)的端到端保证,精确一次是分布式系统里最难的课题之一。在生产上做到“at-least-once + 消费端幂等”就是最可靠的组合。
4.3 顺序消息怎么在Pulsar里实现
顺序性也是一个经常被问到的问题。Kafka的顺序消息是靠分区内有序实现的,同一个Key的消息进同一个分区,消费者按序消费就能保证顺序。Pulsar里同样用的是这个套路,但需要注意订阅类型对顺序的影响。
Pulsar里如果你想让消息严格有序,必须保证两点。第一,生产端用MessageKey指定分片路由规则,Pulsar支持按Key哈希或者按Key取模分发到不同分区,同一个Key的消息会进同一个分区。第二,消费端只能用独占订阅或故障转移订阅,不能用共享订阅。因为共享订阅模式下,消息会被多个消费者并行处理,就算生产端顺序正确,消费端执行顺序也会乱。如果你的业务要求的“顺序”是指同一个订单的支付、发货、完成通知必须按顺序处理,那用订单ID作为MessageKey,并用独占或故障转移订阅,就能得到和Kafka分区内有序一样的效果。
还要补充一个细节,Pulsar的消息是支持延迟投递的。有些场景需要在某个时间点之后消费者才能看到消息,比如订单超时未支付就关单,生产端可以给消息设置一个延迟时间,消息到点之后才允许被消费。这个功能在Kafka里原生的支持很弱,通常要自己造轮子,而Pulsar在API层面直接支持,对这个场景可以说是开箱即用。
4.4 部署形态与关键参数清单
Pulsar的部署有两种主流方式。一种是裸机或虚拟机部署,安装Broker和Bookie服务;另一种是用Kubernetes部署,借助Helm Chart快速拉起。对大多数团队来说,走Kubernetes部署是更省心的路线,因为Pulsar的组件管理、扩缩容、监控都能用K8s原生能力来管。
部署时几个关键参数我列一份清单。Broker端:managedLedgerDefaultEnsembleSize控制副本数,一般设3;managedLedgerDefaultWriteQuorum和managedLedgerDefaultAckQuorum分别控制写入和确认需要的副本数,一般也设2以上;allowAutoTopicCreation建议默认开,但要让运维知道开了之后有Topic爆炸的风险,可以在命名空间级别限制Topic数量上限。Bookie端:journalDirectory和ledgerDirectories一定要分盘放,journal(日志)用小容量高性能盘,ledger(数据)用大容量普通盘,混在一起会在高写入时互相干扰;dbStorage_writeCacheMaxSizeMb和dbStorage_readCacheMaxSizeMb需要根据内存规划,给JVM堆留足空间的同时要给RocksDB缓存留够量。
消费者端还有一个吞吐调优的大杀器,就是批量接收。Pulsar的Consumer支持批量拉取,一次拿到多条消息在客户端本地做缓冲处理。把receiverQueueSize从默认的1000调大,比如5000甚至10000,在高吞吐场景下性能会有立竿见影的提升。但这个值也不能调得过猛,因为消费端进程如果崩溃,缓冲区里未确认的消息可能全部需要重新投递,恢复时间会变长。
5. 生产环境实测与常见问题排查实录
5.1 消费者数量上来了消息却没被消费是怎么回事
我遇到过比较多的情况是,用了共享订阅模式,起了好几个消费者实例,以为消息会被平分到各个实例去处理,结果所有消息都打到了一个实例上,其他实例在空转。排查看下来,发现问题是Topic创建时指定的订阅模式不对,创建的是独占订阅。后来在代码里统一用consumerBuilder.subscriptionType(SubscriptionType.Shared)去指定订阅类型,然后对已有Topic的订阅做迁移才解决。
还有一个类似的问题是,共享订阅模式下消息分配不均衡。Pulsar的共享订阅默认是按消息维度分发,每条消息发送给一个消费者,但如果是批量生产且batch里消息数量大,分发粒度会变成“完整batch发给一个消费者”,从而出现明显的倾斜。解决方案是调整生产端的batchingMaxMessages参数,或者把消费者数量控制在一个合理的范围,不要图省事开上千个消费者去期望消息被均匀摊开。
5.2 BookKeeper磁盘占用异常增长的处理经验
书归正传,BookKeeper磁盘增长快是Pulsar运维最常遇到的问题之一。一个比较容易踩的坑是:消费者下线之后,它的订阅游标停住了,消息就一直在Bookie上保留着,对应Ledger永远不会被删除。尤其是测试环境,频繁建Topic、开消费者,消费者用完就删,但订阅还挂在命名空间里,积压的消息没人消费,磁盘就慢慢被撑满了。
排查思路不复杂。用pulsar-admin topics stats命令看每个Topic的backlog数量,找到backlog大量增长但消费者不在线的Topic。然后用pulsar-admin topics delete把废弃Topic删掉,或者在命名空间层面设置retention策略,控制消息保留时间,超时的自动清理。这里建议从一开始就给测试命名空间设一个比较短的保留策略,比如1小时,防止这类问题造成线上事故。
5.3 性能压测发现的线程模型问题
还有一次压测时发现Pulsar的生产吞吐一直上不去,加机器也没太大改善。排查之后发现瓶颈在Broker的IO线程和Bookie的Journal写入上。Pulsar的Broker默认配置适合常规负载,但在高吞吐场景需要在Broker的conf/broker.conf里调大numIOThreads和numExecutorThreads,这两个参数管的是Broker处理网络写入和消息分发用的线程池大小。机器核心数多的时候,默认值太低就成了瓶颈。
Bookie那边也有类似的并发参数。journalSyncData如果为true,表示每次写入都需要刷盘后才返回,数据安全但性能有损失,在要求低延迟的环境里可以设为false,数据会先落操作系统页缓存,异步刷盘。两个参数怎么取舍,得看你的业务对丢数据的容忍度。金融、订单类业务不敢乱调,日志、通知类业务可以适当放宽。
5.4 消息重复的终极排查思路
最后聊聊排查消息重复这个经典话题。如果线上消息被重复处理了,先别急着怪MQ。按下面这个顺序排查基本能定位问题。
第一步,看生产端有没有重试机制,是不是send超时后客户端自动重发而服务端已写入,产生了重复。第二步,看消费端ackTimeout设置是否合理,一条消息要8秒处理完,ackTimeout设5秒,中间超时触发重投,必然重复。第三步,看消费逻辑里有没有做幂等,如果对数据库只有insert操作,那建一个唯一索引就能挡住;如果做了多次更新,那需要设计好乐观锁或版本号机制。第四步,看订阅模式是不是故障转移,主消费者切换瞬间游标位置可能回退几条消息,短暂重复是正常现象。
这四个步骤完走一遍,90%以上的重复消费问题都能找到根因。剩下的10%属于极端边界情况,比如Broker和Bookie之间的数据复制出现了脑裂场景,这种就只能靠可用区设计和数据校验去兜底了。总之记住一句话:用MQ之前,先确保你的系统能把“至少一次”变成“业务上只处理一次”。
6. 架构决策最终建议:Pulsar适合谁,不适合谁
6.1 适合选Pulsar的场景类型
如果你们团队面临下面任意几种情况,认真考虑Pulsar是值得的。
第一,你打算建设统一的公司级消息平台,业务线多、Topic多、隔离需求强,多租户能力会帮你省很多事。第二,你们有海量消息积压的硬需求,比如每天几十亿条事件数据,消费速度波动大,积压是常态,Pulsar的存储模型显然更匹配。第三,你对延迟没有极端要求,但希望系统能灵活扩缩容,Pulsar的存算分离架构让扩缩容变得很轻松。第四,你们有跨地域复制需求,比如业务本身是多机房的,Pulsar原生支持跨地域复制,可以实现容灾和多活。
6.2 哪些场景我建议暂时别上Pulsar
如果你的业务属于下面的场景,用Pulsar不一定划算。
一是极轻量场景,整个系统就几条消息流转,团队又完全没有Pulsar运维经验,那直接用云上的MQ服务或者Redis Stream就足够了。二是对延迟极其敏感且数据规模又不大,比如交易链路那种对延迟要求到毫秒级且每天都加机器扛峰值的场景,Kafka在性能调优和生态成熟度上更占优势。三是团队对BookKeeper本身就一窍不通,又没有预算和时间去学一套新的存储系统,这种时候贸然上Pulsar运维会很痛苦。
选型和谈恋爱有点像,不是选“最好”的,而是选“最合适”的。Pulsar在架构理念上确实先进,但先进不等于适合所有团队。技术栈切换是有成本的,团队的学习成本、基础设施能力、运维体系的搭配缺一不可。如果你们已经用Kafka用得很顺,也没有强烈的痛点,那没有必要为了追新技术而换。反过来,如果你看到了Kafka在积压、扩容、多租户上的天花板,并且这些天花板已经开始影响业务发展了,那Pulsar值得你马上开始PoC验证。
我个人在实际项目里的体会是,Pulsar最舒服的启动方式不是一上来就搞全量迁移,而是先把一个非核心但真实具备规模特征的业务切过去,比如埋点日志、异步通知这类。攒一个季度的运行数据,验证稳定性、性能、运维工具链都符合预期,再逐步扩大范围。这个思路适合所有中大型系统,稳妥永远是第一位的。