1. 开场:这个周末,去开源集市找点“硬货”
周末的 COSCon'25 开源集市,Apache Pulsar 的展台就在那等着你。如果你正好在会场里转悠,想找一个既能聊实时消息、又能聊存储架构的摊位,Pulsar 那个展位值得多停一会儿。不光是领贴纸、换周边,更重要的是能直接跟项目的 committer 和 Contributor 面对面聊技术——这种机会,平时在线上社区里可不太容易碰到。
Pulsar 是什么,先给还没接触过的朋友一句话说清:它是一个云原生时代的分布式消息和流数据平台,从 Yahoo 内部孵化出来,后来捐给了 Apache 基金会。跟 Kafka 这类老牌消息队列相比,Pulsar 最大的不同是“存储和计算分离”的架构,这也让它天生更适合做多租户、跨地域复制和大规模分层存储。对正在折腾实时数据管道、微服务异步解耦、或者想统一消息和流两种场景的团队来说,Pulsar 是个很值得认真评估的选项。
这篇文章不打算做成活动流水账,我主要想借 “Pulsar 参展 COSCon'25” 这个由头,把 Pulsar 的几个核心技术点拆开讲清楚:它的存储架构强在哪,四种订阅模式到底该怎么选,以及一个最近社区里讨论得很热的真实问题——“key_shared 模式不消费”的 bug。这段排查经验是我在实际环境里踩过的坑,也应该是很多 Pulsar 使用者迟早会撞上的问题,提前知道能帮你省不少时间。
2. 先搞清楚 Pulsar 到底解决了什么问题
2.1 存算分离架构:Broker 和 BookKeeper 的分工
很多刚开始接触 Pulsar 的人会问,Pulsar 和 Kafka 到底有什么本质区别?我一般用一个类比来解释:Kafka 像是把“收银台”和“仓库”绑在一起开分店,每个分店自己管自己的货;Pulsar 则把所有“收银台”集中在前台,后面只有一个巨大的中央仓库,由专门的库管团队(BookKeeper)来负责分拣和存取。
具体到架构上,Pulsar 的 Broker 层只负责消息的路由、权限管理、订阅状态跟踪这些相对轻量的逻辑处理,而真正的数据存储落在 Apache BookKeeper 上。Broker 是无状态的,你可以随时加机器横向扩容,只要把新节点注册进集群,流量就能自动分摊过去。BookKeeper 的数据节点(Bookie)按 Segment(段)来存储消息,每个 Topic 的数据被切成一堆小 Segment,分散存储在多台 Bookie 上,并且每个 Segment 默认写多副本(通常是 3 副本)。这意味着某一台 Bookie 挂了,它的副本会立刻顶上来,数据不丢、服务不断。
这个架构带来一个很实在的好处:Broker 不需要本地磁盘存消息。Kafka 的分区数据是要落盘的,Broker 挂了换新机器,得先把数据重新同步回来,迁移成本高得很。Pulsar 的 Broker 挂了,新 Broker 直接接上,数据还在 BookKeeper 那边躺着,几乎是无感切换。对运维而言,这个差异在故障演练的时候体会最深。
2.2 Topic 可以很大、很多,且支持多租户和分层存储
另一个 Pulsar 让我觉得非常顺手的地方是 Topic 数量可以撑到百万级。早期的消息队列一旦 Topic 多了,整个集群性能都会跟着崩,原因在于每个 Topic 的元数据和索引都要占资源。Pulsar 在应用层把 Topic 抽象成了逻辑概念,实际物理上是一个“Ledger 上的 Segment 集合”,Broker 即使管理上万个 Topic,自身压力也远小于传统方案。
多租户方面,Pulsar 用了三层命名空间结构:Tenant(租户)— Namespace(命名空间)— Topic。比如一个 Topic 的完整名字是persistent://my-tenant/prod-namespace/my-topic。你可以在每个租户级别做隔离、配额、速率控制,几个业务团队共用一套集群,互不干扰。对中大型公司而言,这直接省掉了“每个团队单独搭一套 Kafka”的运维成本。
分层存储(Tiered Storage)也是我很看重的一项功能。默认情况下,数据按保留策略在 BookKeeper 里滚动删除,但只要你配了一个 S3 兼容的对象存储,Pulsar 就会自动把老的 Segment 数据搬到 S3 里。读的时候如果消费者想回溯旧数据,Broker 会从对象存储拉回来。这样你就不必为了“偶尔回溯一下三个月前的数据”而无限扩大 BookKeeper 集群的磁盘容量,成本控制非常直观。
2.3 从 Kafka 迁到 Pulsar 的体验观察
我身边有不少团队做过 Kafka 到 Pulsar 的迁移,反馈基本集中在几个点:一是 Pulsar 的消息消费模型更灵活,同一个 Topic 可以同时跑“队列式消费”和“流式发布订阅”;二是管理端做得比较清晰,pulsar-admin 命令行工具、REST API 和 UI 控制台都能操作;三是社区对生态连接器(比如 Pulsar IO Connector)的支持已经相当丰富,Source 和 Sink 一接就能用。
当然不是说 Pulsar 就面面俱优。它的架构让每个环节都更分布式,这意味着部署和运维的组件变多了:至少要部署 ZooKeeper(或者 etcd,新版已经支持元数据服务抽象)、BookKeeper、Broker 三套东西。小团队如果只有一两台机器,想轻量跑动一个功能完整的 Pulsar 集群,会比单机 Kafka 费劲不少。所以它的优势场景更偏向有一定规模、真正需要弹性扩容和稳定性的生产环境。
3. 消费模型的根本差异:订阅模式选对了,事半功倍
3.1 Pulsar 的四种订阅模式对比
在 Pulsar 里,订阅(Subscription)是消费行为的核心抽象。同一个 Topic,不同的订阅之间是隔离的,各自维护各自的消费进度。Pulsar 支持四种订阅模式:
| 订阅模式 | 消费语义 | 适用场景 |
|---|---|---|
| Exclusive(独占) | 一个订阅只允许一个消费者,其他消费者连接时直接报错 | 严格有序、单消费者场景,比如全局事件表、严格有序的状态变更流 |
| Failover(灾备) | 多个消费者里只有一个 Active,其他作为 Standby,Active 挂了自动切换 | 需要顺序保证且要求一定可用性的场景,比如金融交易事件 |
| Shared(共享) | 消息随机分发给多个消费者,各处理各的,不保证顺序 | 吞吐量大、对顺序不敏感的任务型队列,比如异步通知、批量任务分发 |
| Key_Shared(按键共享) | 相同 key 的消息只会路由到同一个消费者,不同 key 可以分散到多个消费者 | 既需要按 key 保证顺序,又想并行扩容的场景,比如用户维度的实时聚合、订单状态机流转 |
3.2 Key_Shared 到底应该怎么用
我把 Key_Shared 单独拿出来讲,因为它在“顺序”和“并行”之间找到了一个很好的平衡点。举例来说,你在处理一笔订单的实时数据流,订单 ID=10086 的“创建”“支付”“完成”三条消息,如果落到 Shared 模式下,很可能被两个消费者抢着处理,最后出现“支付先于创建”这种乱序;而用 Key_Shared,只要我们把订单 ID 作为消息的 key,这三条消息就一定会进同一个消费者,顺序自然保住了。与此同时,不同的订单 ID 仍然能分给不同消费者处理,吞吐量不会因为单消费者而受限。
在代码上,用 Java Client 指定 Key_Shared 订阅很容易:
Consumer<String> consumer = client.newConsumer() .topic("persistent://tenant/ns/my-topic") .subscriptionName("my-sub") .subscriptionType(SubscriptionType.Key_Shared) .subscribe();有一点必须提醒:Key_Shared 模式下,消息的分发是“尽力按 key 哈希落到消费者”,但它在设计上并不保证一个消费者独占一批 key 之后立即把该 key 的后续消息都交给它。也就是说,如果你既要严格顺序又要重构某些聚合状态,最好还是用 Exclusive 订阅或做补偿机制。生产上跑 Key_Shared,还需要开启消费者端KeySharedPolicy并合理配置哈希分片数量,默认的 2 个分片在小规模消费组下可能不够均匀。
3.3 踩过的订阅模式选型坑
以前接手过一个项目,业务方想用 Key_Shared 同时实现“按用户 ID 有序”和“多个消费者并行”。听起来很完美,跑了一阵子发现消费组里的某个消费者负载明显高于其他节点。查了半天,发现是消息 key 的基数太小——比如 key 只取了业务类型(下单、支付、退款),而业务类型总共就 3 种,哈希来哈希去就落在几个固定消费者上。这不算 bug,但很容易被误判成 Pulsar 的分发策略有问题。后来我把 key 设计成“业务类型_用户ID”组合,分布瞬间均匀了。
像这类负载不均的问题,如果想在测试阶段就发现,可以先用pulsar-admin topic stats看一眼每个消费者的msgRateIn指标,曲线能直接暴露出分发不均。别省这一步。
4. “key_shared 不消费”问题实录:现象、排查与解决
4.1 事故发生时的现场
这个周末开源集市的群里正好有人在聊“key_shared 模式不消费的 bug”,我一下子就想起了几个月前在生产环境踩过的同款情况。当时我们的服务大概是这样的:一个订单系统用 Pulsar Key_Shared 订阅消费订单变更事件,一共起了 8 个消费者,跑了一段时间后,突然监控报警,消费组的 Lag(积压)暴涨。但是诡异的是,你去查消费者的状态,每个消费者都显示“已连接”“正常”,订阅仍然存在,ranges 也还在,就是不出消息。
关键信息我再理顺一遍:不是消费者崩溃,不是订阅不存在,也不是权限问题。broker 端没有任何异常日志,Producer 那边发送消息的速率完全正常,消息也确实写进了 Topic,繁荣的消息就堆在那,但就是没有消费者拉到数据。
这种“连接正常但不消费”的状态,比直接报错难查太多了。
4.2 排查路径:从客户端到 Broker 再到元数据
我当时的排查顺序是这样的,每一步都有明确的排除目的:
第一步,先看客户端日志。在 Pulsar Client 里,消费事件会记录在 debug 级别,包括 Broker 分配的消费队列和 receiver queue 的 fill 情况。把日志级别临时调到 DEBUG,发现消费者端根本没有收到新的MessageReceived事件。这里就说明问题大概率出在“Broker 没把消息匹配给这个订阅”上下,而不是消费回调卡死。
第二步,看pulsar-admin的统计信息。执行:
bin/pulsar-admin topics stats persistent://tenant/ns/orders-events重点看subscriptions里的my-sub对应的msgOutCounter和msgRateOut。当时我看到的msgOutCounter完全不动,而msgBacklog一直在涨,这进一步把问题锁定在 Broker 该订阅的分发逻辑上。什么连接、心跳、TCP 层面的问题都可以先排除掉了。
第三步,查订阅模型和 Topic 策略。Key_Shared 模式有一个很容易被人忽略的约束:它要求 Topic 上的所有订阅都必须是“有序消费类型”,或者至少要让 Key_Shared 能够正常处理。如果同一个 Topic 上既有 Shared 订阅又有 Exclusive 订阅,某些老版本 Broker 对 Key_Shared 消费者的路由表处理会异常。当时我很怀疑是某个排查用的调试脚本在同一个 Topic 上建了另一个普通 Shared 订阅,从而干扰了 Key_Shared 的订阅状态。
第四步,也是最关键的,看元数据里的 consumer 分配情况。当时我们团队先在一个测试 Topic 上复现了这个问题:不用业务侧代码,直接用pulsar-perf consume加上 Key_Shared 订阅去消费,结果同样没消息。这就把问题缩小到了 Broker 源码层面的路由逻辑。查了相关 issue 和源码,终于明白:在较旧的 Pulsar 版本里,Key_Shared 的消息分发依赖 broker 为每个 Key_Shared 订阅维护“hash range 到 consumer”的映射关系。一旦消费者做 rebalance 或有新的消费者加入、退出导致的 hash range 更新不及时,Broker 会认为所有 range 都已经被“冻结”或没有可用的 consumer 而停止分配新消息。表现正是:consumer 连接正常,但消息不出。
4.3 修复方案与规避手段
最终我们升级了 Pulsar Broker 到修复该问题的版本(这里也建议读者在生产环境长期关注 Apache Pulsar Release Notes 中关于 Key_Shared 的 bugfix 条目)。如果短期没法升级,有几种临时的规避方式可以用:
- 在订阅不变的情况下,重启消费者进程,强制触发一次 rebalance。
- 把消费者的数量临时调整为 1,再用后台任务业务补偿把积压消费完,必要时再逐步增加消费者。
- 如果积压时间敏感,直接新开一个 Shared 订阅做兜底消费,把积压先追平,同时注意业务侧接受短时的乱序(或者在下游做 key 维度的窗口排序)。
- 另外要检查客户端的版本和 Broker 版本是不是相差太大,新客户端连老 Broker 的 Key_Shared 订阅在某些协议版本组合下兼容性并不可靠。
这个 bug 给我的教训很深:凡是依赖 Key_Shared 这种“按需有序”模型的业务,必须配套一套“乱序兜底”的应急方案,比如在下游加一个 keyed 合并缓冲层,至少不能让消息彻底卡死在那里。消息系统再可靠,也要对极端 bug 留一手。
4.4 给你的 Key_Shared 使用建议清单
作为一个被坑过的人,我现在总结了一套 Key_Shared 上生产前的检查清单:
- 确认业务的顺序语义真的需要 key 粒度。如果你接受全 Topic 乱序,直接用 Shared,别多此一举。
- 确认消息 key 的基数足够大,避免少数几个 key 占满消费者。
- 不使用同一个 Topic 混合不同订阅模型,尤其不要混用 Shared 和 Key_Shared。
- 所有消费者使用相同版本 Client,并尽量与 Broker 保持同 minor 版本。
- 为订阅配置合理的
maxUnackedMessagesPerConsumer,否则消息不 ack 会把消费卡死。 - 持续关注 Pulsar 社区关于 Key_Shared 的 issue,因为这个功能在多个版本里都有过不同程度的 bug。
5. 从 COSCon 开源集市聊到生产实践里的 Pulsar 细节
5.1 活动现场可以聊什么
如果你周末真的去 COSCon'25 开源集市了,建议别只领贴纸就走。Pulsar 展台通常会有以下几样值得深度交流的东西:
一是项目核心开发者对架构演进方向的最新解读,比如 Pulsar 3.x 里对元数据服务、负载均衡、Broker 自动故障转移的改进。平时在 Release Notes 里看到的文字,远不如现场听作者讲一遍“当时为什么要这么设计”来得通透。
二是关于 Pulsar 生态工具链的演示,比如 Pulsar Functions(轻量流处理)、Pulsar IO Connector、Pulsar SQL(用 Presto 查 Topic 数据),这些组件能让你在不用引入一套 Flink 或 Spark 的情况下,就做不少实时计算和查询的事儿。
三是最重要的——现场“诉苦”环节。开源社区最宝贵的资源,其实是那些踩过坑的真实用户。你可以拿着自己的配置、监控截图、问题栈,直接跟 committer 聊。我上一次在开源大会的 Pulsar 展台,就靠着一张“消费者连接正常但不消费”的截图,让开发者帮我定位到了配置上的问题——那种效率远高于自己在 issue 列表里翻半天。
5.2 生产环境里容易被忽视的 Pulsar 配置细节
借此机会,我再把几个生产环境容易被大家忽视的 Pulsar 参数拎出来聊一聊。
第一个是确认ackTimeout和ackTimeoutRedeliveryBackoff的关系。默认情况下 ack timeout 如果设置得太短,而业务处理偶发超过这个时间,就会造成消息被重复投递,消费者被迫做大量幂等逻辑。如果你对消息时延没那么敏感,建议把 ackTimeout 调到业务处理 P99 时长的 5 倍以上,或者干脆关闭自动 ack timeout,改用negativeAckRedeliveryBackoff配合死信重试。
第二个是managedLedgerDefaultEnsembleSize、managedLedgerDefaultWriteQuorum、managedLedgerDefaultAckQuorum这三个 BookKeeper 参数。默认值(通常是 2/2/2 或 3/3/3)决定了 Segment 的副本数。如果集群里只有 3 台 Bookie,你配了 3 个副本,坏一台机器后可用性会受影响。副本数也直接决定磁盘用量和写入吞吐,先想清楚业务允许丢失多少数据,再定这几个值。
第三个是消费者端的receiverQueueSize。它相当于是你的消费窗口大小。队列太小,消费吞吐上不去,在 Pulsar 上面表现为“平均每秒钟只能消费几千条”;队列太大,一条慢消息会把整个消费窗口卡住几秒,下游 RT 瞬间冲高。这个值一般建议设在 1000 左右,高吞吐场景可以试 5000,但一定要配合maxUnackedMessagesPerConsumer一起调。
第四个是关于“Backlog 管理”。不少人只关注消费 Lag,没想过积压消息如果迟迟不消费,BookKeeper 磁盘会一直不释放。Pulsar 的retention策略和ttl策略需要配合使用:先设 TTL,让过期消息自动清理;再设 retention,保证已经 ack 的消息还能被回溯一段窗口。如果都不设,消息永远删不掉,磁盘迟早撑爆。
5.3 开源集市上的 Pulsar 周边,以及如何带一个问题回家
往年 COSCon 开源集市的惯例是,每个项目展台都会有印章收集或者小任务,集齐几个可以换周边。Pulsar 的周边经常是飞行棋、贴纸、帆布袋一类,质感在同行的开源项目里算不错的。但对我来说,比周边更有价值的,是带着“现场验证一个真实猜想”的心态去聊。
比如你可以在去之前,先在本地环境里起一个 Pulsar Standalone 集群:
# 下载并解压 Pulsar 后,直接启动单机模式 bin/pulsar standalone --num-bookies 1然后建一个 Topic,用 Java 或 Python 客户端分别测试 Exclusive、Shared、Key_Shared 三种订阅模式下的行为。如果你想提前踩一下上面说的那种坑,可以把 Broker 版本固定在某个比较旧的 2.x 版本,然后频繁增删消费者观察 hash range 的变化。现场把这些发现跟 Pulsar 的开发者一聊,你会发现自己对消息系统的理解立刻高了一个level。
6. 常见问题速查:现场最容易听到的 Pulsar 提问
结合活动现场和社区群里听到的高频提问,我做了一个“速查表”,大家遇到类似问题可以直接对着查:
| 现象 | 常见原因 | 解决方案 |
|---|---|---|
Producer 发送报Topic not found | 没有启用自动创建 Topic 或权限不足 | Broker 配置allowAutoTopicCreation=true,赋予对应 role 的生产消费权限 |
| 消费有延迟但 CPU 不高 | ack 太慢,单消费者内积压了太多 unacked 消息 | 增大receiverQueueSize,调大maxUnackedMessagesPerConsumer |
| 消费重复消息明显变多 | ackTimeout 太短,消息未处理完被重复投递 | 调大 ackTimeout,配置重试退避而非立即重投 |
| 单 Topic 吞吐无法突破 | 分区数太少 | 创建分区 Topic,确认分区 Producer 是否按 key 正确路由 |
| 某个消费者永远没消息 | 可能是 key 哈希全部被其他消费者占用,或处于 pending rebalance | 检查msgRateIn和consumerName,触发一次 rebalance 观察 |
bookie磁盘使用率持续增长 | 没有 TTL 或 retention 策略,旧消息不被清理 | 设置ttl和retention,配合pulsar-admin topics delete清理无效 Topic |
在 COSCon 展台现场,另一个问得非常多的问题是“Pulsar ZooKeeper 挂了怎么办”。这个问题的本质是:Pulsar 只用 ZooKeeper 管理元数据,并不承载消息数据的读写。ZooKeeper 短暂不可用,Broker 和 Bookie 之间的现有连接还在工作,消息数据读写不受影响,只是元数据变更(比如 Topic 创建、订阅状态变更)会暂时不可用。所以生产环境把 ZooKeeper 组件单独部署、监控告警配好,一般不会出现整个消息集群不可用的极端情况。
7. 聊聊我对 Pulsar 在开源社区发展的一点感受
在 COSCon 这类开源集市上,我观察到一个很明显的变化:前几年摆摊的大多还是 Web 框架类项目,大家问的是“这个东西能干什么”;现在越来越多像 Pulsar 这样的基础软件项目出现在集市上,来看的人直接拿着一张集群拓扑图或者一堆异常监控截图来问“这里该怎么调”。这个变化背后,是基础架构层面的人才密度在上升,也是大厂开源项目逐渐走向生产可用、走向成熟的信号。
如果你周末到了现场,大概率会看到 Pulsar 展台前围着的人分成两类:一类是刚入行、眼里发光的学生,另一个是带着生产环境难题来的老工程师。这两类人之间的对话,往往才是开源集市上最精彩的环节。中间的技术布道者会反复解释同一件事——为什么消息系统需要存储计算分离,为什么流和队列可以在一个平台里统一,为什么 Pulsar 愿意把自己最核心的 BookKeeper 存储层单独开源出来。这种“愿意剖开自己、让更多人理解底层设计”的氛围,是我觉得国内开源社区最珍贵的东西。
又想起那个 key_shared 的 bug。其实任何一个能扛住大规模生产流量的开源项目,它的成长路径里都写满了类似的 bugfix。能把它拿到台面上,在 issue 里、PR 里、社区活动里被反复讨论,本身就是项目生命力的体现。这也是为什么我更推荐大家去现场和开发者聊聊——你踩到的每一个坑,都有可能是推进下一次版本迭代的一块砖。
周末如果有空,来 COSCon'25 开源集市逛逛吧。可能你站在那里,跟 Pulsar 开发者聊了十分钟,带回的不只是周边,而是一整套关于自己系统架构的新思路。