Pulsar key_shared消息不消费排查与消息中间件选型实践
1. 从开源年会看 Pulsar Developer Day 的定位COSCon‘25 的消息一出来我第一反应是去看同场活动的完整排期。作为一个天天跟消息中间件打交道的后端开发者这几年我对 Pulsar 的关注度其实是持续上升的。今年的 Pulsar Developer Day 作为 COSCon‘25 的同场活动主打的聚焦消息中间件创新实践这个方向恰好踩中了很多团队正在面临的现实问题业务量涨了、Topic 多了、链路长了传统的消息队列开始在某些环节上力不从心而换不换、怎么换、换了之后有哪些坑这些都是需要有人讲清楚的事。这场活动的价值在于它不只是给你讲 Pulsar 的功能特性而是把开发者社区里正在发生的一线实践直接摆出来聊。我自己参加过几届类似的活动最大的感受是与其自己在群里问你们遇到过 key_shared 不消费的问题吗不如去现场听一听别人是怎么排查和解决的。尤其当你发现别人踩过的坑和你踩过的坑高度重合时你会非常确定这不是偶发问题而是有规律、有解法、有前置条件的真问题。对于手里正在做技术选型、或者已经上了 Pulsar 但用得不太顺的团队来说这一场的议程其实是很好的参考坐标。在我个人看来衡量一个消息中间件活动有没有干货就看三个维度有没有讲架构演进的、有没有讲源码和内核机制的、有没有讲真实故障排查的。三者如果齐全基本可以无脑参加。从 Pulsar Developer Day 过往的活动习惯和这次公布的方向来看这三个维度大概率都会被覆盖到。2. 议程背后的消息中间件创新方向由于具体议程条目在正式公布时通常还会有微调我这里结合这类活动的高频议题方向以及 Pulsar 生态近期的发展脉络聊一聊可以在现场重点关注的几个技术方向。这些方向不一定是今年议程的精确复制但它们反映的是 Pulsar 社区在过去一年里真正在做的事。2.1 存算分离架构Pulsar 的核心差异点很多团队初次接触 Pulsar 时最容易困惑的就是它和 Kafka 在架构上的本质区别。Kafka 的 Broker 既管计算又管存储分区副本落在 Broker 本地磁盘上而 Pulsar 把存储层独立出来由 Apache BookKeeper 承担Broker 变成无状态的接入层和调度层。简单类比一下Kafka 像是一个前店后厂的模式生产和仓库在同一个地方Pulsar 则是把门店和中央仓库分开哪个门店倒了换一个就是仓库里的货不会丢。这个架构带来的直接影响是扩容时你不需要做数据重平衡。Kafka 新增 Broker 节点后分区副本需要跨节点迁移这个过程在网络和磁盘 IO 上都有明显开销Pulsar 的 Broker 无状态之后新加节点只需要把流量引过去即可BookKeeper 存储节点也可以独立扩容。在活动的分享里通常会有团队用真实的扩容数据来说明这个优势尤其是夜间高峰或者大促前期这种能力就是实打实的稳定性保障。2.2 分层存储与成本治理我注意到近一年里凡是做 Pulsar 实践分享的团队几乎都会提到分层存储。Pulsar 的 Segment-based 存储天然支持将老数据卸载到对象存储比如 AWS S3、MinIO 或腾讯云 COSBroker 端只需要保留数据的指针信息真正读取老数据时再从对象存储拉回来。这对成本的影响非常大。举个例子假设你的集群每天产生 10TB 消息保留 3 天大约 30TB 的存储需求。如果全部用 SSD 或高性能云盘按云厂商的定价来说是一笔不小的开销。但如果你把保留时间拉长到 30 天而超过 3 天的数据自动沉降到对象存储存储成本可能下降到原来的十分之一甚至更低。关键是你还能继续消费这些老数据只是延迟会高一些。对于有重放历史消息需求的团队比如离线分析、故障现场回溯这个能力远比存三天就删要实用。2.3 多租户隔离和配额管理另一个值得关注的创新方向是多租户体系。Pulsar 的 namespace 天然支持隔离你可以给不同的业务线设置不同的存储配额、消息 TTL、Backlog 配额。这几年社区里也出现了很多在租户隔离基础上做的流量调度、速率限制的实践恰好是很多中大型团队在治理消息链路时的核心痛点。在活动议程里这类分享通常不会停留在我们启用了多租户这种层面而会往下讲清楚Quota 到底怎么设、达到阈值后是阻塞写入还是丢弃消息、不同业务的隔离策略为什么不同。建议你重点关注那些拿出了真实集群拓扑和治理策略的分享这种内容的复用价值最高。2.4 协议兼容与生态融合Pulsar 有一个常被忽视但战略价值很高的能力协议层兼容。通过 Protocol Handler它原生支持 Kafka 协议。换句话说你线上已经有 Java 的 Kafka Producer/Consumer代码一行不改把 broker 地址换成 Pulsar 集群的地址就可以跑通。这给存量系统提供了一条非常平滑的迁移路径。在活动上这类话题通常会伴随着迁移过程中遇到的各种兼容性问题来展开。比如 Kafka 客户端某个版本的行为和 Pulsar 的 Kafka 协议实现不完全一致或者在事务、幂等语义上的差异。对还在观望要不要从 Kafka 迁到 Pulsar的团队来说这是最值得做笔记的部分。3. key_shared 模式消息不消费一次典型的线上排查回顾前面聊了创新的方向现在必须回到一个非常现实的问题上。Pulsar 的 key_shared 模式不消费这个热词我是有切肤之痛的。如果你用过 key_shared 订阅模式大概也遇到过类似的现象消息明明已经生产到 broker 了consumer 连接也正常消费者组也显示 online但消息就是堆在 backlog 里不往下走。我第一次遇到这个情况时差点把集群 manager 叫起来一起排查后来发现根因还得从模式本身的机制说起。3.1 key_shared 的原理按 key 做一致性哈希你首先要理解 key_shared 的定位。它介于 Shared 和 Failover 之间Shared 模式是多个 consumer 一起消费同一个订阅里的消息谁抢到算谁的但有损消息级有序Failover 模式是单 consumer 干活其他 consumer 做备胎保证严格顺序但吞吐受限。key_shared 则允许你把消息按 key 分桶相同 key 的消息始终发送到同一个 consumer这样既保留了一定程度的消息局部有序又让多个 consumer 能并行处理不同 key 的消息。实现机制上Pulsar 对消息的 key 做哈希然后映射到一组 consumer 上。在实际代码里有一个KeySharedPolicy的概念它负责描述怎么在 consumer 之间分配 key 的范围。这是整个模式里最容易被忽略、但也最容易出问题的点。3.2 导致不消费的根因排行根据我自己踩坑和观察社区讨论key_shared 模式下消息不消费通常可以归为以下几类根因类别具体表现排查难度订阅类型配置错误consumer 实际用的是 Shared 模式但期望 key_shared低KeySharedPolicy 配置错误设置了不允许的组合选项策略未生效中consumer 与分区数量关系consumer 数超出分区数部分 consumer 拿不到 key 区间中哈希算法不一致客户端和 broker 之间 key 路由计算不匹配高版本 bug特定 Pulsar 版本中 key_shared 粘滞性异常中我自己遇到的一个典型案例是 KeySharedPolicy 里设置了AutoSplitHashRange同时又不小心把allowOutOfOrderDelivery也配置上了。在当时的版本下这两个选项的组合会让 consumer 侧获取不到正确的 hash 区间消息虽然进了订阅的 backlog但没有任何一个 consumer 认领这批 key。表现上就非常诡异订阅的msgBacklog持续上涨consumer 列表正常client 端的 listener 却毫无动静。3.3 完整的定位链路从现象到根因如果你想在线上快速定位这类问题我建议你按下面的顺序来排查而不是一上来就去翻 broker 日志。第一步确认 subscription 的实际模式。用pulsar-admin或控制台查看订阅类型bin/pulsar-admin topics stats-internal persistent://public/default/your-topic重点看type字段确认它到底是不是Key_Shared。这一步能排除掉大量我以为我用了 key_shared的情况。第二步确认 consumer 的连接状态和 hash 分配。在 topic 的统计数据里找到currentActiveConsumers对比你的实际 consumer 数量。如果两者不一致说明有 consumer 没有成功加入 key 分配。这时候去看 client 侧日志通常会看到类似Cannot get the key hash range或KeySharedPolicy相关的报错。第三步检查消息的 key 是否真实存在。key_shared 依赖消息 key 做路由如果你的 producer 发送消息时根本没设置 key消息会走什么逻辑在 key_shared 模式下无 key 消息会被随机路由到某个 consumer但如果所有 consumer 的 hash 区间都没有认领它们同样会堆住。用pulsar-admin topics peek-messages抽查 backlog 里堆积消息的元数据很容易就能看出 key 字段是不是空的。第四步检查版本兼容性和已知 bug。这一步往往最容易被忽略。Pulsar 的版本迭代中key_shared 相关的 bug 出现过不止一次。我的习惯是直接到 GitHub 的 issues 里搜索key_shared加你的版本号看有没有人报过类似问题。如果确实命中已知 bug优先升级 broker 和小版本客户端通常能解决一大半问题。最后如果以上步骤都查不出问题才需要去翻 broker 端的ManagedLedger日志确认有没有 ledger 的读取异常。但说实话90% 的 key_shared 不消费问题在第三步之前就已经能定位了。3.4 线上规避措施即使你定位并修复了问题我依然建议你在关键链路上加上一层保护。最简单有效的办法是给关键 topic 设置一个backlogQuota告警当堆积阈值触发时立刻触发人工介入而不是等业务侧来反馈消息怎么还不出来。另外生产代码里对 KeySharedPolicy 的构造建议集中封装统一控制参数组合避免不同团队各自为政、埋下隐藏的策略配置差异。4. 结合热点推导出的实用选型与避坑清单正因为在 key_shared 上踩过这么深的坑所以我一直觉得这类活动里案例复盘类内容的价值被严重低估。消息中间件的选型和使用很多时候不是功能堆得越多越好而是你知不知道每条功能背后的代价。4.1 三种共享订阅模式怎么选很多刚接触 Pulsar 的人会在 Shared 和 Key_Shared 之间纠结这里我直接给出我的选择逻辑如果消息之间完全没有顺序要求处理逻辑允许任意并发选 Shared 最省心吞吐也最好。如果消息必须按 key 有序处理比如同一个订单 ID 的后续操作必须严格按照顺序执行选 Key_Shared。如果整个 topic 的消息都必须按严格全局顺序消费那就只能选 Failover 或 Exclusive不要指望 key_shared 能帮你。这里面最需要注意的取舍是Key_Shared 的有序是有前提的。相同 key 的消息必须由同一个 consumer 串行处理但当你动态增减 consumer 时key 与 consumer 的映射关系会重新洗牌过程中可能短时间出现重复或乱序。所以如果业务对顺序极度敏感不要依赖 key_shared 的动态伸缩而应该设计成固定 consumer 数量或者干脆走 Failover。4.2 实验环境验证思路如果你现在没有 Pulsar 环境又想快速验证 key_shared 的行为我建议直接用 Docker 起一个 standalone 集群然后用官方客户端写个小 demo一个 producer 发 1000 条带 key 的消息两个 consumer 用 key_shared 订阅来消费打印各自收到的 key 列表。你会发现 key 的分布和 consumer 数量强相关。然后把 consumer 数量改成 4 个再跑一次看到的 key 区间分布又不一样。这个实验最好在 2.11 以上版本来做因为更早版本里 key_shared 的实现和后续版本差异较大一些 bug 修复也没有同步。你完全可以把这个实验当作上述活动热词的一个配套实操先复现问题再理解机制最后看它在你自己的业务里需要注意什么。4.3 团队采用 Pulsar 的前置准备最后说一点比技术更重要的东西。Pulsar 的架构优势很清晰但它和 Kafka 的操作习惯差异很大。团队如果只是把 Pulsar 当成另一个 Kafka来用很容易在存储治理、BookKeeper 运维、namespace 规划上吃暗亏。我看到的比较成功的接入方式都是先选一两个核心业务跑起来同时配备有 BookKeeper 运维经验的人以及建立完整的 backlog、存储、延迟监控面板。等运行稳定之后再逐步扩展到更多业务。千万别一上来就做大规模全量迁移尤其是线上核心链路消息中间件的切换是最容易出连锁事故的环节。5. 一点个人体会回顾整个 Pulsar 生态的演进再对照key_shared 不消费这类热词背后的真实问题你会发现社区里真正有价值的讨论从来不是单纯报一个 bug 出来而是把现象、版本、配置、排查过程完整地呈现出来。这也是为什么我强烈建议有条件的话去 Pulsar Developer Day 现场或者关注后续演讲录播——你可以看到别人如何从现象出发层层拆解到根因这种思维方式的收获比单纯记下一个解决方案要持久得多。我个人在参加完这类活动之后总会做一件事把现场听到的案例与自己线上系统逐一对照整理成一份我们可能遇到的同类问题清单。像 key_shared 的坑如果你没有提前做过功课遇到时大概率要花掉半天到一天的时间去定位。但如果你已经在活动上看到过完整的排查链路同样的坑可能十几分钟就能解决。这大概就是社区分享最朴实的价值让别人踩过的坑成为你不需要再踩一遍的捷径。