RabbitMQ

一、基础概念

RabbitMQ 是一款基于 Erlang 开发的开源消息中间件。
核心作用:系统之间异步解耦、削峰填谷、异步通知、流量缓冲。
Broker:RabbitMQ 服务实例整体称之为 Broker。

核心内部对象

  1. 虚拟主机 vhost
    逻辑隔离单元,相当于命名空间。
    同一个RabbitMQ服务可以创建多个vhost;每个vhost内部拥有独立交换机、队列,权限配置绑定在vhost上。
    默认虚拟主机:/
    业务规范:不同业务线使用独立vhost隔离。

  2. 连接 Connection
    应用程序与RabbitMQ服务建立的TCP长连接。
    创建销毁成本高,业务应用不要频繁新建连接,建议维持少量长连接。

  3. 通道 Channel
    一条TCP连接之上建立的轻量级逻辑通道。
    所有发布消息、消费消息、声明队列等操作都在Channel执行。

关键特性:Channel不支持多线程并发操作

  1. 交换机 Exchange
    消息进入RabbitMQ的入口。生产者发送消息只能投递到交换机,不会直接发送到队列。
    交换机接收消息后,依据绑定规则将消息路由至队列。

  2. 绑定 Binding
    交换机和队列之间的关联关系,定义路由匹配规则。

  3. 队列 Queue
    消息最终存储容器,消息保存在队列中等待消费者处理。

二、交换机类型

1. Direct 直连交换机

路由规则:消息路由键 和 绑定键完全相等,消息才会路由至对应队列。
适用场景:精准点对点消息、任务分发。

2. Fanout 扇形交换机

路由规则:无视路由键,收到消息后广播给所有绑定该交换机的队列。
适用场景:系统广播通知、多服务同时接收同一条事件。

3. Topic 主题交换机

支持通配符模糊匹配,路由键使用英文句号 . 分割层级。
- *:匹配任意一个层级单词
- #:匹配零个或多个层级单词
示例:
路由键 order.create
绑定键 order.* 可以匹配;order.# 也可以匹配。
适用场景:微服务事件总线、日志分类订阅。

4. Headers 头交换机

不依赖路由键,依靠消息header内键值对进行匹配。
性能较差,业务开发极少使用。

RabbitMQ扩展交换机(非原生四类)

  1. 死信交换机 DLX:承接处理失败、过期、被拒绝的消息
  2. 延迟交换机(依赖插件 rabbitmq-delayed-message-exchange):实现定时消息

三、队列常用配置参数

声明队列时可配置属性:
1. durable:队列持久化。true代表服务重启后队列依然存在;false重启队列消失。
2. exclusive:排他队列。仅创建这条通道对应的连接可见,连接断开自动删除队列,多用于临时回调队列。
3. auto-delete:自动删除。队列存在过消费者,所有消费者断开连接后,自动删除队列。

队列附加扩展参数:
- x-message-ttl:队列内统一消息过期时间,单位毫秒
- x-max-length:队列最大消息数量,超出后触发溢出策略
- x-dead-letter-exchange:指定死信交换机
- x-dead-letter-routing-key:死信转发使用的路由键

四、消息基础结构

一条消息包含两大部分:消息头部属性 + 消息体(业务二进制数据)
常用消息属性:
1. delivery_mode:1非持久,2持久消息
2. expiration:单条消息过期时长(毫秒)
3. correlation_id:消息唯一标识,常用于RPC调用
4. reply_to:RPC场景,消费者返回结果使用的临时队列名称
5. headers:自定义键值对,可以存放业务自定义参数、重试次数等信息

五、五种经典工作模式

1. 简单模式

生产者 → 默认交换机 → 队列 → 单个消费者
默认交换机名称为空字符串 "",系统内置直连类型交换机。隐式规则:队列会自动和默认交换机绑定,绑定键名称等于队列名。

2. 工作模式(Work Queue)

生产者 → 队列 → 多个消费者竞争消费
默认策略:轮询分发,消息均匀分给各个消费者。
开启QoS配置后,可以实现公平分发:消费完成确认之后,Broker才推送下一条消息。

3. 发布订阅模式

生产者 → Fanout交换机 → 多个绑定队列 → 各自消费者
一条消息被投递到全部队列,实现广播效果。

4. 路由模式

生产者 → Direct交换机,携带路由键
只有绑定键与消息路由键完全匹配的队列收到消息。

5. 主题模式

生产者 → Topic交换机,带分层路由键,依靠通配符匹配多个队列。企业项目最常用。

六、可靠性保障机制

1. 生产者投递保障:发布确认 Publisher Confirms

开启确认模式后,Broker成功接收消息完成持久化后,向生产者返回ack;消息异常丢失则返回nack。
业务场景优先使用异步回调确认;不推荐同步阻塞等待。

补充:RabbitMQ支持事务模式实现投递保障,但是性能损耗极大,生产环境禁止使用。

2. 服务宕机消息不丢失:持久化

两个条件必须同时开启:
1. 队列声明设置持久化 durable=true
2. 发送消息设置 delivery_mode=2(持久消息)
只配置其中一项,重启后消息依旧丢失。

3. 消费者消费保障:消息确认机制

消费存在两种模式:
1. 自动确认 autoAck=true
消息推送到消费端,服务立刻直接删除消息。风险:消费者程序崩溃,消息直接丢失。业务系统禁止使用。
2. 手动确认 autoAck=false(标准方案)
Broker保留消息,等待消费端主动发送指令:
- basic.ack:消费正常完成,删除消息
- basic.reject:拒绝单条消息
- basic.nack:支持批量拒绝消息,附带 requeue 参数
- requeue=true:消息重新放回队列重新消费
- requeue=false:消息丢弃或转入死信队列

4. 死信队列(DLX)

消息满足以下任意条件,会被路由至死信交换机:
1. 消息被拒绝,且 requeue=false
2. 消息超过过期时间TTL
3. 队列消息数量达到上限,新消息溢出

用途:收集异常消息,单独排查;配合实现消息有限次数重试逻辑。

七、QoS 流量控制

配置参数 prefetch_count
含义:服务端最多向同一个消费者推送多少条还没有确认ack的消息。
作用:避免短时间大量消息压垮消费应用,控制消费负载。
限制:仅对持续推送消费模式生效,循环主动拉取消息模式不生效。

八、高级特性

1. 消息TTL过期

两种配置方式:
1. 队列级别:队列参数 x-message-ttl,队列内所有消息统一过期时长
2. 消息级别:每条消息设置 expiration 属性
两条同时存在,取时间更小的值。
注意原生TTL存在缺陷:只有队列头部消息过期后,后面到期消息才会被处理。不能直接实现精准定时任务。

2. 延迟消息

安装插件 rabbitmq-delayed-message-exchange
创建类型为 x-delayed-message 的交换机。消息发送后先暂存,到达设定时间再路由至目标队列。用于订单超时关闭等定时场景。

3. 队列类型

  1. Classic 经典队列:默认普通队列
  2. Quorum 仲裁队列(3.8版本后官方主推)
    基于Raft一致性算法,集群环境具备高可用,新项目集群优先选用
  3. Stream 流队列:面向海量日志、持续流式消费场景

九、集群模式

1. 普通集群

队列元数据同步至所有节点;但是消息只存在创建队列的节点。该节点故障后,队列无法读写,无法保证高可用。

2. 仲裁队列集群(生产推荐)

消息副本分散存储在多个集群节点,遵循Raft协议;集群多数节点正常运行即可提供服务,支持自动故障转移。

3. 联邦集群 Federation

多用于跨机房、异地多活场景,普通业务不需要部署。

十、常用插件

插件启动命令:rabbitmq-plugins enable 插件名
1. rabbitmq_management:Web可视化管理后台,端口15672
2. rabbitmq-delayed-message-exchange:延迟消息插件
3. rabbitmq_tracing:消息轨迹追踪,调试排查消息流向
4. rabbitmq_mqtt、rabbitmq_stomp:支持物联网协议接入

端口汇总:
- 5672:客户端通信端口
- 15672:Web管理页面
- 25672:集群节点内部通信端口

十一、线上常见三大问题解决方案

1. 消息丢失

完整方案:生产者开启发布确认 + 队列持久化+持久消息 + 消费端手动ack + 异常消息转入死信队列。

2. 重复消费

网络波动会造成ack信号丢失,RabbitMQ会重新投递消息。
RabbitMQ无法从根源杜绝重复投递。解决方案:消费业务实现幂等性
常见实现:业务唯一编号利用数据库唯一约束、Redis标记消费状态。

3. 消息大量积压

诱因:消费处理速度跟不上生产速度。
处理方案:
1. 横向扩容消费者实例
2. 优化消费内部业务逻辑,耗时操作异步化
3. 业务流量拆分,使用独立队列隔离流量
风险:海量积压消息会持续占用内存、磁盘,触发服务流控,阻塞生产者发消息。

十二、开发最佳实践

  1. Channel 不允许多线程共用;多线程场景,每个线程使用独立Channel。
  2. 复用TCP连接,不要频繁创建Connection。
  3. 业务消费统一采用持续推送模式(basic.consume),禁止循环调用主动拉取接口basic.get,性能极低。
  4. 消息不要携带超大业务数据;超过16MB消息不建议直接投递,传递文件地址。
  5. 消息重试逻辑必须增加最大重试次数,依靠死信队列兜底,避免无限循环消费。
  6. 合理配置连接心跳,防止防火墙长时间空闲断开TCP连接。