返回文章列表
高并发项目
MQRabbitMQ异步通信SpringAMQP

07MQ基础

1.同步通讯与异步通讯

在微服务架构中,服务之间的通讯方式可以分为同步通讯与异步通讯两类:

  • 同步通讯:类似打视频电话,双方交互实时进行,同一时刻只能与一方通讯。
  • 异步通讯:类似发微信聊天,双方交互并非实时,不需要立刻回应,可以同时与多人交互。

同步通讯与异步通讯

以登录场景为例,用户登录后除了校验用户信息,往往还需要风控服务判断是否存在登录风险,并短信通知用户。如果全部采用同步调用,登录链路将变得冗长且耦合度高;而引入 MQ 后,风控判断与短信通知可改为异步消息通知,主链路仅保留核心校验。

登录风控场景中的同步与异步

1.1同步调用及其问题

我们以黑马商城的余额支付为例,业务流程依次为:扣减余额 → 更新支付状态 → 更新订单状态 → 短信通知用户 → 增加用户积分。

余额支付的同步调用链路

同步调用的优势是时效性强,必须等待到结果后才返回。但同时也带来三个问题:

  • 拓展性差:新增功能(如短信通知、积分添加)需要不断追加到原有业务模块,最终变得臃肿。
  • 性能下降:同步通讯逐步执行全流程(支付、用户、交易、通知、积分……),耗时累加,性能下降。
  • 级联失败:分布式事务下,一处失败可能触发全链路回滚,而像通知失败这类业务本不该回滚用户支付。

1.2异步调用及其优势

异步调用通常基于消息通知的方式,包含三个角色:

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

异步调用的三个角色

支付服务不再同步调用业务关联度低的服务,而是发送消息通知到 Broker,具备下列优势:

  • 解除耦合,拓展性强。
  • 无需等待,性能好。
  • 故障隔离,下游服务故障不影响上游业务。
  • 缓存消息,流量削峰填谷。

异步调用改造支付链路

异步调用的不足之处也需注意:不能立即得到调用结果,时效性差;不确定下游业务执行是否成功;业务安全依赖于 Broker 的可靠性。

2.MQ技术选型

MQ(MessageQueue),中文是消息队列,字面来看就是存放消息的队列,也就是异步调用中的 Broker。常见 MQ 对比如下:

RabbitMQ ActiveMQ RocketMQ Kafka
公司/社区 Rabbit Apache 阿里 Apache
开发语言 Erlang Java Java Scala&Java
协议支持 AMQP,XMPP,SMTP,STOMP OpenWire,STOMP,REST,XMPP,AMQP 自定义协议 自定义协议
可用性 高 一般 高 高
单机吞吐量 一般 差 高 非常高
消息延迟 微秒级 毫秒级 毫秒级 毫秒以内
消息可靠性 高 一般 高 一般

综合可靠性、协议支持与生态成熟度,本项目采用 RabbitMQ。

3.RabbitMQ架构与核心概念

RabbitMQ 的整体架构及核心概念如下:

RabbitMQ 整体架构

其中包含几个核心概念:

  • virtual-host:虚拟主机,起到数据隔离的作用。每个虚拟主机相互独立,有各自的 exchange、queue。
  • publisher:消息发送者。
  • consumer:消息的消费者。
  • queue:队列,存储消息。
  • exchange:交换机,负责路由消息。

3.1交换机的三种类型

交换机的作用主要是接收发送者发送的消息,并将消息路由到与其绑定的队列。常见交换机类型有以下三种:

  • Fanout:广播,会将接收到的消息路由到每一个跟其绑定的 queue。
  • Direct:定向,每一个 Queue 都与 Exchange 设置一个 BindingKey,发布者发送消息时指定 RoutingKey,Exchange 将消息路由到 BindingKey 与消息 RoutingKey 一致的队列。
  • Topic:话题,与 Direct 类似,区别在于 routingKey 可以是多个单词的列表,并且以 . 分割。BindingKey 可以使用通配符:
    • #:代指 0 个或多个单词。
    • *:代指一个单词。

4.SpringAMQP集成

Spring AMQP 是基于 AMQP 协议定义的一套 API 规范,提供了模板来发送和接收消息。它包含两部分:spring-amqp 是基础抽象,spring-rabbit 是底层的默认实现。

4.1快速入门

需求:利用控制台创建队列 simple.queue;在 publisher 服务中利用 SpringAMQP 直接向 simple.queue 发送消息;在 consumer 服务中利用 SpringAMQP 编写消费者监听 simple.queue 队列。

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

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

配置 RabbitMQ 服务端信息:在每个微服务中引入 MQ 服务端信息,这样微服务才能连接到 RabbitMQ:

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

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

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

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

@Slf4j
@Component
public class SpringRabbitListener {
 
    @RabbitListener(queues = "simple.queue")
    public void listenSimpleQueueMessage(String msg) throws InterruptedException {
        log.info("spring 消费者接收到消息:【" + msg + "】");
    }
}

4.2Work Queues

Work queues(任务模型)就是让多个消费者绑定到一个队列,共同消费队列中的消息。基本思路如下:

  • 在 RabbitMQ 控制台创建一个队列 work.queue。
  • 在 publisher 服务中定义测试方法,发送 50 条消息到 work.queue。
  • 在 consumer 服务中定义两个消息监听者,都监听 work.queue 队列。
  • 消费者 1 每秒处理 40 条消息,消费者 2 每秒处理 5 条消息。

默认情况下,RabbitMQ 会将消息依次轮询投递给绑定在队列上的每一个消费者,但这并没有考虑到消费者是否已经处理完消息,可能出现消息堆积。因此需要修改 application.yml,设置 preFetch 值为 1,确保同一时刻最多投递给消费者 1 条消息:

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

Work 模型的使用要点:

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

4.3Fanout交换机

Fanout Exchange 会将接收到的消息路由到每一个跟其绑定的 queue,所以也叫广播模式。利用 SpringAMQP 演示 FanoutExchange 的实现思路:

  • 在 RabbitMQ 控制台中声明队列 fanout.queue1 和 fanout.queue2。
  • 在 RabbitMQ 控制台中声明交换机 hmall.fanout,将两个队列与其绑定。
  • 在 consumer 服务中编写两个消费者方法,分别监听 fanout.queue1 和 fanout.queue2。
  • 在 publisher 中编写测试方法,向 hmall.fanout 发送消息。

发送消息到交换机的 API:

@Test
public void testFanoutExchange() {
    // 队列名称
    String exchangeName = "hmall.fanout";
    // 消息
    String message = "hello, everyone!";
    // 发送消息,参数分别是:交换机名称、RoutingKey(暂时为空)、消息
    rabbitTemplate.convertAndSend(exchangeName, "", message);
}

4.4Direct交换机

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

  • 每一个 Queue 都与 Exchange 设置一个 BindingKey。
  • 发布者发送消息时,指定消息的 RoutingKey。
  • Exchange 将消息路由到 BindingKey 与消息 RoutingKey 一致的队列。

Fanout 与 Direct 的差异:Fanout 交换机将消息路由给每一个与之绑定的队列;Direct 交换机根据 RoutingKey 判断路由给哪个队列;如果多个队列具有相同 RoutingKey,则与 Fanout 功能类似。

4.5Topic交换机

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

  • #:代指 0 个或多个单词。
  • *:代指一个单词。

例如:china.news 代表中国的新闻消息;china.weather 代表中国的天气消息;japan.news 代表日本新闻;japan.weather 代表日本的天气消息。可用 china.# 匹配前两个、#.news 匹配第 1、3 个。

4.6声明队列和交换机

声明队列和交换机

SpringAMQP 提供了几个类用来声明队列、交换机及其绑定关系:

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

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

@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个队列,并完成绑定
}

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

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

4.7消息转换器

需求:测试利用 SpringAMQP 发送对象类型的消息。声明一个队列 object.queue,编写单元测试,向队列中直接发送一条消息,消息类型为 Map。

// 准备消息
Map<String,Object> msg = new HashMap<>();
msg.put("name", "Jack");
msg.put("age", 21);

Spring 对消息对象的处理是由 org.springframework.amqp.support.converter.MessageConverter 来处理的,默认实现是 SimpleMessageConverter,基于 JDK 的 ObjectOutputStream 完成序列化。存在下列问题:

  • JDK 的序列化有安全风险。
  • JDK 序列化的消息太大。
  • JDK 序列化的消息可读性差。

建议采用 JSON 序列化代替默认的 JDK 序列化。在 publisher 和 consumer 中都要引入 jackson 依赖:

<dependency>
    <groupId>com.fasterxml.jackson.core</groupId>
    <artifactId>jackson-databind</artifactId>
</dependency>

在 publisher 和 consumer 中都要配置 MessageConverter:

@Bean
public MessageConverter messageConverter(){
    return new Jackson2JsonMessageConverter();
}

5.业务改造

需求:改造余额支付功能,不再同步调用交易服务的 OpenFeign 接口,而是采用异步 MQ 通知交易服务更新订单状态。

  • 支付服务作为消息发送方,向 pay.topic 交换机发送 pay.success 路由的消息。
  • 交易服务监听 mark.order.pay.queue,更新订单状态。
  • 通知服务监听 pay.notify.queue,发送支付成功通知。

至此,原有的同步链路被拆分为支付主链路 + 异步通知/订单更新,性能和拓展性都得到提升。