新闻详情

新闻详情

首页 / 资讯中心 / 详情

MQ架构实战:从双写一致性到Pulsar Key_Shared与消息压缩

发布时间:2026/9/26 5:51:47来源:尧图网络
MQ架构实战:从双写一致性到Pulsar Key_Shared与消息压缩
上半年最忙的一段时间刚过去趁着记忆还新鲜把 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 把存储层抽出来放到 BookKeeperBroker 只负责计算和路由扩容就是加 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 表中 status0 的记录调用 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.5B 占 0.3C 占 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 问题时少走一些我走过的弯路。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

零基础一个月过关软考高项:跟对老师比死磕教材更有效 2026/9/26 6:29:46

零基础一个月过关软考高项:跟对老师比死磕教材更有效

备考时间告急时,我才真正理解“高项”这两个字的重量。作为完全零基础的考生,一开始连信息系统项目管理师考什么、分几科、怎么报名都搞不清楚,翻开教程那一刻直接懵掉——满眼都是“整合管理”“范围确认”“风险定性分析”这些术语&#xf…

阅读更多 →
WorkBuddy接入七牛云大模型广场性能优化实战指南 2026/9/26 6:29:46

WorkBuddy接入七牛云大模型广场性能优化实战指南

1. WorkBuddy 任务执行慢不是“卡”,而是模型调用链路上的多层隐性耗时叠加WorkBuddy 任务执行慢,这个现象在最近两周的开发者社区里高频出现——不是报错、不是崩溃,就是“点下去等三秒才出结果”,用户反复刷新、重试&#xff0c…

阅读更多 →
Kubernetes生产环境部署:从裸机到高可用集群的完整实践 2026/9/26 6:29:45

Kubernetes生产环境部署:从裸机到高可用集群的完整实践

1. 为什么“从零到生产可用”不是一句空话,而是K8s落地最真实的分水岭很多人点开“K8s部署教程”时,心里想的是:装完kubectl、kubeadm、拉起一个master节点、跑通一个nginx Pod,就算“学会了”。我带过三轮K8s内训,每次…

阅读更多 →
高项零基础31天备考攻略:跟对老师,三科一次过 2026/9/26 6:29:44

高项零基础31天备考攻略:跟对老师,三科一次过

朋友发来那条消息的时候,距离考试只剩31天。她是零基础,报名后才翻开官方教材,翻了两天心态就崩了——厚厚一本教程,每一页都像天书,项目管理术语完全看不懂,计算题更是一头雾水。她在消息里连发三个问号&a…

阅读更多 →
Redis二级缓存设计实战:彻底解决热key与缓存穿透 2026/9/26 6:29:44

Redis二级缓存设计实战:彻底解决热key与缓存穿透

上个月我们线上一个查询商品的接口挂了,Redis CPU 飙到 95%,连接数打到上限,数据库的慢查询塞满监控页。排查下来原因很简单:首页和详情页同时刷一批热点商品,每次都是先查 Redis 再查数据库,而重复的 key …

阅读更多 →
Ubuntu 上安装 Codex CLI 与 Claude Code 的完整避坑指南 2026/9/26 6:29:37

Ubuntu 上安装 Codex CLI 与 Claude Code 的完整避坑指南

1. 为什么要在 Linux 上折腾这两个命令行工具如果你最近在关注 AI 辅助编程这个方向,大概率已经反复看到两个名字:Codex CLI 和 Claude Code。前者是 OpenAI 推出的终端编程助手,后者是 Anthropic 出的同类产品,两者都主打"在…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉