
一篇关于消息中间件选型的文章核心重点应该落在“为什么选”和“怎么落地”上而不是堆砌官方文档式的功能清单。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对应的BrokerBroker接到消息后把消息写入BookKeeper的当前LedgerBookKeeper完成多副本写入后返回确认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上的SegmentSegment会按大小滚动系统可以方便地管理这些数据块。即使消息在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控制副本数一般设3managedLedgerDefaultWriteQuorum和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最舒服的启动方式不是一上来就搞全量迁移而是先把一个非核心但真实具备规模特征的业务切过去比如埋点日志、异步通知这类。攒一个季度的运行数据验证稳定性、性能、运维工具链都符合预期再逐步扩大范围。这个思路适合所有中大型系统稳妥永远是第一位的。