本文参考 B 站黑马:https://www.bilibili.com/video/BV1mN4y1Z7t9

1、MQ 介绍

1.1、同步调用

例如:调用支付服务(用户付款成功)时,需要依次调用多个服务(订单服务(更新订单状态)、短信服务(短信通知用户)、积分服务(增加用户积分)等)。

缺点

  • 性能下降。消费者需要等待所有提供者依次执行完成。
  • 级联失败。如果提供者出现故障,则消费者同样出现故障。
  • 耦合度高。如果新增业务需求,则需修改原有代码。

优点

  • 时效性高。可以立即得到结果。

适用场景

  • 对时效性要求高的场景。例如:在查询订单时,同时查询用户信息。

1.2、异步调用

例如:调用支付服务(用户付款成功)时,不用关心后续服务执行情况(订单服务(更新订单状态)、短信服务(短信通知用户)、积分服务(增加用户积分)等)。

优点

  • 性能提升。
  • 故障隔离。
  • 耦合度低。
  • 流量削峰。

缺点

  • 时效性差。不能立即得到结果。
  • 不能确定下游业务是否执行成功。
  • 依赖 Broker(中间件,代理)的可靠性、吞吐量等。

适用场景

  • 高性能、高并发、高可用,对时效性要求不高的场景。

1.3、MQ 常见技术

MQ(Message Queue,消息队列)是消息中间件,用于在分布式系统中实现异步通信。MQ 是异步调用中的 Broker(中间件,代理)。

MQ 常见技术

  • RabbitMQ、RocketMQ:用于可靠性要求高的场景。例如:适合小型项目、或针对具体业务等。
  • ActiveMQ:性能一般,使用较少。
  • Kafka:用于吞吐量要求高、海量数据但可靠性要求低的场景。例如:针对日志、大数据等。
RabbitMQ ActiveMQ RocketMQ Kafka
公司/社区 Rabbit Apache 阿里开发(现为 Apache) Apache
开发语言 Erlang Java Java Scala 与 Java
协议支持(功能) AMQP、XMPP、SMTP、STOMP OpenWire、REST、AMQP、XMPP、STOMP 自定义协议 自定义协议
可用性 一般
单机吞吐量(并发) 一般 非常高
消息延迟 微秒级 毫秒级 毫秒级 毫秒级
消息可靠性 一般 一般

2、RabbitMQ 入门

2.1、RabbitMQ 安装

RabbitMQ 是基于 Erlang 语言开发的消息中间件,用于在分布式系统中实现异步通信。

官网:https://www.rabbitmq.com/。

安装(Linux 版,使用 Docker)

  • 一、上传 mq.tar 到 /tmp 目录,执行以下命令。
docker load -i mq.tar # 加载镜像
docker images # 查看镜像

# 运行
docker run \
  --name mq \
  --hostname mq1 \
  -p 15672:15672 \
  -p 5672:5672 \
  -d \
  rabbitmq:3.8-management

# 查看
docker ps
----------------
docker run \ # 运行
  --name mq \ # 名称
  --hostname mq1 \ # 主机名,以后集群使用
  -p 15672:15672 \ # 控制台端口,用于控制台登录
  -p 5672:5672 \ # 监听端口,用于消息通信
  -d \ # 后台运行
  rabbitmq:3.8-management # 镜像名称
  • 二、登录。访问 http://192.168.112.131:15672。账号/密码:guest/guest。

整体架构与核心概念

  • Vitual Host:虚拟主机,用于数据隔离。
  • Publisher:消息发布者。
  • Consumer:消息消费者。
  • Queue:队列,用于存储消息。
  • Exchange:交换机,用于路由消息。
              【RabbitMQ Server】
         →→→→【Exchange→→→→Queue】→→→→Consumer
Publisher→→→→【Exchange→→→→Queue】→→→→Consumer
         →→→→【Exchange→→→→Queue】→→→→Consumer

2.2、RabbitMQ 控制台操作

需求

在控制台中,创建 Queue、使用默认交换机 amq.fanout 发布消息。

步骤

  • 一、在控制台的 Queue 中,添加一个 Queue。(输入 Name 然后点击 【Add queue】)
Name: hello.queue1
Name: hello.queue2
Name: hello.queue3
  • 二、在控制台的 Exchange 中,绑定 Queue,发布消息。(在默认交换机 amq.fanout 中,输入 Name 然后点击【Bind】,输入 Payload 然后点击【Publish message】)
Payload: Hello,RabbitMQ!
  • 三、在控制台的 Queue 中,查看消息。(点击新添加的 Queue,点击 【Get messages】)

2.3、数据隔离机制

每个 User 都有自己的 Vitual Host。

每个 Vitual Host 都有自己的 Exchange、Queue。

通常不同项目或服务,使用自己的 User 与 Vitual Host,用于数据隔离。

需求

  • 一、创建一个 User。
  • 二、退出当前 User,使用新 User 创建一个 Vitual Host。
  • 三、测试不同 Vitual Host 之间的数据隔离现象。

示例

  • 一、创建一个 User。
Username: xiaofeng
Password: xiaofeng
          xiaofeng
Tag: administrator // 权限标签
  • 二、退出当前 User,使用新 User 的账号密码登录,并为该 User 创建一个 Vitual Host。
Name: /xiaofeng
  • 三、测试不同 Vitual Host 之间的数据隔离现象。

3、Spring AMQP 入门

3.1、Spring AMQP 介绍

AMQP(Advanced Message Queuing Protocol,高级消息队列协议)是微服务之间传递业务消息的开放标准(规范)。该协议与语言或平台无关,符合微服务对技术独立性的要求。

Spring AMQP 是基于 AMQP 定义的一套 API 规范,提供模板(RabbitTemplate)用于发布与接收消息。包括:Spring AMQP 是接口规范,Spring Rabbit 是底层实现(封装 RabbitMQ)。

官网:https://spring.io/projects/spring-amqp。

特点

  • 提供 RabbitListener 消息监听器,用于异步处理消息。
  • 提供 RabbitTemplate,用于发布与接收消息。
  • 提供 RabbitAdmin,用于自动声明队列、交换与绑定。

3.2、环境准备

步骤

  • 一、创建父项目 rabbitmq,导入 Spring AMQP 依赖。
  • 二、创建子模块 consumer,创建启动类,配置 YML。
  • 三、创建子模块 publisher,创建启动类,配置 YML。

示例

  • 一、创建父项目 rabbitmq,导入依赖。
<properties>
    <java.version>1.8</java.version>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
    <spring-boot.version>2.6.13</spring-boot.version>
</properties>

<dependencies>
    <!--Spring AMQP-->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    <!--Test-->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-test</artifactId>
        <scope>test</scope>
    </dependency>
</dependencies>

<dependencyManagement>
    <dependencies>
        <!--Spring Boot 管理依赖-->
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-dependencies</artifactId>
            <version>${spring-boot.version}</version>
            <type>pom</type>
            <scope>import</scope>
        </dependency>
    </dependencies>
</dependencyManagement>
  • 二、创建子模块 consumer,创建启动类,配置 YML。
@SpringBootApplication
public class ConsumerApplication {
    public static void main(String[] args) {
        SpringApplication.run(ConsumerApplication.class, args);
    }
}
spring:
  rabbitmq: # RabbitMQ 配置
    host: 192.168.112.131 # 主机
    port: 5672 # 监听端口
    virtual-host: /xiaofeng # 虚拟主机
    username: xiaofeng # 用户名
    password: xiaofeng # 密码
  • 三、创建子模块 publisher,创建启动类,配置 YML。
@SpringBootApplication
public class PublisherApplication {
    public static void main(String[] args) {
        SpringApplication.run(PublisherApplication.class, args);
    }
}
spring:
  rabbitmq: # RabbitMQ 配置
    host: 192.168.112.131 # 主机
    port: 5672 # 监听端口
    virtual-host: /xiaofeng # 虚拟主机
    username: xiaofeng # 用户名
    password: xiaofeng # 密码

3.3、消息队列模式

官网教程:https://www.rabbitmq.com/tutorials。

模式 交换机 说明
Hello World 简单模式 一个发布者,一个消费者。
Work Queues 工作模式 一个发布者,多个消费者互相竞争,可以提高消息处理速度。
Publish/Subscribe 发布/订阅模式 Fanout 一个发布者,多个消费者收到相同消息。
Routing 路由模式 Direct 一个发布者,多个消费者选择性的接收消息。
Topics 主题模式 Topic 一个发布者,多个消费根据主题选择性的接收消息。
RPC 请求/回复模式
Publisher Confirms 发布者确认模式

3.3、简单模式

需求

  • 一、在控制台中,创建队列 simple.queue。
  • 二、在 consumer 中,使用 @RabbitListener 监听 simple.queue,然后启动 Spring Boot
  • 三、在 publisher 中,使用 RabbitTemplate 向 simple.queue 发布消息。

示例

  • 一、在控制台中,创建队列 simple.queue。
  • 二、在 consumer 中,使用 @RabbitListener 监听 simple.queue,然后启动 Spring Boot
@Component
public class SpringRabbitListener {
    // Simple
    @RabbitListener(queues = "simple.queue")
    public void listenSimpleQueue(String msg) {
        System.out.println("【收到消息】:" + msg);
    }
}
  • 三、在 publisher 中,使用 RabbitTemplate 向 simple.queue 发布消息,然后运行测试方法
@SpringBootTest
public class PublisherTest {
    @Autowired
    private RabbitTemplate rabbitTemplate;

    @Test
    void send2SimpleQueue() {
        String queueName = "simple.queue";// 队列名称
        String msg = "hello,world";// 消息
        rabbitTemplate.convertAndSend(queueName, msg);// 发布消息
    }
}

3.4、工作模式

需求

  • 一、在控制台中,创建队列 work.queue。
  • 二、在 consumer 中,定义两个方法,同时监听 work.queue。consumer1 每秒处理 50 条消息,consumer2 每秒处理 5 条消息。
  • 三、在 publisher 中,1 秒发布 50 条消息到 work.queue。

示例

  • 一、在控制台中,创建队列 work.queue。
  • 二、在 consumer 中,定义两个方法,同时监听 work.queue。consumer1 每秒处理 50 条消息,consumer2 每秒处理 5 条消息。
// Work
@RabbitListener(queues = "work.queue")
public void listenWorkQueue1(String msg) throws InterruptedException {
    Thread.sleep(20);
    System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "work.queue")
public void listenWorkQueue2(String msg) throws InterruptedException {
    Thread.sleep(200);
    System.out.println("【Consumer2 收到消息】========:" + msg);
}
  • 三、在 publisher 中,1 秒发布 50 条消息到 work.queue。
@Test
void send2WorkQueue() throws InterruptedException {
    String queueName = "work.queue";
    for (int i = 1; i <= 50; i++) {
        String message = "消息" + i;
        rabbitTemplate.convertAndSend("work.queue", message);
        Thread.sleep(20);
    }
}

3.5、限制消息预取

在默认情况下,RabbitMQ 会将消息轮询投递给与队列绑定的每个消费者,并未考虑每个消费者的处理消息能力。例如:Consumer1、Consumer2 都是处理 25 条消息。

限制消息预取,设置 prefetch 为 1,可以确保每次仅投递给每个消费者 1 条消息,处理完成后再从队列中获取。例如:Consumer1、Consumer2 根据处理消息能力,能者多劳。

示例

在 consumer 中,配置 prefetch 的值为 1。

spring:
  rabbitmq: # RabbitMQ 配置
    host: 192.168.112.131 # 主机
    port: 5672 # 监听端口
    virtual-host: /xiaofeng # 虚拟主机
    username: xiaofeng # 用户名
    password: xiaofeng # 密码
    listener:
      simple:
        prefetch: 1 # 限制消息预取,默认没有限制

3.6、发布/订阅模式(Fanout 交换机 )

Fanout ExChange 会将接收到的消息发布到所有与其绑定的队列,也称为广播模式

需求

  • 一、在控制台中,创建队列 fanout.queue1、fanout.queue2,创建交换机 xiaofeng.fanout 并绑定这个两个队列。
  • 二、在 consumer 中,创建两个方法,分别监听 fanout.queue1、fanout.queue2。
  • 三、在 publisher 中,创建测试方法,向 xiaofeng.fanout 发布消息。

示例

  • 一、在控制台中,创建队列 fanout.queue1、fanout.queue2,创建交换机 xiaofeng.fanout 并绑定这个两个队列。

  • 二、在 consumer 中,创建两个方法,分别监听 fanout.queue1、fanout.queue2。

// Fanout
@RabbitListener(queues = "fanout.queue1")
public void listenFanout1(String msg) {
    System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "fanout.queue2")
public void listenFanout2(String msg) {
    System.out.println("【Consumer2 收到消息】:" + msg);
}
  • 三、在 publisher 中,创建测试方法,向 xiaofeng.fanout 发布消息。
@Test
void sendFanout() {
    String exChangeName = "xiaofeng.fanout";
    String msg = "重要通知,明天所有人加班";
    rabbitTemplate.convertAndSend(exChangeName, null, msg);
}

3.7、路由模式(Direct 交换机)

Direct Exchange 会将接收到的消息根据规则发布到指定的队列,也称为定向模式

需求

  • 一、在控制台中,创建队列 direct.queue1、direct.queue2,创建交换机 xiaofeng.direct 并绑定这个两个队列。
  • 二、在 consumer 中,创建两个方法,分别监听 direct.queue1、direct.queue2。
  • 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.direct 发布消息。

示例

  • 一、在控制台中,创建队列 direct.queue1、direct.queue2,创建交换机 xiaofeng.direct 并绑定这个两个队列。
direct.queue1	blue
direct.queue1	red
direct.queue2	red
direct.queue2	yellow
  • 二、在 consumer 中,创建两个方法,分别监听 direct.queue1、direct.queue2。
// Direct
@RabbitListener(queues = "direct.queue1")
public void listenDirect1(String msg) {
    System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "direct.queue2")
public void listenDirect2(String msg) {
    System.out.println("【Consumer2 收到消息】:" + msg);
}
  • 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.direct 发布消息。
@Test
void sendDirect() {
    String exChangeName = "xiaofeng.direct";
    String msg = "西安暴雨预警,注意防范";
    //        rabbitTemplate.convertAndSend(exChangeName, "red", msg);
    //        rabbitTemplate.convertAndSend(exChangeName, "blue", msg);
    rabbitTemplate.convertAndSend(exChangeName, "yellow", msg);
}

3.8、主题模式(Topic 交换机)

Topic Exchange 与 Direct Exchange 类似,区别是 Routing Key 可以使用通配符。

  • “#”:表示零或多个单词。例如:Routing Key 为 china.# 时, 可以匹配到 china.weather、china.news,不能匹配到 japan.weather、japan.news。
  • “*”:表示一个单词。

需求

  • 一、在控制台中,创建队列 topic.queue1、topic.queue2,创建交换机 xiaofeng.topic 并绑定这个两个队列。
  • 二、在 consumer 中,创建两个方法,分别监听 topic.queue1、topic.queue2。
  • 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.topic 发布消息。

示例

  • 一、在控制台中,创建队列 topic.queue1、topic.queue2,创建交换机 xiaofeng.topic 并绑定这个两个队列。
topic.queue1	china.#
topic.queue2	#.news
  • 二、在 consumer 中,创建两个方法,分别监听 topic.queue1、topic.queue2。
// Topic
@RabbitListener(queues = "topic.queue1")
public void listenTopic1(String msg) {
    System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "topic.queue2")
public void listenTopic2(String msg) {
    System.out.println("【Consumer2 收到消息】:" + msg);
}
  • 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.topic 发布消息。
@Test
void sendTopic() {
    String exChangeName = "xiaofeng.topic";
    String msg = "今天天气晴,平均温度 28°C";
    //        rabbitTemplate.convertAndSend(exChangeName, "china.weather", msg);
    //        rabbitTemplate.convertAndSend(exChangeName, "china.news", msg);
    rabbitTemplate.convertAndSend(exChangeName, "japan.news", msg);
}

4、创建队列与交换机

4.1、使用配置类(不太建议)

Spring AMQP 提供工厂类 ,用于创建队列与交换机、及其绑定关系。

  • Queue:用于创建队列,可以使用工厂类 QueueBuilder。
  • Exchange:用于创建交换机,可以使用工厂类 ExchangeBuilder。
  • Binding:用于绑定队列与交换机,可以使用工厂类 BindingBuilder。

需求

使用配置类,实现 Fan 交换机中的示例。

示例

  • 一、在控制台中,删除相关队列、交换机,及其绑定关系。
  • 二、在 Consumer 中,使用配置类。
@Configuration
public class FanoutConfig {
    // 队列
    @Bean
    public Queue fanoutQueue1() {
        return QueueBuilder.durable("fanout.queue1").build();
//        return new Queue("fanout.queue1");
    }
    @Bean
    public Queue fanoutQueue2() {
        return QueueBuilder.durable("fanout.queue2").build();
//        return new Queue("fanout.queue2");
    }

    // 交换机
    @Bean
    public FanoutExchange xiaofengFanout() {
        return ExchangeBuilder.fanoutExchange("xiaofeng.fanout").build();
//        return new FanoutExchange("xiaofeng.fanout");
    }

    // 绑定关系
//    @Bean
//    public Binding binding1() {
//        return BindingBuilder.bind(fanoutQueue1()).to(xiaofengFanout());
//    }
//    @Bean
//    public Binding binding2() {
//        return BindingBuilder.bind(fanoutQueue2()).to(xiaofengFanout());
//    }
    @Bean
    public Binding binding1(Queue fanoutQueue1, FanoutExchange xiaofengFanout) {
        return BindingBuilder.bind(fanoutQueue1).to(xiaofengFanout);
    }
    @Bean
    public Binding binding2(Queue fanoutQueue2, FanoutExchange xiaofengFanout) {
        return BindingBuilder.bind(fanoutQueue2).to(xiaofengFanout);
    }
}

4.2、使用注解(建议)

需求

使用注解,实现 Topic 交换机中的示例。

示例

  • 一、在控制台中,删除相关队列、交换机,及其绑定关系。
  • 二、在 Consumer 中,使用注解。
    // Topic
//    @RabbitListener(queues = "topic.queue1")
//    public void listenTopic1(String msg) {
//        System.out.println("【Consumer1 收到消息】:" + msg);
//    }
//    @RabbitListener(queues = "topic.queue2")
//    public void listenTopic2(String msg) {
//        System.out.println("【Consumer2 收到消息】:" + msg);
//    }

    // Topic
    @RabbitListener(bindings = @QueueBinding(
            value = @Queue(name = "topic.queue1"),// 队列
            exchange = @Exchange(name = "xiaofeng.topic", type = ExchangeTypes.TOPIC),// 交换机
            key = {"china.#"}// Routing key
    ))
    public void listenTopic1(String msg) {
        System.out.println("【Consumer1 收到消息】:" + msg);
    }
    @RabbitListener(bindings = @QueueBinding(
            value = @Queue(name = "topic.queue2"),// 队列
            exchange = @Exchange(name = "xiaofeng.topic", type = ExchangeTypes.TOPIC),// 交换机
            key = {"#.news"}// Routing key
    ))
    public void listenTopic2(String msg) {
        System.out.println("【Consumer2 收到消息】:" + msg);
    }

5、消息转换器

5.1、发布 Object 类型的消息

需求

使用 RabbitTemplate 的 convertAndSend() 方法发布 Map 类型的消息,在控制台中查看消息。

示例

  • 一、在 consumer 中,创建队列 object.queue。
@Component
public class QueueConfig {
    @Bean
    public Queue objectQueue() {
        return QueueBuilder.durable("object.queue").build();
    }
}
  • 二、在 publisher 中,发布 Map 类型的消息。
@Test
void send2SimpleQueue() {
    String queueName = "object.queue";
    Map<String, String> msg = new HashMap<>();
    msg.put("name", "张三");
    msg.put("age", "22");
    rabbitTemplate.convertAndSend(queueName, msg);
}
  • 三、在控制台中查看消息。
content_type: application/x-java-serialized-object
Payload: rO0ABXNyABFqYXZhLnV0aWwuSGFzaE1hcAUH2sHDFmDRAwACRgAKbG9hZEZhY3RvckkACXRocmVzaG9sZHhwP0AAAAAAAAx3CAAAABAAAAACdAAEbmFtZXQA
BuW8oOS4iXQAA2FnZXQAAjIyeA==

这里使用的是 JDK 序列化。

Spring 使用 org.springframework.amqp.support.converter.MessageConverter 处理对象,默认实现是 SimpleMessageConverter,最终使用 ObjectOutputStream 完成序列化。

5.2、消息转换器

示例

  • 一、在控制台中,删除队列 object.queue。
  • 二、在父项目中,导入 JSON 依赖。(jackson-dataformat-xml 包含 jackson-databind,两者均可)
        <!--JSON-->
        <dependency>
            <groupId>com.fasterxml.jackson.core</groupId>
            <artifactId>jackson-databind</artifactId>
        </dependency>
<!--        <dependency>-->
<!--            <groupId>com.fasterxml.jackson.dataformat</groupId>-->
<!--            <artifactId>jackson-dataformat-xml</artifactId>-->
<!--        </dependency>-->
  • 三、在 consumer 与 publisher 的启动类中,注入 Jackson2JsonMessageConverter。(或使用配置类)
@Bean
public MessageConverter messageConverter() {
    return new Jackson2JsonMessageConverter();
}
  • 四、在控制台中查看消息。
content_type: application/json
Payload: {"name":"张三","age":"22"}
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐