消息中间件这玩意儿,平时在业务代码里看不见摸不着,但一旦流量上来,你就知道它有多重要。所以看到 COSCon‘25 同场活动 Pulsar Developer Day 的议程正式发布时,我第一反应是:这届开源大会是真懂开发者的痛点。作为在消息队列坑里摸爬滚打好几年的人,我打算借这个活动的话题,把我对 Pulsar 和消息中间件的一些实践思考梳理一遍,既是给要去现场的朋友划重点,也是给没空参会的同学一份技术笔记。
先说清楚,这篇文章不是议程复读机,而是围绕“消息中间件创新实践”这个核心展开。不管你是刚接触 Pulsar 的新手,还是已经在生产环境里折腾过 Kafka、RabbitMQ 的老手,都能从里面找到一些可以直接拿走的思路和避坑经验。
1. 为什么消息中间件值得一场专属开发者日
1.1 同场活动的含金量
COSCon 在国内开源圈的地位不用多说,它是开源社主办的一年一度的大会,覆盖的议题从操作系统到 AI 基础设施都有。而 Pulsar Developer Day 能作为同场活动出现,背后传达的信息其实很明确:Apache Pulsar 在中文开源社区里已经有足够多的用户基础,而且这些用户不只是拿来跑个 demo,是真的在生产环境里把它当核心基础设施用。
我接触过的不少团队,早期选消息中间件都是从 Kafka 入手的,毕竟资料多、生态熟。但业务发展到一定阶段,会开始遇到几个躲不开的问题:存储成本越来越高、扩容要动集群整体、多团队共用集群时权限和隔离特别难搞。这时候 Pulsar 的“存储计算分离”和“多租户原生支持”就会显得格外有吸引力。这场开发者日之所以值得关注,就是因为它把 Pulsar 社区里最活跃的一批实践者聚在一起,聊的全是这些真实问题。
1.2 议程背后是真实行业痛点
如果说 COSCon 主会场是开源生态的“全景图”,那 Pulsar Developer Day 就是消息中间件领域的“专题深挖”。从标题透出的“创新实践”四个字,结合这几年社区的热门话题,我推测议程大概率会覆盖这样几个方向:多租户隔离方案、分层存储降本、跨地域复制容灾、Pulsar 与 Flink/Spark 的流批一体集成,还有生产环境下的性能调优。
这些方向不是拍脑袋想出来的,而是 Pulsar 用户社群中讨论频率最高、最影响落地决策的话题。比如分层存储,很多团队数据量上来了以后,Kafka 的存储成本会让他们半夜惊醒,而 Pulsar 可以把老数据自动卸载到对象存储,这一块的实践分享对成本敏感型业务来说就是真金白银。再比如多租户,中大型公司内部往往有几十个业务线,如果每条线都单独部署一套集群,运维规模会爆炸,Pulsar 原生就能在同一个集群里划出隔离的命名空间,这种实践特别适合被拿到开发者日上展开聊。
2. Pulsar 架构里最值得吃透的几个核心设计
2.1 存储与计算分离:一条消息的完整旅程
Pulsar 的架构精髓可以浓缩成一句话:Broker 无状态,BookKeeper 管存储。什么意思呢?我拿 Kafka 做个对比你就明白了。Kafka 的 Broker 自己既处理请求又要管磁盘上的分区数据,扩容的时候数据要跟着分区迁移,这个过程既慢又容易出问题。Pulsar 则是把“接收消息、处理订阅”和“持久化消息”这两件事彻底拆开。
一条消息从 Producer 发出来,先打到 Broker,Broker 立刻把它交给底层的 BookKeeper 集群完成持久化,然后 Broker 再把消息推给 Consumer,或者等 Consumer 来拉取。这里的 Broker 只保留非常短暂的缓存状态,真正的数据住在 BookKeeper 里。想象一下,这就像餐厅里点菜的服务员和后厨的厨师完全分开:服务员只管接单、传菜,菜怎么做、怎么保存是后厨的事。客人越来越多,你只要多加服务员就行,后厨那边不需要跟着一起动。
这个设计带来的直接好处就是扩容变得非常优雅。你想给集群加 Broker,只要把新节点挂上去,不需要做任何数据搬迁;你想给存储扩容,那就加 BookKeeper 节点。存储和吞吐各自独立伸缩,这是单体架构的消息中间件做不到的。
2.2 多租户与资源隔离:集群共享不吵架
Pulsar 的资源模型是三层:Tenant(租户)→ Namespace(命名空间)→ Topic。这个抽象层级很值得玩味。Tenant 通常对应一个部门或一个大业务线,Namespace 对应具体业务模块,Topic 就是你我熟悉的队列通道。
我见过很多团队在 Kafka 里做隔离是怎么做的?要么一个业务线一套集群,要么大家挤在一个集群里用不同 topic 前缀硬分,前者浪费机器,后者出事就是全链路雪崩。Pulsar 的多租户是原生设计,租户之间有认证鉴权和配额管理,命名空间之间可以单独设置消息保留策略、Backlog 上限、生产消费速率限制。
实操中最直观的场景是这样的:公司有交易、日志、推荐三个业务线,它们在同一个 Pulsar 集群里各用各的 Tenant。交易线流量突增,最多影响它自己那个租户下的资源配额,日志和推荐线稳稳当当。这一点在大规模内部基础设施共享时特别重要,它直接把“运维多少个集群”和“有多少条业务线”这两个问题解耦了。
2.3 分层存储:给无限增长的消息找个廉价仓库
消息中间件有个让运维头疼的问题:消息数据只能持续增长,不像数据库还能归档删表。Kafka 的日志保留策略是把数据留在 Broker 本地磁盘,想要保留时间长一点就得堆昂贵的高性能盘。Pulsar 给出了一个思路完全不同的方案,叫分层存储。
Pulsar 允许你给命名空间设置一个“卸载”策略:消息在 BookKeeper 里读写的频率高,就留在高性能存储里;当这条消息超过设定时间(比如三天)或者积压量达到阈值,数据会被自动搬到配置好的对象存储里,比如 AWS S3、阿里云 OSS、腾讯云 COS。消费者再读老消息时,Broker 会自动去对象存储里捞,对客户端完全透明。
这个能力对于“消息要长时间保留给离线分析用”的场景可以说是刚需。Pulsar 官方宣称可以做到无限流存储,听起来玄乎,实际上就是用了对象存储几乎无限扩容的特性。我在实践里的体会是,这个功能不一定所有业务都用得上,但只要用上,每月云盘账单直接下降一个数量级。
2.4 跨地域复制:容灾不是靠运气
大厂做容灾讲究的是“两地三中心”,消息中间件要支持跨机房数据同步。Pulsar 的跨地域复制是一个非常成熟的功能,它允许你配置两个集群之间的命名空间镜像,消息在一个集群写入后会异步复制到另一个集群,而且这个复制是双向的,每个集群既能读也能写。
这个机制说起来不复杂,但它的价值在于容灾切换时你不需要改客户端配置。因为 Pulsar 的复制是端到端的,客户端连接的是本地集群的地址,本地集群故障了,DNS 切换到另一个集群,消息还在,因为 Replication 已经把数据同步过去了。相比有些方案要用 MirrorMaker 这样的额外组件去搬数据,Pulsar 内置机制省掉了一层运维复杂度。我在现场听同行分享时,最常听到的反馈就是“跨地域复制配置简单,但真要演练故障切换时,才体会到它的好”。
3. 从原理到参数,Pulsar 落地必看的实操要点
3.1 五分钟本地跑起一个 Pulsar 实例
刚开始接触 Pulsar 最友好的方式,是用 Docker 跑 standalone 模式,一条命令就能把 Broker 和 BookKeeper 一起拉起来:
docker run -d \ --name pulsar \ -p 6650:6650 \ -p 8080:8080 \ apachepulsar/pulsar:3.3.0 \ bin/pulsar standalone6650 是消息协议端口,8080 是管理接口和 admin API 的端口。跑起来之后,可以用官方命令行工具测一下消息收发:
# 消费(先起一个终端) docker exec -it pulsar bin/pulsar-client consume my-topic \ --subscription-name my-sub \ --num-messages 0 # 生产(再开一个终端) docker exec -it pulsar bin/pulsar-client produce my-topic \ --messages "hello-pulsar"这条路径走通之后,你对 Pulsar 的“客户端→Broker→BookKeeper”链路会有一个直观的体感。本地环境的用途是学 API 和验证想法,真要搞生产集群,建议用官方推荐的 Helm Chart 部署在 Kubernetes 上,或者用裸机部署,但无论哪种方式,都要记住一个原则:Broker 节点不要和 BookKeeper 节点混部,否则存储计算分离的意义就打了折扣。我见过有人为了省机器把两者混在一起,结果一个大流量的 Topic 直接把整台机器的 IO 打满,Broker 和 BookKeeper 互相争资源,最后全部超时,这就是没搞懂架构设计初衷的典型反面案例。
3.2 关键配置项的正确打开方式
配置 Pulsar 最核心的文件是broker.conf和bookkeeper.conf,但真正决定集群行为的是你在命名空间上设置的策略。我重点说三个最常被忽视的配置:
第一,消息保留策略。Pulsar 默认会消费完就删除数据,但对很多业务来说,消费完的数据还要留给离线分析,这时候要设置“Retention”,让消息在消费完之后继续保留一段时间。注意,Retention 和 TTL 是两回事:TTL 是消息没被消费多久之后可以被跳过,Retention 是已经被确认的消息再留多久。我见过很多新手把这两个概念搞混,结果消息全被清了还不知道为什么。
第二,Backlog 配额。这是用 Pulsar 多租户必须设置的护栏。如果不给命名空间设置maxBacklog或者maxSize,某个下游系统故障导致消息堆积时,集群的 BookKeeper 磁盘会被慢慢吃满,最终影响同一集群里的其他业务。设置配额之后,超过限制的 Pulsar 可以自动降级处理,比如拒绝新消息写入,这样故障的影响范围就被限制住了。
第三,生产端的 Batching 参数。Pulsar 客户端默认有消息批处理,把多条小消息合成一批发送,提升吞吐。但你在低延迟场景下一定要注意batchingEnabled和batchingMaxPublishDelayMs这两个参数。默认的延迟是 10ms,如果业务对延迟特别敏感,比如要毫秒级响应,那就要关闭批处理或者调低这个延迟阈值,代价是吞吐会降一些。这个取舍必须基于业务指标来定,不做压测就调参属于自找麻烦。
3.3 四种订阅类型到底该怎么选
Pulsar 的订阅模型是我认为它比 Kafka 做得好用的地方。Kafka 的消费组模式在互斥消费场景下很自然,但一旦你需要做“广播 + 工作组”的混合模式,或者想让同一条消息被多个消费者各自独立处理,就要把一条消息复制到多个 topic 里,非常别扭。Pulsar 直接在订阅类型上解决这个问题。
Exclusive 模式是排他的,一个订阅只能挂一个消费者,适合要求消息严格顺序处理的场景。Failover 模式是主备模式,主消费者挂了,备消费者顶上,顺序性也有保证。Shared 模式是共享队列,消息会轮询发给订阅下的所有消费者,适合快速处理大量消息但不在乎顺序的场景。Key_Shared 模式是前两者的折中:相同 key 的消息总是发给同一个消费者,既保证了同一业务实体的顺序,又实现了水平扩展。
选型建议很简单:需要严格全局顺序就用 Exclusive 或 Failover,需要吞吐和并行就用 Shared,需要兼顾业务单实体顺序和并行度就用 Key_Shared。这里不需要纠结,你只要在场景里想清楚“这条消息如果不按顺序处理会发生什么”,答案立刻就出来了。
4. 消息堆积、磁盘飙升与连接抖动排查实录
4.1 消费端处理慢导致的堆积问题
消息堆积是消息中间件最常见的事故,也是 Pulsar 比较让人省心的地方。为什么说省心?因为 Pulsar 的 Backlog 统计非常直观,你可以用一条命令看到每个 Topic 的堆积情况:
docker exec -it pulsar bin/pulsar-admin topics stats my-topic重点看几个字段:msgBacklog表示积压消息数,backlogSize表示积压消息大小,storageSize是存储占用。有一次我在生产环境碰到处理延迟飙到几分钟,一查 stats,发现有一个消费组的 Backlog 从几百涨到了几十万。
当时我的排查路径是:先确认消费组是否在线,再看消费者有没有大量 Redelivery 消息。结果发现是消费逻辑里调用了一个外部第三方接口,那个接口偶尔会超时,而消费代码没有做超时保护,导致一批消息反复被重试。解决方案是在消费逻辑里加断路器,超时的消息先落到本地表,等外部服务恢复了再补偿处理。这类问题其实和 Pulsar 本身关系不大,但 Pulsar 的重投递机制会放大外部依赖的抖动,这是所有消息消费端都要记住的教训。
4.2 BookKeeper 的磁盘与性能问题
BookKeeper 是 Pulsar 的存储底座,它的性能瓶颈往往出现在磁盘写入上。BookKeeper 有两种磁盘角色:Journal 和 Ledger。Journal 类似数据库的 Write-Ahead Log,要求低延迟、高吞吐,最好用 SSD,并且挂载时要使用独立目录;Ledger 是实际数据块存储,可以用机械盘或者对象存储做冷备替换。
我在排查磁盘 IO 飙升问题时,通常会做三件事:第一,检查journalSyncData配置,如果追求极致吞吐可以设为 false,但一旦机器断电就可能丢少量数据,这个参数要结合实际断电风险来权衡;第二,检查服务端是否频繁执行 GC,因为 BookKeeper 的读写路径上有大量 ByteBuf 分配,JVM 参数不调优的话 Full GC 会直接影响写入毛刺;第三,检查数据写入副本数,默认的 Ensemble Size 和 Write Quorum 都是 3,如果为了省存储改成 2,要明白这意味着允许一台机器故障,但如果本是三副本的地盘只有两副本,数据冗余度是降级的。
4.3 客户端连接抖动与鉴权配置
很多初次用 Pulsar 的人会遇到一个奇怪的现象:生产客户端偶尔报超时,但集群监控一切正常。我当时排查到最后发现,问题出在客户端和服务端之间的 TCP 长连接上。Pulsar 默认会让空闲连接保活,但有些公司的防火墙策略会把空闲连接断开,客户端不知道,还傻乎乎地在断掉的连接上发消息,自然就超时了。
解决方案其实不复杂:在客户端配置里开启keepAliveIntervalSeconds,并且实现重连逻辑。Pulsar 官方 Java 客户端自带重连机制,但生产端最好在代码里捕获PulsarClientException,对可重试的异常做指数退避重试。另外关于鉴权,Pulsar 支持 JWT 和 TLS 双向认证,我建议即使是内网环境也要开启基础鉴权,消息中间件一旦裸奔,任何人都能往 topic 里塞数据,轻则数据污染,重则被塞满磁盘,这种事在开发环境里见过不止一次了。
4.4 一份可收藏的故障速查表
| 现象 | 常见原因 | 排查命令/手段 |
|---|---|---|
| 消费延迟持续升高 | 消费者处理慢、外部依赖超时 | 检查 Redelivery 和消费异常日志 |
| 消息生产吞吐上不去 | 批处理未开或缓存队列太小 | 调整 Batching 参数和pendingQueueSize |
| BookKeeper 磁盘告急 | 保留策略过长或未设 Backlog 配额 | pulsar-admin namespaces get-retention |
| 客户端频繁断连 | 防火墙空闲连接回收 | 开启 keepalive,检查网络设备策略 |
| 消费重复但顺序对不上 | 订阅类型误用 Shared | 切换 Key_Shared,确认消息 key |
| 集群 OOM | JVM 堆内缓存过大 | 压测时盯 GC 日志,调整堆外内存 |
5. 写在使用 Pulsar 两年之后的一点心里话
大概两年前我第一次在生产环境引入 Pulsar,当时犹豫了很久,毕竟团队里已经有人熟悉 Kafka 了,再引入一套新系统,学习和运维成本都摆在那里。但业务压力逼着我们做选择,现在回头看,这个决定是对的。Pulsar 的多租户让我们的集群数量从十几套缩减到三套,运维负担大幅下降;分层存储让消息保留时长从 3 天延长到 30 天,云成本几乎没怎么涨;跨地域复制替我们解决了容灾演练时最头疼的数据同步问题。当然,它也不是银弹,某些场景下 BookKeeper 的运维复杂度和 JVM 调优成本是实实在在的,你必须有专门的精力去维护它,这不是一个装完就忘的组件。
如果你正准备从零评估消息中间件选型,或者想去 COSCon‘25 现场听听 Pulsar Developer Day 的内容,我的建议是:不要只看谁的吞吐量数字更高、谁的 Kafka 生态更丰富,先把你自己的业务场景列清楚——你需要顺序消费吗?需要广播模式吗?需要多部门共享一套集群吗?需要保留多长的消息历史?这些问题想清楚之后,你会发现选择并不难。最后再分享一个我个人的小习惯:每次看这类技术活动的议程,优先看带“实践”“调优”“踩坑”字样的议题,因为原理文档到处都能看,而真正让你少加班的东西,全在别人踩过的坑里。