欢迎光临
我们一直在努力

RabbitMQ 基础篇:从零开始掌握消息队列核心概念

目录

  • 一、课程背景
  • 二、初识 MQ:同步调用 vs 异步调用
    • 2.1 同步调用
    • 2.2 异步调用
    • 2.3 小结
  • 三、MQ 技术选型:RabbitMQ / ActiveMQ / RocketMQ / Kafka
  • 四、RabbitMQ 基本介绍
    • 4.1 整体架构与核心概念
    • 4.2 官方工作流概览
  • 五、SpringAMQP 快速入门
    • 5.1 引入依赖
    • 5.2 配置 RabbitMQ 服务端信息
    • 5.3 发送消息
    • 5.4 接收消息
  • 六、Work Queues(工作队列)
    • 6.1 任务模型
    • 6.2 消费者消息推送限制
    • 6.3 小结
  • 七、交换机类型
    • 7.1 Fanout 交换机(广播)
    • 7.2 Direct 交换机(定向路由)
    • 7.3 Topic 交换机(话题)
  • 八、声明队列和交换机
    • 8.1 基于 Bean 的方式
    • 8.2 基于 @RabbitListener 注解的方式
  • 九、消息转换器
  • 十、写在最后

一、课程背景

我们以一个常见的登录场景为例:

用户登录 → 用户微服务查询并校验用户信息 → 风控微服务记录登录信息并判断登录风险 → 短信微服务发送风险短信 → 通知用户。

如果登录事件希望做风控和短信提醒,传统方式就是登录完成后依次调用风控、短信两个服务。但在真实业务里,"风控"和"短信通知"跟登录主链路关系并不紧密,完全可以异步处理——这就轮到 MQ(消息队列)登场了。

整个 RabbitMQ 课程分为基础篇和高级篇:

基础篇

高级篇

同步和异步

发送者重连

MQ 技术选型

发送者确认

数据隔离

MQ 持久化

SpringAMQP

LazyQueue

Work 模式

消费者确认

MQ 消息转换器

失败重试

发布订阅模式

业务幂等

消息堆积问题处理

延迟消息

本文聚焦于基础篇的全部内容。


二、初识 MQ:同步调用 vs 异步调用

2.1 同步调用

以黑马商城的"余额支付"为例。一次支付要经过 5 个步骤:

  • 扣减余额(用户服务)
  • 更新支付状态(支付服务)
  • 更新订单状态(交易服务)
  • 短信通知用户(通知服务)
  • 增加用户积分(积分服务)
  • 每一步 50ms,串行下来一次请求需要 300ms。这种链式同步调用会出现三个问题:

    • 拓展性差:每加一个下游业务就要改支付服务的代码;
    • 性能下降:链路越长,整体 RT 越高;
    • 级联失败:任何一环出问题都会让整条链路失败。

    小结:

    • 同步调用的优势是时效性强,调用方等到结果后才返回;
    • 缺点是拓展性差、性能下降、级联失败。

    2.2 异步调用

    异步调用本质上是基于消息通知的方式,一般包含三个角色:

    • 消息发送者:投递消息的人,就是原来的调用方;
    • 消息代理:管理、暂存、转发消息,可以理解成"消息服务器";
    • 消息接收者:接收并处理消息的人,就是原来的服务提供方。

    用异步改造余额支付,支付服务只同步做完"扣减余额 + 更新支付状态"这两件强关联的事,剩下的订单状态、短信、积分通过 Broker 异步分发。改造后访问时长从 300ms 缩短到 100ms,并且具备四个明显优势:

    • 解除耦合,拓展性强
    • 无需等待,性能好
    • 故障隔离:下游服务故障不影响上游业务
    • 缓存消息、流量削峰填谷:把脉冲式的 QPS 削平成平缓曲线

    小结:

    • 异步调用的优势是耦合度低、拓展性强、性能好、故障隔离、流量削峰;
    • 缺点是不能立即得到调用结果,时效性差;不确定下游业务是否执行成功;业务安全依赖于 Broker 的可靠性。

    2.3 小结

    维度

    同步调用

    异步调用

    时效性

    强,等结果

    弱,后置处理

    拓展性

    性能

    随链路下降

    稳定

    耦合度

    故障影响

    级联失败

    隔离

    流量处理

    难以削峰

    削峰填谷

    适用场景

    强一致性、强依赖

    最终一致性、可异步化的非关键路径

    ⚠️ 什么时候用异步? 主链路完成后才需要的非关键业务(风控、通知、积分、日志、审计…),都适合异步化。但涉及金钱、需要强一致性的核心链路,依然建议同步 + 本地事务。


    三、MQ 技术选型:RabbitMQ / ActiveMQ / RocketMQ / Kafka

    MQ(MessageQueue),中文"消息队列",本质上就是异步调用中的 Broker。业界常见的有四款产品:

    维度

    RabbitMQ

    ActiveMQ

    RocketMQ

    Kafka

    公司/社区

    Rabbit

    Apache

    阿里

    Apache

    开发语言

    Erlang

    Java

    Java

    Scala & Java

    协议支持

    AMQP, XMPP, SMTP, STOMP

    OpenWire, STOMP, REST, XMPP, AMQP

    自定义协议

    自定义协议

    可用性

    一般

    单机呑吐量

    一般

    非常高

    消息延迟

    微秒级

    毫秒级

    毫秒级

    毫秒以内

    消息可靠性

    一般

    一般

    怎么选?

    • 强业务可靠、对延迟敏感:RabbitMQ(微秒级延迟 + 高可靠,电商订单场景首选);
    • 日志、大数据、流计算、IoT 等超高吞吐场景:Kafka;
    • 阿里系、有顺序消息、事务消息要求:RocketMQ;
    • 老系统维护、AMQP 协议强需求:ActiveMQ。

    本课程选 RabbitMQ 作为讲解对象,覆盖绝大多数业务场景。


    四、RabbitMQ 基本介绍

    4.1 整体架构与核心概念

    RabbitMQ 的整体架构里有五个核心概念,先记牢:

    • virtual-host:虚拟主机,起到数据隔离的作用(类似 MySQL 的 database);
    • publisher:消息发送者;
    • consumer:消息的消费者;
    • queue:队列,用于存储消息;
    • exchange:交换机,负责路由消息。

    publisher 把消息发给 exchange,exchange 按规则路由到一个或多个 queue,最后由 consumer 消费。每个 VirtualHost 是一个独立的小型 RabbitMQ 实例。

    4.2 官方工作流概览

    Publisher → Exchange(按 Binding 规则)→ Queue → Consumer,这是 RabbitMQ 的标准消息流。理解清楚"exchange 只负责转发、不存储消息"这一点,对后面的代码实现非常关键。


    五、SpringAMQP 快速入门

    官方地址:Spring AMQP

    Spring AMQP 是基于 AMQP 协议定义的一套 API 规范,包含两部分:

    • spring-amqp:基础抽象;
    • spring-rabbit:底层默认实现(也就是我们用的 RabbitMQ)。

    下面用 4 步把一个最简单的"生产者发送字符串、消费者接收字符串"跑起来。

    5.1 引入依赖

    在父工程中引入 spring-amqp 依赖,这样 publisher 和 consumer 服务都可以使用:

    <dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>

    5.2 配置 RabbitMQ 服务端信息

    在每个微服务的 application.yml 中配置,这样服务才能连到 RabbitMQ:

    ```yaml
    spring:
    rabbitmq:
    host: 192.168.150.101 # 主机名
    port: 5672 # 端口
    virtual-host: /hmall # 虚拟主机
    username: hmall # 用户名
    password: 123 # 密码
    ```

    5.3 发送消息

    SpringAMQP 提供了 RabbitTemplate 工具类,方便我们发送消息:

    ```java
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Test
    public void testSimpleQueue() {
    // 队列名称
    String queueName = "simple.queue";
    // 消息
    String message = "hello, spring amqp!";
    // 发送消息
    rabbitTemplate.convertAndSend(queueName, message);
    }
    ```

    5.4 接收消息

    SpringAMQP 提供声明式的消息监听,我们只需通过 @RabbitListener 注解在方法上声明要监听的队列名称,将来 SpringAMQP 就会自动把消息传递给当前方法:

    ```java
    @Slf4j
    @Component
    public class SpringRabbitListener {

    @RabbitListener(queues = "simple.queue")
    public void listenSimpleQueueMessage(String msg) throws InterruptedException {
    log.info("spring 消费者接收到消息:【" + msg + "】");
    if (true) {
    throw new MessageConversionException("故意的");
    }
    log.info("消息处理完成");
    }
    }
    ```

    💡 这一节抛出的 MessageConversionException 是为高级篇里讲"失败重试 / 死信队列"故意埋的伏笔,基础篇先不用深究。


    六、Work Queues(工作队列)

    6.1 任务模型

    Work Queues 又叫任务模型,简单说就是让多个消费者绑定到一个队列,共同消费队列中的消息。结构如下:

    ```
    publisher → queue → consumer1
    ↘ consumer2
    ```

    它的特点是:

    • 同一消息只会被一个消费者处理;
    • 通过多消费者能显著加快消息处理速度。

    6.2 消费者消息推送限制

    默认情况下,RabbitMQ 会把消息依次轮询投递给绑定在队列上的每一个消费者。但它并不会等消费者"处理完"再投递下一条,而是按"平均分配消息数量"的方式分发,可能造成消息堆积。

    解决办法是修改 application.yml,设置 prefetch 值为 1,确保同一时刻最多投递给消费者 1 条消息:

    ```yaml
    spring:
    rabbitmq:
    listener:
    simple:
    prefetch: 1 # 每次只能获取一条消息,处理完成才能获取下一个消息
    ```

    这样无论是消费者 1 还是消费者 2,处理完一条再处理下一条,就能做到"能者多劳"——机器性能好的自然消费得快。

    6.3 小结

    Work 模型的使用要点:

    • 多个消费者绑定到一个队列,可以加快消息处理速度;
    • 同一条消息只会被一个消费者处理;
    • 通过 prefetch 控制消费者预取的消息数量,处理完一条再处理下一条,实现能者多劳。

    七、交换机类型

    真正生产环境里,消息都是经过 exchange 发送的,而不是直接发送到队列。交换机共有三种类型:

    • Fanout:广播;
    • Direct:定向;
    • Topic:话题。

    7.1 Fanout 交换机(广播)

    Fanout Exchange 会将收到的消息广播到每一个跟其绑定的 queue,所以也叫广播模式。

    ```
    publisher → Fanout exchange
    ├─→ queue1 → consumer1(订单服务)
    └─→ queue2 → consumer3(日志服务)
    ```

    典型场景:下单成功 → 通知订单服务更新状态 + 日志服务记录 + 积分服务加积分,各业务方互不影响、各自处理。

    小结:

    • 交换机的作用:接收 publisher 发送的消息、按规则路由到与之绑定的队列;
    • FanoutExchange 会把消息路由到每个绑定的队列。

    7.2 Direct 交换机(定向路由)

    Direct Exchange 会将接收到的消息根据规则路由到指定的 Queue,因此称为定向路由。

    规则有三条:

  • 每一个 Queue 都与 Exchange 设置一个 BindingKey;
  • 发布者发送消息时,指定消息的 RoutingKey;
  • Exchange 将消息路由到 BindingKey 与消息 RoutingKey 一致的队列。
  • ```
    publisher → Direct exchange
    ├─ queue1 (BindingKey: red, blue) → consumer1 收 blue/red
    └─ queue2 (BindingKey: yellow, red) → consumer2 收 yellow/red
    ```

    发送消息时 RoutingKey 为 red 时,queue1 和 queue2 都会收到;RoutingKey 为 yellow 时只有 queue2 收到。

    7.3 Topic 交换机(话题)

    TopicExchange 与 DirectExchange 类似,区别在于 routingKey 可以是多个单词的列表,并且以 . 分割。Queue 与 Exchange 指定 BindingKey 时可以使用通配符:

    • #:代表 0 个或多个单词;
    • *:代表 1 个单词。

    比如要区分中国/日本 + 新闻/天气四种消息:

    ```
    china.news 代表中国的新闻消息;
    china.weather 代表中国的天气消息;
    japan.news 代表日本新闻;
    japan.weather 代表日本的天气消息;
    ```

    ```
    publisher → Topic exchange
    ├─ queue1 (china.#) → consumer1 收 china.*
    ├─ queue2 (japan.#) → consumer2 收 japan.*
    ├─ queue3 (#.weather) → consumer3 收 *.weather
    └─ queue4 (#.news) → consumer4 收 *.news
    ```

    Direct 与 Topic 的差别:

    • Topic 交换机接收的消息 RoutingKey 可以是多个单词,以 . 分隔;
    • Topic 交换机与队列绑定时的 BindingKey 可以指定通配符;
      • #:代表 0 个或多个单词;
      • *:代表 1 个单词。

    🎯 Topic 模型是业务最常用的,能用一套规则覆盖很多场景(省份+业务类型、商品类目+操作类型…)。


    八、声明队列和交换机

    SpringAMQP 提供几种方式声明队列、交换机及其绑定关系:

    • Queue:声明队列,可以用工厂类 QueueBuilder 构建;
    • Exchange:声明交换机,可以用 ExchangeBuilder 构建;
    • Binding:声明队列和交换机的绑定关系,可以用 BindingBuilder 构建。

    下面两种写法都要掌握。

    8.1 基于 Bean 的方式

    例如,声明一个 Fanout 类型的交换机,并且创建两个队列与它绑定:

    ```java
    @Configuration
    public class FanoutConfig {
    // 声明 fanoutExchange 交换机
    @Bean
    public FanoutExchange fanoutExchange(){
    return new FanoutExchange("hmall.fanout");
    }
    // 声明第 1 个队列
    @Bean
    public Queue fanoutQueue1(){
    return new Queue("fanout.queue1");
    }
    // 绑定队列 1 和交换机
    @Bean
    public Binding bindingQueue1(Queue fanoutQueue1, FanoutExchange fanoutExchange){
    return BindingBuilder.bind(fanoutQueue1).to(fanoutExchange);
    }
    // … 略,以相同方式声明第 2 个队列,并完成绑定
    }
    ```

    8.2 基于 @RabbitListener 注解的方式

    除了写 @Configuration,SpringAMQP 还提供了基于 @RabbitListener 注解来声明队列和交换机的方式:

    ```java
    @RabbitListener(bindings = @QueueBinding(
    value = @Queue(name = "direct.queue1"),
    exchange = @Exchange(name = "itcast.direct", type = ExchangeTypes.DIRECT),
    key = {"red", "blue"}
    ))
    public void listenDirectQueue1(String msg){
    System.out.println("消费者 1 接收到 Direct 消息:【" + msg + "】");
    }
    ```

    8.3 小结

    声明队列、交换机、绑定关系的 Bean 是什么?

    • Queue;
    • FanoutExchange / DirectExchange / TopicExchange;
    • Binding。

    基于 @RabbitListener 注解声明队列和交换机时常用的注解:

    • @Queue;
    • @Exchange。

    九、消息转换器

    Spring 对消息对象的处理是由 org.springframework.amqp.support.converter.MessageConverter 负责的。默认实现是 SimpleMessageConverter,它基于 JDK 的 ObjectOutputStream 完成序列化。

    这种默认实现有三个致命问题:

  • JDK 序列化有安全风险:反序列化漏洞一直存在;
  • JDK 序列化的消息太大:把类名、字段全写进去,浪费带宽;
  • JDK 序列化的消息可读性差:抓包看到的是二进制,没法直接看。
  • 解决思路是把 MessageConverter 替换成更轻量、更通用的格式,比如 JSON。生产环境最常用的就是 Jackson2JsonMessageConverter(这部分会在进阶篇中展开)。


    十、写在最后

    整篇基础篇其实讲的是一件事:消息从 publisher 出发,经过 exchange 和 queue,最终被 consumer 消费。围绕这条主线,我们掌握了:

  • 为什么要用异步:解耦 + 削峰 + 提速;
  • 怎么选型:业务可靠性选 RabbitMQ,海量日志选 Kafka;
  • RabbitMQ 核心组件:virtual-host / publisher / consumer / queue / exchange;
  • SpringAMQP 入门:依赖、配置、发送、接收一气呵成;
  • Work 模型:多消费者 + prefetch 实现能者多劳;
  • 三种交换机:Fanout 广播、Direct 定向、Topic 通配符;
  • 声明方式:Bean 形式 vs 注解形式;
  • 消息转换器:JDK 序列化不可用,要换成 JSON 之类更友好的方案。
  • 下一步进入高级篇,继续解锁发送者确认、消费者确认、失败重试、延迟消息、消息堆积处理等"踩坑必备"技能。

    如果本文对你有帮助,欢迎点赞 👍、收藏 ⭐、评论 💬 三连。后续会更新 RabbitMQ 高级篇,敬请期待。

    赞(0)
    未经允许不得转载:171主机测评 » RabbitMQ 基础篇:从零开始掌握消息队列核心概念
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址