消息中间件RabbitMQ学习总结

发布时间:2026/8/21 10:32:46
消息中间件RabbitMQ学习总结 一什么是消息队列消息队列这个词由两部分组成第一消息message就是应用间传输的数据可是是简单的文本字符串也可以有复杂的嵌入对象。第二队列queue数据结构中有一章专门介绍队列最鲜明的特点就是先进先出基于此特点也是更加适用于消息的处理先来到消息先发放后续的消息等前面的消息发放完才执行符合真实的逻辑与预期。二消息队列有哪些作用和应用场景2.1 应用解耦理解应用解耦就是先清楚耦合是什么当多个系统直接调用强依赖关系时一旦其中有任何一个服务宕机就直接影响整个业务流程。举例就是经典的下单业务在没有MQ的情况下用户下单之后需要创建订单后台收到订单后需要扣减库存核实完之后需要发送拼单成功的短信 最后还可以添加优惠卷发放的功能当这四个功能高度耦合时一步接一步一个功能调用下一个功能当其中又任何一步出现问题时就会阻塞整个业务后续再新加任务时还需要修改前面已有代码这显然有违程序员的精简精神。于是MQ便登上舞台为大型的复杂业务解耦主业务也就是生产者只负责把消息发送给MQ发完结束自己的主流程各个下游服务也就是消费者各自监听MQ队列自己拉取消息异步的执行各自相对应的业务。MQ作为中间缓冲层上下游系统不再直接调用对方而是只和MQ进行交互实现解耦2.2 异步提速还是回到刚才的场景在高耦合的状况下用户下单成功需要等待后续一些操作全部完成后前端才会返回通知下单成功添加的业务越多链路越长接口响应的就越慢用户等待的时间就越久这显然是不合理的也与我们具体的体验不同。那我们实际下单流程是什么呢一般付完款之后会马上显示付款成功等待商家打包之类的后续过一段时间可能才会有短信提示下单成功商品号是什么走什么快递。这里其实就是引入MQ的具体体现生产者也只完成核心业务创建订单即可然后直接返回下单成功其余的一系列附属操作直接封装成一条条消息发送到MQ后续的库存短信通知等非核心操作交给消费者异步慢慢执行这就可以大大缩短接口响应时间提升系统的吞吐量。2.3 流量削峰这一部分大家应该很熟悉比如在某个固定的时间点演唱会抢票铁路12306抢票景区门票预约大部分时间都是一直在等待页面半天没反应退出刷新重进之后票没了。。。。。。如果没有MQ的情况瞬间的海量请求会直接打到后端服务数据库很可能就直接数据库断联整合系统宕机卡顿正常的其他功能也跟着一块死机MQ就相当于缓冲蓄水池高并发的瞬间所有的请求先转换成消息存入消息队列后端消费者按照自己能承受的工作量匀速从拉取消息处理其实这个速度很快但奈何请求量太大大部分请求仍然需要排队对于系统的好处就是能抵御瞬时高并发保护核心业务防止系统雪崩。三RabbitMQRabbitMQ 是一款开源、轻量、可靠的消息中间件消息队列采用 Erlang 语言开发遵循 AMQP高级消息队列协议。作用实现进程之间、不同服务之间的异步消息传递也就是我们前面讲的应用解耦、异步提速、流量削峰。组成生产者Producer消费者Consumer交换机Exchange队列Queue绑定Binding四RabbitMQ 的安装4.1 文件下载Windows 下载地址Erlang 官网https://www.erlang.org/downloadsRabbitMQ 官网https://www.rabbitmq.com/download.html4.2 安装文件Windows 安装步骤运行 Erlang 安装包一路下一步安装完成后配置系统环境变量ERLANG_HOME运行 RabbitMQ 安装包默认路径下一步即可以管理员的身份打开终端# 进入rabbitmq插件目录 cd D:\app\RabbitMQ Server\rabbitmq_server-xxx\sbin启动 / 停止命令# 前台启动 .\rabbitmq-server # 安装为系统服务(只需一次) rabbitmq-server --install-service # 启动服务 net start rabbitmq # 停止服务 net stop rabbitmq #开启 Web 管理面板 .\rabbitmq-plugins enable rabbitmq_management访问管理页面http://localhost:15672默认账号guest密码guest⚠️ guest 默认只能本机访问外网 / 虚拟机访问需要新建管理员用户4.3 添加一个新的用户# 1. 创建用户 rabbitmqctl add_user admin Admin123 # 2. 设置为管理员角色 rabbitmqctl set_user_tags admin administrator # 3. 赋予 / 虚拟主机全部权限 rabbitmqctl set_permissions -p / admin .* .* .*参数说明administrator管理员最高权限vhost/默认虚拟主机 权限三段配置权限、写权限、读权限4.4端口说明5672程序连接 RabbitMQ 通信端口代码连接用15672Web 管理后台页面端口浏览器访问五RabbitMQ 的工作原理一条消息的完整“生死之旅”工作流程假设你的订单服务生产者要给短信服务消费者发一条消息。阶段一生产者发消息发件建立物理连接TCP订单服务与 RabbitMQ Broker 建立 TCP 长连接Handshake。开辟轻量通道Channel在 TCP 连接上创建一个 Channel。原理TCP 创建销毁极耗资源Channel 复用 TCP是并发操作的最小单位。声明交换机与队列通过 Channel 发送指令确保 Broker 上存在目标 Exchange 和 Queue若不存在则自动创建。发送并确认生产者将消息含 Routing Key发送给 Exchange。同时开启Publisher Confirm模式。Broker 收到消息并成功路由到队列后会异步回传一个ACK给生产者。生产者必须收到这个 ACK才能认为“消息已安全送达”否则重发。阶段二Broker 内部处理中转与存储交换机路由Exchange RouterExchange 接收到消息根据 Binding 关系和 Routing Key 进行匹配Direct/Fanout/Topic。入队与持久化Enqueue Persist匹配成功的消息进入 Queue。关键原理如果队列开启了持久化durabletrue消息会先刷入磁盘再返回 ACK 给生产者确保断电不丢。等待消费消息在 Queue 中排队FIFO 顺序但优先级队列可调整等待消费者来拉取。阶段三消费者收消息收件拉取或推送Pull/Push消费者通过 Channel 订阅 Queue推荐basic.consume推模式长连接。业务处理与手动 ACK最重要消费者拿到消息执行发短信逻辑。原理执行完毕后消费者必须发送basic.ack给 Broker。删除消息Broker 收到消费者的 ACK 后才将 Queue 中的这条消息标记删除。如果消费者没发 ACK 就宕机Broker 会将该消息重新放回队头发给下一个消费者保证不丢。六如何实现生产者和消费者这里以我自己的工单分配的小项目举例想要达到的效果就是前端创建工单之后不必再等待数据落库而是创建之后直接返回一个“工单创建中请稍后查看”后续再工单详情列表可以再查看创建的新工单虽然这个功能看起来挺简单实际确实能感受到耗时上的差异加入MQ之后可以说是点击即完成而未加之前大概可能要等个半秒到一秒之间。下面这部分与大家看到的一些讲解视频可能不一样教程视频里是用的原生RabbitMQ Java客户端而我这里因为是项目工程化直接通过Spring封装代理了避免写重复代码。6.1 引入包依赖pom文件不再单独配置版本RabbitMQ默认遵循AMQP协议starter会自动选择适配的版本dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependencyyml文件rabbitmq: host: localhost port: 5672 username: guest password: guestSpringBoot starter自动读取 yml自动创建 Connection、Channel创建一个RabbitConfig配置类Configuration EnableRabbit public class RabbitConfig { /** 工单创建队列 */ public static final String TICKET_QUEUE ticket.queue; /** 工单交换机Direct */ public static final String TICKET_EXCHANGE ticket.exchange; /** 工单创建路由键 */ public static final String TICKET_ROUTING_KEY ticket.create; Bean public Queue ticketQueue() { // durabletrue消息持久化重启不丢 return new Queue(TICKET_QUEUE, true); } Bean public DirectExchange ticketExchange() { // durabletrue / autoDeletefalse交换机持久化 return new DirectExchange(TICKET_EXCHANGE, true, false); } Bean public Binding ticketBinding() { return BindingBuilder.bind(ticketQueue()).to(ticketExchange()).with(TICKET_ROUTING_KEY); } /** * JSON 消息转换器RabbitTemplate 发送与 RabbitListener 接收统一使用 * 收发双方均为 Ticket JSON 而非 Java 原生序列化。 */ Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); } }1.EnableRabbit开启Spring AMQP的RabbitListener消费者注解支持2. starter自动读取yml中rabbitmq连接配置自动装配 ConnectionFactory、RabbitTemplate 放入Spring容器3. 开发者使用Bean 定义 Queue队列、Exchange交换机、Binding绑定交给Spring容器4. 项目启动后Spring会自动向RabbitMQ服务声明这些队列、交换机不存在则创建5. RabbitTemplate用于发送消息RabbitListener 监听队列接收消息6.2 生产者代码// 投递 MQ真正的落库由监听器异步执行接口立即返回 rabbitTemplate.convertAndSend(RabbitConfig.TICKET_EXCHANGE, RabbitConfig.TICKET_ROUTING_KEY, ticket); log.info(工单已投递MQ等待异步落库: id{}, ticketNo{}, customerId{}, ticket.getId(), ticket.getTicketNo(), operator.userId()); return 工单创建中请稍后查看;在yaml配置文件里配置的连接信息Spring内部自动new ConnectionFactory 交给容器管理RabbitTemplate作为已经封装好的工具类底层内置连接池自动复用 Connection、Channel自动关闭回收防止连接泄漏。convertAndSend方法需要传入三个必填参数1.exchange交换机名称,消息先发给这个交换机传空字符串 使用 Rabbit 默认直连交换机。2.routingKey路由 key,交换机根据这个 key匹配绑定关系把消息投递到对应队列。3.message消息载体,Spring 自动调用消息转换器默认 Jackson把对象转为 JSON 字节数组6.3 消费者代码RabbitListener(queues RabbitConfig.TICKET_QUEUE) Transactional(rollbackFor Exception.class) public void handleCreate(Ticket ticket) { // 幂等保护消息可能被重复投递消费者 ack 前宕机等已存在则跳过 if (ticketMapper.selectById(ticket.getId()) ! null) { log.warn(工单已存在跳过重复消费: id{}, ticket.getId()); return; } // 异步消费线程没有登录上下文重建创建人信息保证操作日志留痕 restoreLoginUser(ticket.getCustomerId()); try { // 1. AI 智能分析开关关闭时跳过保留客户填写的分类/优先级 // 2. 真正落库 CREATE 操作日志同事务BEFORE_COMMIT 原子提交 ticketMapper.insert(ticket); stateMachine.recordCreate(ticket, 客户创建工单); // 3. 自动分配开关开启且有在线客服时PENDING_ASSIGN - IN_PROGRESS log.info(工单异步落库成功: id{}, ticketNo{}, status{}, category{}, priority{}, aiEnabled{}, degraded{}, ticket.getId(), ticket.getTicketNo(), ticket.getStatus(), ticket.getCategory(), ticket.getPriority(), aiEnabled, aiDegraded); } finally { LoginUserHolder.clear(); } }这里我剪切过来MessageListener到相关代码并以此解释一下消费者消费消息的流程1.在RabbitListener()里填入queues RabbitConfig.TICKET_QUEUE监听刚才生产者生产消息的队列Transactional(rollbackFor Exception.class)事务注解中间任何流程报错全部操作回滚防止脏数据。2.RabbitMQ 有可能一条消息重复推送比如消费者处理一半宕机 这里先查库如果工单已经存过 → 直接 return不再重复插入避免重复生成数据。MQ 消费是独立异步线程没有前端传过来的登录用户信息通过自定义方法手动塞入登录上下文可暂时忽略3.主体在try-finally里真正的落库在此并输出日志方便排查错误最终理线程里的登录用户信息防止内存泄漏、用户上下文错乱。七RabbitMQ 交换机类型在发送信息时交换机和路由键需要进行绑定这里就涉及到Routing key的绑定规则不同的匹配程度对应的不同的交换机类型消息发放的规则也不同。7.1 direct最常用精确完全匹配创建queueBind时routingKey必须跟bindingKey一样发送消息时如果routingKey一样在同一交换机下会同时向所有routingKey相同的队列发消息。7.2 fanout忽略routingKey全部广播只要在同一交换机的队列都会收到同一份消息不做任何匹配7.3 topic主题交换机主要是模糊匹配支持通配符模糊匹配routingKeyroutingKey规范用 . 来分割多个单词比如key1.key2.key3。通配符* 表示匹配1个单词# 表示匹配0个或多个单词这里举例说明一下比如现在交换机跟queue_topic1通过key1.key2.key3.*连接跟queue_topic2通过key1.#连接跟queue_topic3通过*.key2.*.key4连接跟queue_topic4通过#.key3.key4连接。现在比如生产者发送消息routingkey的值为“key1”看下RabbitMQ控制台只有queue_topic2接受到消息再发送routingkey值为“key1.key2.key3.key4”可以看到四个队列都有消息其中queue_topic2中有两条其他为一条这几说明刚才的生产者其实向所有队列都发送了一条消息7.4 headers几乎用不到不看 routingKey对比消息头部headers属性键值对是否匹配。 绑定队列时设置一组 header 参数消息 headers 满足条件才路由。(这里其实还有点内容不过因为不太重要就不再过多详细解释可再自行了解)八RabbitMQ 集群搭建单机的RabbitMQ有两个致命问题1本台服务器宕机整个MQ不可用消息功能失效业务中断。 2单台机器的CPU内存磁盘都有限消息量一旦太大后就容易阻塞堆积。于是RabbitMQ集群便产生了核心目标是高可用负载分散以及横向扩容这里因为我是windows操作系统直接搭建很多坑而生产环境一般都用Linux。可以通过Docker Desktop搭建 3 节点 RabbitMQ 集群因为我还对Docker不太熟悉下一篇文章打算学习一下Docker到那时再回来补ovo九镜像队列9.1使用镜像队列的原因先看一下普通集群遇到的困境一个集群比如有三个节点当某一个节点生产消息时这个消息只存在原始创建队列的那台节点上如果此时这台机突然宕机或者关闭服务了这个消息无法消费其他节点也看不到重启之后消息直接丢失了。镜像队列是在普通队列的基础上配置队列镜像队列的消息会主动复制到多个节点当主节点宕机时自动挑选镜像节点升级为主节点消息不丢失服务也不中断。9.2 核心角色Master主节点生产者发消息只会投递到 master所有消费也只能从 master 读取。镜像节点不承担读写分流Slave镜像副本异步同步 master 的消息平时只做备份master 挂了才晋升为主节点。9.3镜像队列策略 Policy核心配置通过rabbitmqctl set_policy定义规则哪些队列要镜像、存几份副本例如这个规则意思就是创建一条名叫ha-two的策略匹配所有的队列每个队列保存两份副本自动同步历史消息rabbitmqctl set_policy ha-two ^ {ha-mode:exactly,ha-params:2,ha-sync-mode:automatic}ha-mode(镜像模式)可选all镜像到全部节点exactly指定副本总数量nodes手动写死哪些节点做镜像。ha-params队列副本总量配合exactly使用这里填2代表一个主master一个镜像slave一次一共存两份数据ha-sync-mode消息同步模式automatic自动和manual手动十负载均衡‑HAProxy10.1作用RabbitMQ集群有多个节点我们自己随便指定一个当主节点都可能会挂掉HAProxy则是一个高性能的负载均衡软件有四大主要功能1统一入口所有生产者消费者只连接HAProxy地址不用维护MQ节点IP2负载均衡客户端连接均衡分发到后端各个RabittMQ节点3健康检查自动探测MQ节点是否存活故障节点自动剔除不转发连接4支持长连接TCP代理适配RabbitMQ的AMQP协议10.2核心概念frontend前端对外暴露端口接收客户端连接监听 5672backend后端配置所有 RabbitMQ 真实节点地址balance 负载策略roundrobin轮询默认均匀分发连接MQ 场景常用source根据客户端 IP 哈希固定客户端连接同一个节点health check 健康检测定期探测 MQ 节点宕机自动下线十一消息队列面试题Q1为什么你的系统要用消息队列三大经典用途参考答案我主要是基于解耦、异步、削峰这三个核心场景来用的。异步处理对应你的超时场景比如用户发起复杂查询同步查 DB 会超时。改成 MQ 后主线程只负责把任务扔进队列立即返回消费者慢慢处理彻底规避了 HTTP 超时。解耦订单系统只需要把消息发给 MQ下游的短信、积分、物流系统各自订阅。订单系统不需要知道谁在消费代码零侵入。流量削峰秒杀时瞬间百万请求直接打 DB 必挂。MQ 作为缓冲后端消费者按照自己能承受的速度比如每秒 2k 条拉取处理保证系统不崩溃。Q2RabbitMQ 怎么保证消息不丢生产级核心保障参考答案消息丢失可能发生在生产者、Broker、消费者三个环节我的应对方案是“三道防线”生产者端防丢开启Publisher Confirm发布确认机制。消息到达 Exchange 或落盘后Broker 会异步回传 ACK。若生产者没收到 ACK 或收到 NACK则触发重试或本地补偿日志。Broker 端防丢持久化Durable三件套。交换机持久化、队列持久化并且发消息时设置delivery_mode2消息持久化。消息先刷盘再回 ACK哪怕 RabbitMQ 宕机重启数据也能从磁盘恢复。消费者端防丢关闭自动 ACKautoAckfalse采用手动确认。业务逻辑执行成功如数据库事务提交后再调用basicAck告诉 Broker 删消息。如果业务抛异常或服务宕机Broker 收不到 ACK会把消息重新入队发给其他消费者。Q3如何处理消息重复消费幂等性参考答案因为 Q2 中我们开启了重试机制MQ 保证的是At-least-once至少一次所以重复消费是必然存在的。解决思路是消费端自己做幂等处理我们项目常用这三种方案方案一数据库唯一键利用数据库唯一约束。比如订单消息直接拿order_id作为数据库唯一索引重复插入会报错直接捕获异常回传 ACK 即可不业务处理。方案二Redis 分布式锁/状态位消费前用业务ID作为 Key 查询 Redis。若Key存在说明已处理过直接回 ACK 丢弃若不存在处理业务成功后将 Key 存入 Redis 并设置合理过期时间。方案三业务状态机更新数据时在 SQL 中加状态判断。例如UPDATE order SET status SUCCESS WHERE order_id ? AND status INIT。如果影响行数为 0说明已被更新属于重复请求。Q4如果消费者处理不过来消息堆积严重怎么办参考答案这是生产环境最容易出事故的问题。我的处理思路分“临时救火”和“根本解决”临时扩容救火先排查消费者是否被阻塞如 DB 连接池耗尽。如果是性能瓶颈直接快速扩容消费者实例并增加并发消费线程数。但注意队列是 FIFO 的如果新增消费者绑定到同一队列普通模式下它们会竞争消费如果是 Topic 模式可以临时新增一个队列绑定同一交换机启动新的消费者专门处理积压。紧急转移降级如果积压了几百万条且持续写入干脆把旧消息转移到另一个临时队列用专门的后台程序慢慢分析处理。主业务切回正常。根本解决排查 SQL 或下游接口是否响应慢优化消费者逻辑如开启批量拉取basicQos(100)或者评估是否需要用 MQ 存储海量数据——如果涉及大数据量这可能更适合改用 Kafka。Q5如何保证消息的顺序性顺序消费参考答案RabbitMQ 本身只能保证同一个队列Queue内部的消息是严格 FIFO先进先出的但无法保证多个队列之间的顺序。解决思路依业务而定方案一单队列单消费者如果全局顺序要求极高例如数据库 Binlog 同步就把所有相关消息路由到同一个 Queue并且只绑定一个消费者。这样性能会下降但绝对保序。方案二分区顺序最推荐利用一致性哈希Consistent Hash或取模路由。例如用order_id % 10决定消息进入哪个队列。相同的订单号永远进入同一个队列由一个消费者处理。这样既保证了单个订单的创建、支付、发货顺序又利用了多队列并发提升吞吐量。十二结语本文是我作为一个Java小白初学者的学习总结所有不懂的地方基本也是反复询问AI讲解之后才敢网上写的对于集群搭建以及负载均衡的部分其实我并没有实际操作这部分在文章内主要只是理论打算下一步先看一下另一个大名鼎鼎的消息中间件Kafka再进一步学一下Docker到时候回来重新对比两者的区别重新自己实操部署一下。最近实习参与一下组内的项目验收帮忙做一些功能的测试和需求规格说明书的编写都是一些基础活吧果然如网上所说实习都是打杂不过应该也不用太着急我对自己的定位还是比较清晰的我们学校开学早下周一正式上课公司离学习挺近打算继续实习到10月底吧11月回去准备再考考六级刷刷分准备一下哈寒假能不能找到杭州的中厂实习一步一步来。文章有任何错误或者不足的地方欢迎各位佬讨论。