IT袋

当前位置:主页 > 经验教程 > 建站编程 >

Spring

Spring Boot项目集成RabbitMQ实战以及坑点讲解(3)

时间:2024-01-30 15:00:18 来源:IT袋 作者:马勇
导读:Spring,交换机列表 队列列表 2. 生产者配置 生产者的消息发送确认主要包含两部分, producter - rabbitmq broker exchange - queue 消息从 producte( 生产者)发送到 rabbitmq

Spring

Spring

交换机列表

Spring

队列列表

2. 生产者配置

Spring

生产者的消息发送确认主要包含两部分,

producter -> rabbitmq broker exchange -> queue

  • 消息从 producte( 生产者)发送到 rabbitmq broker(RabbitMQ 服务器)的交换机中,发送后会触发confirmCallBack回调
  • 消息从 exchange 发送到 queue,投递失败则会调用returnCallBack回调

waynboot-mall 项目的 yml 中关于 RabbitMQ 的相关配置如下,

spring:
  # 配置rabbitMq 服务器
  rabbitmq:
    host: 127.0.0.1
    port: 5672
    username: guest
    password: guest
    # 消息确认配置项
    # 确认消息已发送到交换机(Exchange)
    publisher-confirm-type: correlated
    # 确认消息已发送到队列(Queue)
    publisher-returns: true
    # 虚拟主机名称
    virtual-host: /
publisher-confirm-type 属性

可以看到,我们设置了 publisher-confirm-type 属性为 correlated,表示开启发布确认模式,用来确认消息已发送到交换机,publisher-confirm-type 有三个选项:

  • NONE:禁用发布确认模式,是默认值
  • CORRELATED:发布消息成功到交换器后会触发回confirmCallBack回调方法
  • SIMPLE:经测试有两种效果,其一效果和 CORRELATED 值一样会触发回调方法,其二在发布消息成功后使用 rabbitTemplate 调用 waitForConfirms 或 waitForConfirmsOrDie 方法等待 broker 节点返回发送结果,根据返回结果来判定下一步的逻辑,要注意的点是 waitForConfirmsOrDie 方法如果返回 false 则会关闭 channel,则接下来无法发送消息到 broker。
publisher-returns 属性

在 RabbitMQ 中,消息发送到交换机中也不代表消费者一定能接收到消息,所以我们还需要设置 publisher-returns 为 true 来表示确认交换机中消息已经发送到队列里。true 表示开启失败回调,开启后当消息无法路由到指定队列时会触发ReturnCallback回调。

接着是 RabbitTemplateConfig 的代码,这里面会定义前面提到的 confirmCallBackreturnCallBack 相关代码,

@Slf4j
@Component
public class RabbitTemplateConfig {
    @Bean
    public RabbitTemplate rabbitTemplate(CachingConnectionFactory connectionFactory) {
        RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
        // 设置开启Mandatory,才能触发回调函数,无论消息推送结果怎么样都强制调用回调函数
        rabbitTemplate.setMandatory(true);
        // 交换机收到消息回调
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> log.info("消息发送成功:correlationData({}),ack({}),cause({})", correlationData, ack, cause));
        // 队列收到消息回调,如果失败的话会进行 returnCallback 的回调处理,反之成功就不会回调。
        rabbitTemplate.setReturnsCallback(returned -> {
            log.info("returnCallback:     " + "消息:" + returned.getMessage());
            log.info("returnCallback:     " + "回应码:" + returned.getReplyCode());
            log.info("returnCallback:     " + "回应信息:" + returned.getReplyText());
            log.info("returnCallback:     " + "交换机:" + returned.getExchange());
            log.info("returnCallback:     " + "路由键:" + returned.getRoutingKey());
        });
        return rabbitTemplate;
    }
}

在 RabbitTemplateConfig 类代码里,我们可以设置confirmCallBackreturnCallBack回调函数后,监控生产者发送消息是否被交换机接收、以及交换机是否把消息发送到队列中。

3. 使用 RabbitTemplate 发送消息

在 Spring Boot 项目中,集成了 spring-boot-starter-amqp 依赖后,就可以直接注入 RabbitTemplate 来发送消息。

这里用 waynboot-mall 项目中的异步下单流程举例,代码如下,

@Slf4j
@Service
@AllArgsConstructor
public class OrderServiceImpl extends ServiceImpl<OrderMapper, Order> implements IOrderService {
    private RabbitTemplate rabbitTemplate;
    @Override
    public R asyncSubmit(OrderVO orderVO) {
        OrderDTO orderDTO = new OrderDTO();
        ...
        // 开始异步下单
        String uid = IdUtil.getUid();
        // 1. 创建消息ID,确认机制发送消息时,需要给每个消息设置一个全局唯一 id,以区分不同消息,避免 ack 冲突
        CorrelationData correlationData = new CorrelationData(uid);
        // 2. 创建消息载体 Message ,AMQP 规范中定义的消息承载类,用来在生产者和消费者之前传递消息
        Map<String, Object> map = new HashMap<>();
        map.put("order", orderDTO);
        map.put("notifyUrl", WaynConfig.getMobileUrl() + "/callback/order/submit");
        try {
            Message message = MessageBuilder
                    .withBody(JSON.toJSONString(map).getBytes(Constants.UTF_ENCODING))
                    .setContentType(MessageProperties.CONTENT_TYPE_TEXT_PLAIN)
                    .setDeliveryMode(MessageDeliveryMode.PERSISTENT)
                    .build();
            // 3. 发送消息到 RabbitMQ 服务器,需要指定交换机、路由键、消息载体以及消息ID
            rabbitTemplate.convertAndSend(MQConstants.ORDER_DIRECT_EXCHANGE, MQConstants.ORDER_DIRECT_ROUTING, message, correlationData);
        } catch (UnsupportedEncodingException e) {
            log.error(e.getMessage(), e);
        }
        return R.success().add("actualPrice", actualPrice).add("orderSn", orderSn);
    }
}

waynboot-mall 项目中在使用 rabbitTemplate 发送消息时,按照如下步骤,大家可以参考

  1. 创建消息 ID,确认机制发送消息时,需要给每个消息设置一个全局唯一 id,以区分不同消息,消费者消费时出现 ack 冲突。
  2. 创建消息载体 Message ,AMQP 规范中定义的消息承载类,用来在生产者和消费者之前传递消息。
  3. 发送消息到 RabbitMQ 服务器,需要指定交换机、路由键、消息载体以及消息 ID。

以上就是生产者发送消息时所有相关代码了,接着我们看下消费者处理消息的相关代码。

消费者处理消息

在 waynboot-mall 项目中,还是用订单消息来举例,消费者 yml 配置如下,

1. 消费者配置

在 RabbitMQ 的消息消费环节,需要注意的一点就是,如果需要确保消费者不出现漏消费,则需要开启消费者的手动 ack 模式。

spring:
  rabbitmq:
    host: 127.0.0.1
    port: 5672
    username: guest
    password: guest
    ...
    listener:
      simple:
        # 消息确认方式,其有三种配置方式,分别是none、manual(手动ack) 和auto(自动ack) 默认auto
        acknowledge-mode: manual
        # 一个消费者最多可处理的nack(未确认)消息数量,默认是250
        prefetch: 250
        # 设置消费者数量
        concurrency: 1
acknowledge-mode 属性

在 yml 文件的消费者配置中,acknowledge-mode 属性用于指定消息确认模式,有三种模式:

  1. 手动确认 manual,在该模式下,消费者消费消息后需要根据消费情况给 Broker 返回一个回执,是确认 ack 使 Broker 删除该条已消费的消息,还是失败确认返回 nack,还是拒绝该消息。开启手动确认后,如果消费者接收到消息后还没有返回 ack 就宕机了,这种情况下消息也不会丢失,只有 RabbitMQ 接收到返回 ack 后,消息才会从队列中被删除。
  2. 自动确认 none,rabbitmq 默认消费者正确处理所有请求(不设置时的默认方式)。
  3. 根据请况确认 auto,主要分成以下几种情况:
    • 如果消费者在消费的过程中没有抛出异常,则自动确认。
    • 当消费者消费的过程中抛出 AmqpRejectAndDontRequeueException 异常的时候,则消息会被拒绝,且该消息不会重回队列。
    • 当抛出 ImmediateAcknowledgeAmqpException 异常,消息会被确认。
    • 如果抛出其他的异常,则消息会被拒绝,但是与前两个不同的是,该消息会重回队列,如果此时只有一个消费者监听该队列,那么该消息重回队列后又会推送给该消费者,会造成死循环的情况。
prefetch 属性

消费者配置中,prefetch 属性用于指定消费者每次从队列获取的消息数量。

每个 customer 会在 MQ 预取一些消息放入内存的 LinkedBlockingQueue 中进行消费,这个值越高,消息传递的越快,但非顺序处理消息的风险更高。如果 ack 模式为 none,则忽略。

prefetch 默认值以前是 1,这可能会导致高效使用者的利用率不足。从 spring-amqp 2.0 版开始,默认的 prefetch 值是 250,这将使消费者在大多数常见场景中保持忙碌,从而提高吞吐量。

不过在有些情况下,尤其是处理速度比较慢的大消息,消息可能在内存中大量堆积,消耗大量内存;以及对于一些严格要求顺序的消息,prefetch 的值应当设置为 1。

对于低容量消息和多个消费者的情况(也包括单 listener 容器的 concurrency 配置)希望在多个使用者之间实现更均匀的消息分布,建议在手动 ack 下并设置 prefetch=1。

如果要保证消息的可靠不丢失,当 prefetch 大于 1 时,可能会出现因为服务宕机引起的数据丢失,故建议将 prefetch=1。

concurrency 属性

消费者配置中,concurrency 属性设置的是对每个 listener 在初始化的时候设置的并发消费者的个数。在上面的 yml 配置中,concurrency=1,即每个 Listener 容器将开启一个线程去处理消息。在 2.0 以后的版本中,可以在注解中配置该参数,实例代码如下,

@RabbitListener(queues = MQConstants.ORDER_DIRECT_QUEUE, concurrency = "2")
public void process(Channel channel, Message message) throws IOException {
    String body = new String(message.getBody());
    log.info("OrderPayConsumer 消费者收到消息: {}", body);
    ...
}

2. 使用 RabbitListener 注解消费消息

在 waynboot-mall 项目中,消费者监听队列代码如下,

@Slf4j
@Component
public class OrderPayConsumer {
    @Resource
    private RedisCache redisCache;
    @Resource
    private MobileApi mobileApi;
    @RabbitListener(queues = MQConstants.ORDER_DIRECT_QUEUE)
    public void process(Channel channel, Message message) throws IOException {
        // 1. 转换订单消息
        String body = new String(message.getBody());
        log.info("OrderPayConsumer 消费者收到消息: {}", body);
        // 2. 获取消息ID
        String msgId = message.getMessageProperties().getHeader("spring_returned_message_correlation");
        // 3. 获取发送tag
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        // 4. 消费消息幂等性处理
        if (redisCache.getCacheObject(ORDER_CONSUMER_MAP.getKey()) != null) {
            // redis中包含该 key,说明该消息已经被消费过
            log.error("msgId: {},消息已经被消费", msgId);
            channel.basicAck(deliveryTag, false);// 确认消息已消费
            return;
        }
        try {
            // 5. 下单处理
            mobileApi.submitOrder(body);
            // 6. 手动ack,消息成功确认
            channel.basicAck(deliveryTag, false);
            // 7. 设置消息已被消费标识
            redisCache.setCacheObject(ORDER_CONSUMER_MAP.getKey(), msgId, ORDER_CONSUMER_MAP.getExpireSecond());
        } catch (Exception e) {
            channel.basicNack(deliveryTag, false, false);
            log.error(e.getMessage(), e);
        }
    }
}

waynboot-mall 项目中在使用 RabbitListener 注解消费消息时,按照如下步骤,大家可以参考

  1. 将 message 参数转换成订单消息。
  2. 从 message 参数中获取消息唯一 msgId。
  3. 从 message 参数中获取消息发送 tag。
  4. 幂等性处理,根据第二步获取的 msgId ,消费消息时需要先判断 msgId 是否已经被处理。
  5. 调用 mobile-api 服务,进行下单逻辑处理,在 mobileApi.submitOrder(body) 方法中使用 Spring-Retry 的 @Retryable 注解,进行自动重试。
  6. 手动 ack,basicAck(long deliveryTag, boolean multiple)。basicAck 方法表示成功确认,使用此方法后,消息就会被 rabbitmq 服务器删除。

其中参数 long deliveryTag 为消息的唯一序号也就是第三步获取的发送 tag,第二个 boolean multiple 参数表示是否一次消费多条消息,false 表示只确认该序列号对应的消息,true 则表示确认该序列号对应的消息以及比该序列号小的所有消息,比如我先发送 2 条消息,他们的序列号分别为 2,3,并且他们都没有被确认,还留在队列中,那么如果当前消息序列号为 4,那么当 multiple 为 true,则序列号为 2、3 的消息也会被一同确认。

  1. 幂等性处理,消息已经被成功消费后,根据第二步获取的 msgId 设置幂等标识。

总结一下

这篇文章给大家讲解了在 Spring Boot 项目中如何集成消息队列 RabbitMQ 用于业务逻辑解耦,有架构介绍、应用场景、坑点解析、代码实战 4 个部分,能带领大家比较全面的了解一波 RabbitMQ。

大家在自己的项目中如果需要引入 RabbitMQ 时,都可以参考本文的代码实战配置,帮助大家快速集成、避免踩坑。

以上IT袋网网介绍的Spring 跟 Boot项目集成RabbitMQ实战以及坑点讲解的具体介绍,供网友们借鉴参考。

相关阅读

  • 服务器查询网站入口在哪 查询网站服务器的办法

    服务器查询网站入口在哪 查询网站服务器的办法

    小编为你讲解服务器查询网站入口在哪和查询网站服务器的办法的介绍,接下来分享详细内容。 随着科学技术的飞速发展,越来越多的用户接触到服务器,对服务器有了一定的了解。服务器配

  • 如何通过 Nebula Exchange 导入数据

    如何通过 Nebula Exchange 导入数据

    小编带来的是及Nebula的相关话题,接下来IT袋网小编为大家介绍。 Nebula Exchange 是一款 Apache Spark 应用,用于在分布式环境中将集群中的数据批量迁移到 Nebula Graph 中,能支持多种不同格式(CS

  • 内网ftp服务器搭建教程 简述ftp服务器搭建过程

    内网ftp服务器搭建教程 简述ftp服务器搭建过程

    今天带来的IT技巧小经验内网ftp服务器搭建教程和简述ftp服务器搭建过程的方法内容,接下来一起来看看吧。 1 文件传输协议 一般来讲,人们将计算机联网的首要目的就是获取资料,而文件传

  • 香港BGP服务器丢包率高与什么有关? 服务器丢包的原因有哪些

    香港BGP服务器丢包率高与什么有关? 服务器丢包的原因有哪些

    相对于大多数人香港BGP服务器丢包率高与什么有关的内容,接下来一起来看看吧。 当前,BGP线的香港机房比较普及,假设用户所在机房不是BGP机房,而是联通或电信等其他机房,这时当用户通