ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

消息队列两大核心模式:点对点与发布/订阅全解析

消息队列两大核心模式:点对点与发布/订阅全解析 做了这么多年后端消息队列几乎是无处不在的基础设施。以前跟同事讨论方案最常被问到的就是这里用队列还是用 topic说白了就是点对点模式和发布/订阅模式怎么选。很多人对这两种模式的理解停留在一个是一条消息只能被消费一次另一个是能被消费多次但真正落到项目里涉及 ACK 机制、消费组、消息堆积、重复消费这些问题时光靠这句话远远不够。这篇内容我打算从实际项目视角出发把消息队列的两大核心模式拆开讲透——先说清楚消息队列到底解决了什么问题再分别讲点对点和发布/订阅的底层模型、适用场景最后给出一份可以直接参考的选型建议和避坑清单。不管你是刚接触消息队列的初级开发还是正在设计中间件方案的资深工程师这篇都能帮你把基础概念和工程落地连成一条线。1. 三个根本的业务痛点消息队列为什么值得引入聊两种模式之前必须先想清楚一个问题我们为什么要引入消息队列如果连这个都没想明白后面选模式就是空中楼阁。业界常说的三大作用——异步、削峰、解耦——不是空话它们对应着真实业务场景里的具体痛点。1.1 异步把必须立刻完成和可以稍后完成分开同步调用最直观的问题是响应时间被链条上最慢的环节拖死。举个例子用户下单后系统要扣库存、发优惠券、送积分、发短信通知如果全是同步调用哪怕每个服务只要 50ms串起来也要 200~300ms。而用户真正关心的只是下单成功短信晚 3 秒发没有任何影响。引入消息队列后主链路只做核心操作写入订单、扣减库存然后把发优惠券发短信这类非核心动作变成消息投递到队列里由下游服务异步消费。用户看到的响应时间从 300ms 降到 80ms体验提升是实打实的。这里用到的通常是点对点模式——每个任务只需要一个处理方来完成谁抢到谁处理。1.2 削峰用缓冲换稳定秒杀是削峰最典型的场景。平时订单系统每秒处理 1000 个请求足够但秒杀瞬间流量可能是平时的 50 倍。如果让下游直接扛 5 万 QPS数据库大概率直接打满。消息队列在这里充当了一个缓冲池上游把请求快速写入队列下游按照自己的最大处理能力慢慢消费流量高峰被削平下游系统不会因为瞬时压力过载。这里的关键点是队列的吞吐量要远大于下游的消费能力否则削峰变成堵车。实际项目中我会关注消费端积压情况一旦积压量超过阈值就告警。这个场景下点对点模式仍然是主流因为每个秒杀请求只需要一个消费者处理。1.3 解耦让上下游互不相识没有消息队列时A 系统要把数据给 B、C、D 三个系统就得在 A 里写三个调用逻辑。明天新增一个 E 系统要接数据A 又得改代码、发版。更麻烦的是如果 B 系统宕机A 的调用还会阻塞甚至失败。引入消息队列后A 只需要把数据发布到 TopicB、C、D、E 各自订阅即可。A 不认识也不关心谁订阅了新增或下线一个消费者都不需要动 A 的代码。这里用的就是发布/订阅模式——一份消息广播给所有关心它的系统。解耦带来的维护收益在系统数量超过三个后非常明显。2. 点对点模式拆解一条消息一生只会被一个消费者领走点对点模式Point-to-Point是所有消息队列最基础的形态。它对应的模型是 Queue队列核心语义是生产者把消息放入队列消息只会被一个消费者取走并处理处理完毕即从队列移除。用一句大白话概括一份工作几个人竞争谁抢到谁来干。2.1 核心组成与消息流转逻辑点对点模式里有三个角色Producer生产者、Queue队列、Consumer消费者。流转逻辑如下Producer 把消息发送到指定 Queue消息持久化到队列存储Consumer 主动拉取Pull或等待推送Push消息队列按投递策略选择其中一个 Consumer 投递消息Consumer 处理完毕后发送 ACK 确认队列收到 ACK 后删除这条消息。这个模型看起来简单但有一个细节容易误解多个消费者监听同一个队列时每条消息只会被其中一个消费者消费而不是所有消费者都消费。这是点对点模式和发布/订阅模式最本质的区别。实际项目中队列选择消费者的策略通常有几种轮询Round Robin、公平分发Fair Dispatch、按权重路由。RabbitMQ 的 Work Queue 默认用的是轮询但如果没有配合 ACK 和 prefetch 限制容易出现某个消费者处理慢、其他的在空转的情况。所以生产环境我一般会设置prefetch_count让消费者每次只取一条、处理完再取下一条这样处理慢的消费者不会被持续塞入新消息。2.2 ACK 确认机制为什么消费成功必须举手回报ACKAcknowledgment是点对点模式的核心机制。消费者从队列拿到消息后如果什么都不说队列怎么知道消息处理成功还是失败所以约定消费者处理完业务后必须向队列发送 ACK队列收到 ACK 才会把消息标记为已消费并删除。这里有一个常见的工程陷阱如果消费者在业务处理完成后、发送 ACK 之前崩溃了队列会认为这条消息没有被成功消费于是重新投递给其他消费者或稍后重投。结果就是业务已经被处理了一次又会被处理第二次——这就是重复消费问题的根源之一。所以在使用点对点模式时消费者的处理逻辑必须设计成幂等的。所谓幂等就是同一个操作执行一次和执行多次的结果一致。比如扣减余额操作不能简单地balance - amount要先判断这个订单是否已经扣过款。我在项目中常用方案是在业务表里加一个唯一的message_id字段消费时先查这个 ID 是否已经处理过处理过就直接返回成功不重复执行。2.3 多消费者场景下的负载均衡与顺序问题点对点模式天然支持消费端的水平扩展。生产者把大量消息写入队列队列按策略把消息分散到多个消费者实例上从而提升整体消费吞吐量。这个特性让点对点模式成为任务分发类场景的默认选择订单处理、邮件发送、图片处理、日志写入都可以用多个消费者并行处理。但有得必有失多消费者并行带来一个副产品消息顺序可能乱。假设生产者依次发送了 A1、A2、A3 三条消息它们被队列轮询投递给了三个不同的消费者实例由于每个实例的处理速度不同可能出现 A2 先处理完、A1 后处理完的情况。如果业务对顺序有严格要求比如必须先创建用户再发送欢迎邮件就不能简单依赖多消费者并行。业界常见方案有几种一是把需要保证顺序的消息路由到同一个消费者比如按用户 ID hash 后投递到固定分区二是在业务层做顺序校验乱序的后续操作直接拒绝并重试。Kafka 用分区Partition来保证同 key 消息的顺序就是后一种思路的典型实现。3. 发布/订阅模式拆解一条消息广播给所有订阅者发布/订阅模式Pub/Sub解决的是另一种问题一份消息需要同时给多个消费者。对应的模型是 Topic主题核心语义是生产者把消息发布到 Topic所有订阅了该 Topic 的消费者都能收到这条消息的副本。用生活中类比就像电台广播只要收音机调到同一个频率谁都能收到。3.1 Topic 模型与订阅关系发布/订阅模式里的核心概念是 Topic。生产者往 Topic 发消息消息不会像点对点模式那样一条只给一个人而是被复制成多份分别投递给每个订阅者。订阅关系是动态的一个消费者随时可以订阅一个 Topic也随时可以取消订阅生产者和消费者互相无感知。这里要澄清一个高频误解在发布/订阅模式下如果一个 Topic 没有任何订阅者生产者发出去的消息会怎样答案是直接丢弃对非持久化订阅而言。很多初学者在这里踩坑——测试时先启动生产者发消息再启动消费者订阅结果发现消费者什么都收不到就是因为订阅关系成立之前的消息已经被丢弃了。理解这一点后就能明白为什么实时性要求高的场景更适合发布/订阅而离线消费场景必须依赖持久化订阅。3.2 持久化订阅与非持久化订阅两条完全不同的规则发布/订阅模式内部还分为两种订阅方式在实际工程中非常关键。非持久化订阅Non-Durable Subscription的逻辑很简单消费者在线时能收到消息一旦消费端断开消息就丢了。它适合对消息可靠性要求不高的实时通知场景比如前端实时收到新订单提醒断线重连后丢失几条提醒可能无所谓。持久化订阅Durable Subscription则不同。消费者在订阅时注册了自己的持久化订阅 IDBroker 会为这个订阅保存消息进度即使消费者离线消息也不会被丢弃等它重新上线后会从上次消费的位置继续拉取。JMS 里的 Durable Topic Subscription、Kafka 的 Consumer Group 配合 log retention 都属于这类。我实际工作中遇到过一个典型的持久化订阅应用订单系统把订单变更事件发布到 Topic下游的数据同步服务、搜索索引服务、数据分析服务分别订阅。如果数据同步服务临时停机半小时重启后依然能消费到停机期间产生的全部订单事件不会丢数据。这就是持久化订阅的价值。3.3 消费者组组内竞争消费组间同时广播Kafka 在发布/订阅基础上做了一层非常重要的抽象——Consumer Group消费组这在实际项目里几乎绕不开。一个消费组里可以有多个消费者实例同一个 Topic 的一条消息只会被组内的一个消费者实例消费但不同消费组之间每个组都会完整消费到这条消息。用一句话概括就是组内是点对点组间是广播。这个抽象非常巧妙它让同一个 Topic 既能支持多个消费者并行分摊处理组内竞争又能支持多个业务系统各自独立消费组间广播。举个例子订单 Topic 可以被两个消费组订阅订单处理组有 5 个消费者实例每条订单消息只会被其中 1 个实例处理5 个实例分摊全部消息提升吞吐量数据分析组有 2 个消费者实例每条订单消息也只会被分析组内的 1 个实例处理但这个组消费到的消息是和订单处理组完全独立的副本。这种设计带来的直接好处是如果想增加一种新的消费用途只需要新加一个消费组去订阅同一个 Topic完全不用改动生产者和已有消费者。这个特性让 Kafka 成为大数据场景下数据分发的首选本质就是发布/订阅模式加上了分区与消费组机制。4. 两种模式的核心差异对照前面分别讲了两者的模型和流转逻辑这一节我用表格把核心差异集中列出来方便读者在方案设计时快速对照。很多技术方案争论到最后回归到这些基础差异上答案往往就清晰了。4.1 一张表理清两种模式的关键区别对比维度点对点模式发布/订阅模式核心模型Queue队列Topic主题消息消费次数一条消息只被一个消费者消费一次一条消息可被所有订阅者各消费一次消费关系一对一的竞争关系一对多的广播关系消息生命周期消费成功并 ACK 后删除每个订阅者都有独立消费进度需各自消费完负载均衡多消费者分摊队列消息消费组内部分摊组间各拿全量副本离线消息消息保存在队列中消费者上线后仍可消费非持久化订阅会丢持久化订阅可恢复典型场景任务分发、订单处理、日志异步写入事件通知、系统解耦、数据分发、实时广播代表实现RabbitMQ Work Queue、ActiveMQ Queue、RocketMQ 普通消息Kafka、RabbitMQ Pub/Sub、MQTT、JMS Topic这个表格里最值得停下来想一行的是消息消费次数。点对点模式下队列天然保证一条消息只被一个消费者领走所以消息不需要复制发布/订阅模式下消息必须复制多份每个订阅者一份。这条差异直接决定了后续所有设计——包括消息存储结构、消费进度管理、可靠投递策略。4.2 为什么容易混淆同一个产品同时支持两种模式很多人在学习时混淆两种模式一个重要原因是主流消息队列产品往往同时支持两种模式API 表面上差别也不大。以 RabbitMQ 为例它本质上是 AMQP 协议的实现通过 Exchange交换机和 Binding绑定规则可以做出队列模型也可以做出发布/订阅模型。具体来说RabbitMQ 用默认 Exchange 直连 Queue就是点对点模式用 Fanout Exchange 绑定多个 Queue每个 Queue 对应一个消费者就是发布/订阅模式。外面的壳都是发送消息、接收消息但内部的消息路由和消费语义完全不同。所以在看代码或者设计文档的时候不要只看用了什么消息中间件要去看它的消息模型。同一个 RabbitMQ 集群里可能一部分队列跑的是点对点另一部分走的就是发布/订阅。判断标准始终是这条消息被写完后是一个消费者处理还是多个消费者各处理一遍。5. 重复消费问题的根因与幂等方案消息队列相关热词里重复消费是被问到最多的问题之一。无论用点对点还是发布/订阅模式重复消费都可能发生。很多人以为消息队列应该保证不重复消费这个认知需要纠正。绝大多数消息队列提供的是 at least once至少一次语义而不是 exactly once恰好一次。也就是说消息可能重复但通常不会丢失。5.1 重复消费到底是怎么发生的重复消费的根源通常不是消息队列自己多发了消息而是分布式环境下的不确定性和重试机制。典型路径如下消费者拉取到消息开始执行业务逻辑业务逻辑执行成功比如已写入数据库消费者在发送 ACK 前突然宕机或网络发生抖动导致 ACK 包丢失消息队列没有收到 ACK判定消费失败重新投递这条消息消费者重启后再次收到同一消息重复执行业务逻辑。这个路径在点对点模式中最常见在发布/订阅模式的持久化订阅中同样会出现。只要采用了先处理业务、后确认消费这种常规顺序就必然存在上述重复窗口。还有一种情况消费者处理超时消息队列主动重投。比如消费端处理一条消息耗时 5 分钟超过了 Broker 配置的 ack 超时时间Broker 强制重新投递消费者端就会同时收到两条内容相同的消息。这种问题不是业务逻辑 bug而是配置和实际处理耗时严重不匹配造成的。5.2 幂等设计不能把希望寄托在恰好一次既然消息队列没法根除重复那解决方案只有一个让消费端具备幂等性重复消费不产生副作用。这是我做消息消费设计时优先级最高的一条原则和选哪种模式无关。几个我在项目中用过的幂等方案唯一业务键 去重表每条消息携带一个全局唯一的消息 ID或由业务键如订单号派生消费时先查去重表已存在则直接 ACK否则执行业务并写入去重表。注意去重操作和业务操作要在一个事务里否则并发下可能双双穿透。数据库乐观锁业务表里加版本号字段更新时带上版本号条件UPDATE t SET amount amount - 100, version version 1 WHERE id ? AND version ?更新行数为 0 说明已被执行过直接忽略。状态机校验订单状态从待支付改成已支付只有一次合法流转重复消费时如果发现当前状态不是待支付直接拒绝处理。这里我想给一个非常具体的建议不要过度信任消息中间件自带的幂等功能它通常有很强的限定条件。比如有些队列支持 deduplication 插件但只能保证在某个时间窗口内去重窗口外的重复消息依然防不住。可靠的方案一定是在业务侧做幂等控制消息队列只负责投递和确认。6. 选型建议与实际落地经验究竟是选点对点还是发布/订阅最终要回到业务场景上。没有绝对的好坏只有合适不合适。结合我的项目经验给出几条可落地的选型思路和踩坑记录。6.1 什么场景选点对点什么场景选发布/订阅判断标准可以简化为三个问题这条消息需要几个独立的业务系统处理如果只需要一个选点对点如果需要多个各自处理选发布/订阅。这些处理之间是什么关系如果多个处理是上下游依赖关系先 A 后 B用点对点加消息链路如果是并列关系A、B 同时执行且互不依赖用发布/订阅让它们各拿各的副本。如果增加一个新的消费需求能不能接受改生产者代码如果能接受点对点也能凑合如果不能接受发布/订阅是唯一合理解。具体到业务场景订单创建后需要异步扣库存、发短信、更新统计报表库存扣减和短信发送互不依赖但统计报表系统可能是新接入的——这种我建议用发布/订阅三个订阅者各取所需后续加一个积分服务也不用改订单主流程。用户上传图片后需要异步生成缩略图这个任务只需要一个消费者处理重试也只需要针对这一个任务——点对点模式更直接。日志采集所有应用把日志发到一个 Topic多个消费组分别做实时告警、离线分析、冷存储归档——这是发布/订阅的经典用法。6.2 真实项目中我踩过的几个坑第一个坑发布/订阅模式下的消费者临时出问题消息丢失。以前我们用 RabbitMQ 的 Fanout Exchange 做事件广播消费者是纯内存订阅服务一重启就丢失部分消息。后来改成持久化订阅每个消费者注册固定订阅名Broker 端保存消费进度问题才解决。这个教训让我养成了一个习惯凡是涉及事件通知一律优先考虑持久化订阅不给自己留丢数据的隐患。第二个坑点对点模式下消费者处理太慢导致消息大量积压。当时我们给图片处理服务配了 10 个消费者实例以为吞吐量跑满就没问题。结果发现生产端的批量任务一次性投递了百万条消息消费者每个实例每秒钟只能处理 20 条积压量一路飙升。后来我们加了一层批量拉取 批量处理的逻辑每条消息的处理成本降下来积压才逐渐消化。这个案例提醒我消费者扩容不是万能的处理逻辑本身的性能才是瓶颈。第三个坑同一个队列既当点对点用又想广播导致业务逻辑混乱。有同事为了省事让多个消费者监听同一个队列以为每个消费者都会收到全部消息结果消息被轮询分发每个消费者各拿一部分触发了一个隐蔽的 bug。排查了两个小时才反应过来——这不是消息队列的问题是根本用错了模式。设计阶段多花 10 分钟澄清消息语义远比出问题后排查来得划算。消息队列的两种模式在概念上并不难难的是在真实业务场景里辨别清楚自己的需求属于哪一种。我的个人经验是方案设计时先画清楚消息流转图标出生产者和消费者的关系一比照上面那张对比表用哪种模式基本就呼之欲出了。消息队列的坑大多不在中间件本身而在我们对模式的理解和业务边界的划分。希望这篇内容能帮你少走一些弯路。
返回列表