上半年最忙的一段时间刚过去,趁着记忆还新鲜,把 COSCon‘25 和 Pulsar Developer Day 2025 合办的专场里那些让我印象深刻的议题,结合我自己在生产环境折腾 MQ 的实战经验,系统地梳理一篇。这次活动最核心的几个话题,其实都是消息中间件领域这些年最扎心的痛点:先写数据库还是先发 MQ、Pulsar 的 Key_Shared 模式为什么会出现不消费的怪现象、消息补发要怎么做才优雅,以及很多人忽略的算术编码在消息压缩里的原理。如果你正在做消息架构选型,或者已经被线上消息积压、重复消费折磨过,这篇文章值得你花 15 分钟读完。
1. COSCon 2025 现场:大家为什么盯着 Pulsar 和 MQ 不放
1.1 一个合办专场背后的技术风向
这次 COSCon 2025 和 Pulsar Developer Day 2025 合办的消息专场,来的人比我想象中多得多。会场里一问,基本都是已经在生产环境用了三四年 Kafka、RabbitMQ,最近开始认真调研 Pulsar 的团队。过去大家聊 MQ 都是问"怎么保证不丢消息",现在话题明显往"怎么在复杂架构里把消息用得优雅"这个方向倾斜了。
一个很明显的信号是,现场提问环节被问得最多的问题不再是"Pulsar 和 Kafka 谁吞吐高",而是"我们想把核心交易链路迁到 Pulsar,双写一致性怎么处理"以及"Key_Shared 模式扩容后为什么有的分区不消费了"。这类问题在官方文档里能找到答案的不多,基本都得靠实战踩坑总结,所以大家宁愿挤在专场里听从业者聊真实案例,也不愿意回去翻那些讲得过于理想的架构图。
我自己这几年做消息中台,也经历了从单一 Kafka 集群到多 MQ 组件共存、最后把核心链路逐步收敛到 Pulsar 的过程。这次专场分享的很多议题,和我实际踩过的坑几乎完全重合,尤其是 Key_Shared 模式和双写一致性问题,几乎每个做订单、支付、库存相关系统的团队都会遇到。
1.2 存算分离架构为什么成了 MQ 的转折点
Pulsar 被讨论得越来越多,核心原因是它的存算分离架构确实切中了大规模消息系统的痛点。传统 MQ 把数据存储在 Broker 本地磁盘,扩容一台 Broker 往往意味着要搬迁数据,分区数、副本数、磁盘水位全都要重新评估。Pulsar 把存储层抽出来放到 BookKeeper,Broker 只负责计算和路由,扩容就是加 Broker 节点,存储层独立扩展,两者互不拖累。
这个架构带来的最直观好处是,你不再需要为"未来半年要涨多少流量"提前囤机器了。我之前维护过一个 Kafka 集群,大促前要做容量评估,因为 Broker 和存储绑定,扩容窗口期长,经常要提前两三周做准备。换到 Pulsar 之后,流量预估偏差的影响小了很多,Broker 层两三天就能扩完,存储层按 BookKeeper 集群的容量规划走,节奏从容得多。
当然存算分离也有代价。BookKeeper 本身是一套分布式存储系统,运维复杂度明显比单机磁盘高,网络延迟、节点故障、Ensemble 大小这些参数都要重新学习。但对于日消息量在十亿级以上的核心链路,这个代价是值得的。专场里不少分享者都提到,他们的迁移过程不是一次切换,而是通过 Pulsar 的跨集群复制功能先做双跑,逐步把流量切过去,这个思路我自己也非常认同。
2. 先写数据库还是先写 MQ:双写一致性没有银弹
2.1 两种写入顺序各自的风险模型
这是消息场景里被问烂了、但又永远绕不开的问题。先说结论:无论先写数据库还是先发 MQ,单靠一种顺序都无法保证绝对一致。我们来仔细拆一下两种方案的风险。
先写数据库,再发 MQ,这是大多数团队默认的选择。理由很朴素:数据库是事实来源,业务数据不能丢。但这个顺序有个致命风险——如果数据库事务提交成功,但 MQ 发送失败或者 Broker 响应超时,这条消息就再也没有机会发出去了。下游系统不会收到事件,缓存不会更新,搜索引擎也不会同步,你只能在业务代码里 catch 异常做重试,但进程崩溃、网络分区这些场景下,重试逻辑本身也不可靠。
先发 MQ,再写数据库,风险更明显。如果消息发出去了,但数据库写入失败,下游消费者已经拿到消息开始处理,它去查数据库会发现这笔业务数据根本不存在,轻则报错,重则产生脏数据。我见过一个支付回调场景,团队为了让回调更快返回,把发消息放在了写库之前,结果一次数据库连接池耗尽导致订单没落库,但支付成功的消息已经广播出去了,库存扣减、积分发放全部执行,最后人工对账才发现问题,折腾了一整天才把数据修正。
所以不该问"先写哪个",而应该问"写失败之后怎么办"。双写一致性问题的本质,不是顺序问题,而是补偿机制的设计问题。
2.2 本地消息表:最朴素但最可靠的兜底方案
本地消息表是我个人最推荐的双写一致性方案,也是这次专场里至少两位分享者给出的答案。它的核心思想很简单:把写业务表和写消息表放进同一个本地事务,消息发送状态独立跟踪,用定时任务补偿未发送成功的消息。
-- 业务事务里同时写入消息表 START TRANSACTION; INSERT INTO order (id, user_id, amount, status) VALUES (1001, 88, 99.00, 'CREATED'); INSERT INTO outbox (msg_id, biz_type, biz_id, payload, status) VALUES (uuid(), 'ORDER_CREATED', 1001, '{"orderId":1001}', 0); COMMIT;这样做的关键点是,业务数据和待发消息要么同时成功,要么同时失败,不再存在"数据库有但 MQ 没有"的中间状态。之后由一个后台任务定期扫描 outbox 表中 status=0 的记录,调用 MQ 发送接口,发送成功再把状态改成 1。
实际落地时要注意两个细节。第一,消息表需要设计独立的主键或唯一索引(比如用 biz_type + biz_id 做唯一键),保证定时任务重复扫描时不会重复发送。第二,定时任务的扫描频率要和生产消息量匹配,量大的话每秒钟扫一次,量小可以放宽到一分钟一次,但需要接受一定的发送延迟。
这个方案的缺点也很明显,消息表本身对数据库有一定压力,而且引入了额外的表结构和定时任务代码。但对于大多数业务系统来说,它是可控、可维护、可解释的方案,出问题排查起来路径清晰,比引入分布式事务中间件容易得多。
2.3 从时序角度重新看双写一致性问题
分享会上有个观点让我印象很深:双写问题不该从"谁先谁后"看,而该从"事件时序"的角度看。数据库里的业务状态是结果,MQ 里的事件是过程,两者本身就不是同一类东西,硬要追求"同时发生"本来就是错的。
更合理的思路是,以业务数据为准,把 MQ 消息当作数据变更的"通知"而不是"数据本身"。消费者收到消息后,不要直接信任消息体里的内容,而是用消息里的业务 ID 去查询最新的业务数据。这样即使消息发送延迟,甚至消息丢失后通过补偿机制补发,消费者都能拿到一致的数据视图。
这个思路也解释了为什么很多团队最终会选择"先写数据库 + 订阅变更日志(比如 Debezium 同步 binlog)+ 发 MQ"的架构。它本质上不是双写,而是把数据库日志作为唯一的事实来源,MQ 消息只是日志的投影。这套方案的一致性最强,但复杂度也最高,适合业务体量大、对一致性要求极高的核心系统。
对于中小团队,我的建议很直接:先用本地消息表把全部逻辑跑通,等真的遇到性能瓶颈了,再考虑引入更复杂的方案。不要一上来就上分布式事务,那是给具体场景的解法,不是默认配置。
3. Key_Shared 模式不消费的 Bug:消费模型与坑点全解
3.1 Pulsar 的四种订阅模式对照
Pulsar 之所以比传统 MQ 灵活,很大程度上是因为它支持四种订阅模式。很多团队刚接触时只用了默认的 Exlusive 模式,完全没发挥出它真正的能力。我整理了一张对比表,方便你快速理解各自适用场景:
| 订阅模式 | 消费方式 | 顺序保证 | 适用场景 |
|---|---|---|---|
| Exclusive | 单消费者独占订阅 | 严格有序 | 顺序要求极高的单消费者场景 |
| Failover | 多消费者,主备切换 | 严格有序(主消费者) | 高可用场景下仍需严格顺序 |
| Shared | 多消费者共享消息 | 无顺序保证 | 吞吐优先、顺序无要求的场景 |
| Key_Shared | 按 key 哈希路由到固定消费者 | 相同 key 内有序 | 需要共享吞吐,又要求相同 key 有序 |
从表格能看出来,Key_Shared 是 Shared 和 Exclusive 的折中方案。它既要共享模式的高吞吐,又要保留相同 key 消息的顺序性,所以在设计上必然会引入额外复杂度。这个复杂度就是坑的来源。
3.2 Key_Shared 不消费问题的根因与排查路径
这次专场里被讨论最多的一个生产事故,就是 Key_Shared 模式下消费者扩容之后,部分消息长时间不被消费。表面现象是消息积压持续上涨,消费者集群看起来正常,但就是有消息卡住不动。
根因要从 Key_Shared 的路由机制说起。Pulsar 在 Key_Shared 模式下,会根据消息 key 做哈希计算,把相同 key 的消息路由到同一个消费者。但消费者列表发生变化时(扩容、缩容、某个消费者宕机),Pulsar 会将 key 的映射关系重新分配。问题在于,这个重分配过程中,如果消息已经写入了旧的消费者对应的 backlog 里,而新的 mapping 又把后续消息路由给了新消费者,就可能出现旧 backlog 里的消息等待旧消费者继续消费,但旧消费者已经不在了,或者负载已经转移,导致这些消息永远卡在积压队列中。
我踩过一次很类似的坑。当时线上集群从三个消费者扩容到六个,扩容完成后发现一个分区的 backlog 一直不下降,消费者日志里也没有任何报错。查了很久才意识到,是 Key_Shared 模式下某些 key 的消息在扩容时被分配到了已经退出的消费者分支路径上,而新的消费者实例的订阅游标并没有接续那部分 backlog。最终我们是靠重置订阅位置,让消费者从积压位置开始重新消费才恢复。
排查这类问题的路径,我给个清单:
- 先看消费者组的订阅列表,确认所有消费者实例是否都处于 Active 状态
- 检查每个消费者的 backlog 分布,找出不消费的分区或 key 范围
- 查看 broker 日志中关于 Key_Shared 重分配的事件,确认扩容时是否发生过映射切换
- 用 pulsar-admin 命令查看 topic 的订阅状态,确认 cursor 是否推进到了最新位置
3.3 如何通过监控与治理避免同类问题
Key_Shared 的问题光靠排查不够,更要在架构层面做好治理。我总结了几个有效的措施。
第一,严格控制 Key_Shared 的 key 分布。如果某个 key 的消息量特别大,这个 key 对应的消费者必然成为热点,其他消费者再空闲也帮不上忙。针对这种情况,最好是拆 key、加一层二次 sharding,或者在业务层面做 key 粒度拆分。
第二,合理设置消费者的限制参数。Pulsar 的 Key_Shared 模式下,消费者堆积会导致消息重排,建议配置 maxUnackedMessagesPerConsumer 限制单个消费者未确认消息数量,避免某个消费者积压过深。同时开启 dead letter topic,让重试次数超限的消息进入死信队列,而不是永远卡住正常消费流程。
第三,做好定时巡检。我在实践中会写一个脚本,定时扫描所有生产重要的 subscription 指标,如果发现某个订阅的 backlog 超过阈值且增长率持续为正值,就触发告警,推送到值班群。这类问题很多都是慢性的,不会立刻爆雷,但放任不管迟早会变成故障。提前发现,在业务低峰期做一次消费位移重置,往往就能解决。
4. 消息补发、幂等消费与最终一致性
4.1 事务消息与延迟队列:补发的标准姿势
消息补发不是"重发一遍"那么简单,它要回答两个问题:补发什么时间段的消息,怎么保证补发的消息不会和已消费的消息冲突。
第一种场景是"先写数据库,但 MQ 发送失败",这种补发用本地消息表就能解决,定时任务扫描状态为待发送的记录即可。第二种场景是"消息已经发送成功,但消费端处理失败",这种补发要依赖 MQ 自带的重试机制。Pulsar 支持消息重试队列和延迟队列,消费失败的消息会进入 retry topic,按延迟级别重新投递,超过最大重试次数进入死信队列。
事务消息则是另一个思路。它的核心是两阶段:先发半消息,等本地事务提交后再确认发送;如果本地事务回滚,半消息会被删除。这个机制解决的是"发送方无法确定数据库事务是否成功"的问题。我在生产环境用过 RocketMQ 的事务消息做订单状态变更通知,体验是:机制本身可靠,但需要额外处理半消息的超时回查逻辑,复杂度确实不低。
延迟队列本身也很实用。比如支付超时未回调、订单超过 15 分钟未付款自动关闭,这类场景用延迟消息比用定时任务扫表优雅得多。Pulsar 的 delayed message 是基于时间戳分桶实现的,延迟消息不会阻塞普通消息的消费,适合大量延迟任务的场景。
4.2 消费幂等:消息重复才是常态
不管你怎么设计补发逻辑,都必须接受一个现实:消费者收到的消息可能会重复。原因是多方面的,发送方重试导致重复投递,消费方处理成功后还没来得及提交 offset 就宕机,重启后重新消费,这些都会造成重复。
解决重复的唯一方式是消费端幂等。我的经验是三层幂等设计:
第一层,利用业务唯一 ID。比如订单创建事件里带上 orderId,消费端处理前先查一下当前订单状态,如果已经是终态就跳过。这个逻辑简单,但大多数场景下够用。
第二层,利用数据库唯一索引或 redis 分布式锁。比如库存变更记录设计唯一键 (order_id, sku_id),数据库层面保证同一条变更只落一次。用了唯一索引后,即使消息重复投递,第二次插入也会报错,代码里 catch 住这个异常即可。
第三层,状态机兜底。核心状态变更尽量设计成显式状态机,比如订单从 CREATED 到 PAID 到 SHIPPED 单向流转,消费端强制按状态机规则推进,不合法状态转换直接拒绝。这样即使消息乱序或重复,也不会把状态改错。
我想强调的是,幂等不是某个方法的事,是整条消费链路的约束。任何消费端代码都要假设"这条消息我之前可能已经处理过了"。
4.3 消息轨迹和全链路追踪怎么配合
消息补发和幂等设计得再好,如果没有观测能力,出了问题你依然是盲人摸象。我在消息中台建设中体会最深的一句话是:可靠的 MQ 系统有三根支柱——消息不丢、消费不重、问题能看到。第三根支柱往往最容易被忽视。
建议在生产环境接入消息轨迹功能。Pulsar 支持通过 Broker 拦截器采集消息的 生产时间、存储位置、消费时间、消费结果,并且可以输出到外部存储做查询分析。云上 MQ 产品一般也自带消息轨迹控制台。
同时要和全链路追踪结合。每一条业务消息在生产者侧生成时,都会携带 traceId 和 spanId,消费者侧解析这些 ID 并把处理链路接续上去。这样一旦消息处理失败,你能直接从 trace 系统跳到 MQ 的消息轨迹,看到消息从哪台机器产生、经过哪些节点、在哪个消费环节出了差错。我在实际故障排查里,用这套组合把一次需要两小时的问题定位缩短到了十分钟以内。
5. 算术编码原理:MQ 背后被忽视的压缩技术
5.1 算术编码的核心思想:用一个小区间表示整个序列
讨论完消息可靠性和消费模型,回到一个偏原理但很值得懂的话题:算术编码。很多人看到"算术编码"会以为是数学课的内容,但实际上它在 MQ 的消息压缩领域有非常实际的应用。
先解释它的核心思想。如果我们要编码一段符号序列,比如 "ABCAB",传统 Huffman 编码是为每个符号分配一个整数二进制码,然后把码拼接起来。算术编码换了一个角度——它把整段符号序列映射到 0 到 1 之间的一个实数区间上,区间越小,需要的二进制小数位数越多,也就越精确地表达了这个序列。
它的工作过程是这样的:预先知道每个符号的出现概率(比如 A 占 0.5,B 占 0.3,C 占 0.2),然后从一个初始范围 [0, 1) 开始,每读入一个符号,就按这个符号的概率比例缩小当前范围。处理完所有符号后,得到一个最终的小数区间,从区间中任意选一个二进制小数作为编码结果就可以。
我用一个生活化的例子说明。想象你有一把一米长的尺子,尺子上按概率划分了 A、B、C 三个区段。第一个符号如果是 B,你就把目光锁定在 B 对应的那一段刻痕上,再把这段刻痕继续按概率划分成三份。读下一个符号,继续往里缩。处理完整段消息后,你在最终锁定的那一小段上取一个刻度数字,这个数字就是整段消息的编码。解码的时候,你看到这个数字,反向在尺子上不断判断它落在哪个区段,就能逐个还原出原始符号。
5.2 算术编码在消息压缩和网络协议中的应用
算术编码最核心的优势是压缩率接近信息熵的理论极限。Huffman 编码每个符号至少要用 1 位二进制表示,遇到概率极低的符号,短的码字根本不够用。算术编码没有这个限制,它本质上是用"越来越准确的小数"来编码整个序列,所以对符号频率分布的适应性更好。
在 MQ 领域,大批量消息在传输前通常会做压缩。常见的 LZ4、Snappy、Zstd 主要利用字典匹配和熵编码(Zstd 内部会用到有限状态熵编码,原理上也可以归到这类),它们在平衡压缩速度和压缩率方面做了很多工程优化。算术编码本身计算量偏大,很少直接用于逐条消息的实时压缩,但理解它的原理,能帮你更好地理解为什么有些压缩算法在特定数据分布下表现更好。
我实际见过的一个应用是定制化的消息压缩策略:消息内容如果具有明显的字符概率不均匀分布(比如超长的 JSON 里 status 字段只有几个固定值,大部分内容高度重复),团队会先对消息做字段拆分,对高重复字段走字典编码,对低频率长尾字段结合算术编码思路做熵编码,整体压缩率比直接上 Snappy 提高了 30% 左右。
不过对于大多数业务场景,我建议还是优先用好现成的压缩算法。Pulsar 生产端配置 compression type 就能选择 LZ4 或 Zstd,改动成本低且稳定。算术编码这类技术,更多是当你需要对压缩做专项优化时,理解其背后的原理才能找到改进方向。
5.3 消息压缩的工程取舍:什么时候值得手动介入
这里补充一个容易踩坑的点:压缩不是无脑开就一定好。消息体很小(比如小于 1KB)时,压缩算法的开销和增加的数据包复杂度可能超过节省的带宽,反而得不偿失。我在某次压测里发现,开启 LZ4 后,单条 500 字节左右的消息虽然体积减小了 20%,但 CPU 消耗上升了约 8%,在吞吐量极高的链路上,这个 CPU 开销是不可忽略的。
我的建议是分场景测试:大消息(大于 4KB)优先开启压缩;小消息建议保持不压缩,除非带宽成为瓶颈。另外,还要监控压缩率指标,如果线上消息的压缩率始终接近 1,说明数据本身重复度不高,压缩策略需要重新优化。
关于算术编码,最后再提一句。它和 MQ 的关联除了消息压缩,还有一层更底层的逻辑:任何需要把大量数据变成尽量小的比特流的场景,本质上都是在做熵编码。理解了算术编码,你对"数据为什么能压缩"以及"压缩的极限在哪"这两个问题就建立了很本质的认知,这对做消息系统的底层优化会有长期帮助。
6. 我的实操心得和避坑建议
6.1 消息治理三板斧:规约、巡检、混沌
基于这些年运维消息系统的经验,我把最核心的心得总结成三板斧:规约、巡检、混沌。
规约是指消息命名的规范。每条消息都要带上明确的业务类型、版本号、幂等键和链路 ID。别小看这些字段,消息系统出问题时,排查效率很大程度上取决于你消息字段是否齐全。我们内部做过一个统计,规范前故障平均定位时间约 1.5 小时,规范后缩短到 25 分钟,差距非常明显。
巡检是指有一套自动化的指标巡检机制。消息积压量、消费延迟、重试次数、死信数量这四项指标,建议至少按分钟粒度采样,超过阈值自动告警。另外每季度做一次消费链路梳理,清理无效订阅和过期 topic,避免消息中台里积累大量僵尸 topic。
混沌是指定期做故障演练。比如人为停掉一个 BookKeeper 节点,看看消息生产是否受影响;或者让某个消费者实例宕机,观察 Key_Shared 模式的重分配是否符合预期。演练的目的不是证明系统能抗住故障,而是让团队在真实故障来临时不慌。我见过太多平时状态良好、一故障就手忙脚乱的团队,问题就出在没演练过。
6.2 给新手的三个建议
如果看到这里,你刚准备把核心链路迁到 Pulsar,或者正准备优化自己团队的 MQ 架构,我有三个建议。
第一,不要迷信单一组件。没有最好的 MQ,只有当前业务阶段最合适的选择。团队维护成本、协议生态、客户端成熟度,这些都要纳入考量。我见过不少团队把 Kafka 换到 Pulsar 后并没有获得预期的收益,原因是他们根本没用上多租户和存算分离的优势,只是平迁,那迁移的意义就很小了。
第二,迁移要有回退方案。任何一次 MQ 变更都要考虑如何回退,不仅仅是"切回旧集群",还包括消费进度对齐、消息轨迹衔接、双跑期间的重复消费策略。把回退方案设计得足够细致,你才敢在业务高峰期放心操作。
第三,先保证可靠性,再追求性能。消息系统的第一优先级永远是不丢不重、可追踪。在可靠性没验证充分之前,不要为了撑高吞吐去调整批量参数、关闭持久化,这些优化动作引发的问题往往比收益更难看。
6.3 再分享一个小技巧:把 MQ 当数据库来治理
做消息中间件这些年,我积累的最深体会是:要把 MQ 当数据库来治理。所谓"当数据库治理",是指 MQ 里的每一条消息都是事实、是记录,应该有生命周期管理、有访问审计、有数据质量管理。很多人只把 MQ 当成一个"传数据的管道",用完即弃,这是最大的误区。管道是会堵的,数据是会产生价值的。我在实际工作中,会把消息的 schema 纳入统一管理,给消息 producer 和 consumer 都做应用级鉴权,重要 topic 的消�息保留时长设置得足够长,以备审计和追溯。这套治理理念,比任何具体工具都能帮你减少故障。
回到这次大会的主题"Make MQ Great Again",应该说 MQ 从来不需要被"拯救",它只是需要被理解和被认真对待。每个 Key_Shared 的坑、每份双写的妥协、每次幂等设计的思考,都让消息系统的实践经验更扎实。希望这篇回顾能让你在下次遇到 MQ 问题时,少走一些我走过的弯路。