最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

SpringBoot+RabbitMQ實現(xiàn)消息可靠投遞+防重復(fù)消費(可直接落地)

 更新時間:2026年04月20日 08:30:44   作者:喝汽水的貓^  
本文主要介紹了Spring Boot + RabbitMQ 實戰(zhàn):消息可靠投遞+防重復(fù)消費,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

提示:文章寫完后,目錄可以自動生成,如何生成可參考右邊的幫助文檔

前言

Spring Boot + RabbitMQ 實戰(zhàn):消息可靠投遞+防重復(fù)消費(可直接落地)
在高并發(fā)業(yè)務(wù)場景中,RabbitMQ 作為消息中間件,核心作用是削峰填谷、解耦服務(wù),但最關(guān)鍵的兩個問題的是:消息不丟失、不重復(fù)消費。

基于 Spring Boot 整合 RabbitMQ,提供一套可直接復(fù)制、生產(chǎn)環(huán)境可用的實戰(zhàn)代碼,涵蓋「生產(chǎn)者 Confirm 確認、消息持久化、消費者手動 ACK、Redis 冪等防重」全流程,避開所有常見坑,新手也能快速落地。

適用場景:訂單異步創(chuàng)建、短信/通知推送、物流狀態(tài)同步等所有需要保證消息可靠性的業(yè)務(wù),尤其適配高并發(fā)下單、秒殺等場景。

一、核心需求與技術(shù)選型

1. 核心需求(必滿足)

  • 生產(chǎn)者:消息必須送達 RabbitMQ Broker,失敗可重試,杜絕生產(chǎn)者丟消息
  • 消息本身:Broker 重啟、服務(wù)宕機后,消息不丟失
  • 消費者:業(yè)務(wù)處理成功后再確認消息,異??芍匦氯腙牐沤^消費端丟消息
  • 冪等性:避免因網(wǎng)絡(luò)重試、消息重入隊導致的重復(fù)消費(比如重復(fù)創(chuàng)建訂單、重復(fù)扣庫存)

2. 技術(shù)選型

  • 框架:Spring Boot 2.x(兼容 3.x,只需微調(diào)依賴)
  • 消息中間件:RabbitMQ 3.9+
  • 冪等校驗:Redis(高效判重)+ 數(shù)據(jù)庫唯一索引(兜底)
  • 核心依賴:spring-boot-starter-amqp、spring-boot-starter-data-redis

二、環(huán)境配置(application.yml)

核心配置:開啟生產(chǎn)者 Confirm 機制、Return 機制,消費者手動 ACK,同時配置限流防止數(shù)據(jù)庫被沖垮,注釋清晰可直接復(fù)制。

spring:
  # RabbitMQ 核心配置
  rabbitmq:
    host: 127.0.0.1  # 本地環(huán)境,生產(chǎn)環(huán)境替換為服務(wù)器地址
    port: 5672       # RabbitMQ 默認端口
    username: guest  # 默認用戶名,生產(chǎn)環(huán)境需修改為自定義賬號
    password: guest  # 默認密碼,生產(chǎn)環(huán)境需修改
    virtual-host: /  # 虛擬主機,默認即可
    connection-timeout: 10000  # 連接超時時間,避免無限等待
    # 1. 生產(chǎn)者確認機制:確保消息到達 Broker
    publisher-confirm-type: correlated  # correlated:異步回調(diào),獲取確認結(jié)果
    # 2. 消息回退機制:消息無法路由時返回生產(chǎn)者,避免消息丟失
    publisher-returns: true
    # 3. 消費者配置
    listener:
      simple:
        acknowledge-mode: manual  # 手動 ACK(關(guān)鍵!避免自動確認丟消息)
        concurrency: 5            # 消費者核心并發(fā)數(shù)
        max-concurrency: 10       # 消費者最大并發(fā)數(shù)
        prefetch: 10              # 限流:每次只獲取10條消息,防止消費過快沖垮數(shù)據(jù)庫

三、核心代碼實現(xiàn)(全可復(fù)制)

1. 依賴導入(pom.xml)

無需額外配置,導入 Spring Boot 整合 RabbitMQ 和 Redis 的 starter 即可。

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!--  lombok 簡化代碼,可選但推薦 -->
<dependency>
    <groupId>org.projectlombok</groupId>
    <artifactId>lombok</artifactId>
    <optional>true</optional>
</dependency>

2. 實體類(OrderCreateDTO)

訂單消息傳輸?shù)膶嶓w類,根據(jù)自身業(yè)務(wù)調(diào)整字段,需實現(xiàn) Serializable 接口(RabbitMQ 消息傳輸要求)。

import lombok.Data;
import java.io.Serializable;
/**
 * 訂單創(chuàng)建消息DTO
 */
@Data
public class OrderCreateDTO implements Serializable {
    // 訂單唯一編號(用于業(yè)務(wù)冪等)
    private String orderSn;
    // 用戶ID
    private Long userId;
    // 訂單金額
    private BigDecimal orderAmount;
    // 商品ID(多個可改為List)
    private Long productId;
    // 購買數(shù)量
    private Integer quantity;
}

3. 生產(chǎn)者:消息可靠投遞(Confirm + 持久化)

核心邏輯:

  • 生成全局唯一 msgId(用于后續(xù)冪等判重)
  • 設(shè)置消息持久化(Broker 重啟后消息不丟失)
  • 開啟 Confirm 回調(diào),監(jiān)聽消息是否成功送達 Broker,失敗可重試/落庫
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import java.util.UUID;
/**
 * 訂單消息生產(chǎn)者(可靠投遞)
 */
@Component
@Slf4j
@RequiredArgsConstructor
public class OrderProducer {
    // 注入RabbitTemplate,用于發(fā)送消息
    private final RabbitTemplate rabbitTemplate;
    // 交換機名稱(需與消費者隊列綁定)
    private static final String ORDER_EXCHANGE = "order.exchange";
    // 路由鍵(需與隊列綁定,確保消息能路由到指定隊列)
    private static final String ORDER_CREATE_ROUTING_KEY = "order.create";
    /**
     * 發(fā)送訂單創(chuàng)建消息
     * @param dto 訂單創(chuàng)建DTO
     */
    public void sendOrderMsg(OrderCreateDTO dto) {
        // 1. 生成全局唯一消息ID,用于冪等判重(UUID保證唯一性)
        String msgId = UUID.randomUUID().toString().replace("-", "");
        // 2. 關(guān)聯(lián)消息ID,用于Confirm回調(diào)獲取消息標識
        CorrelationData correlationData = new CorrelationData(msgId);
        // 3. 發(fā)送消息(設(shè)置持久化 + 攜帶消息ID)
        rabbitTemplate.convertAndSend(
                ORDER_EXCHANGE,          // 交換機
                ORDER_CREATE_ROUTING_KEY,// 路由鍵
                dto,                     // 消息內(nèi)容
                message -> {
                    // 設(shè)置消息持久化(MessageDeliveryMode.PERSISTENT:持久化)
                    message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                    // 將msgId存入消息屬性,供消費者獲取
                    message.getMessageProperties().setMessageId(msgId);
                    return message;
                },
                correlationData          // 關(guān)聯(lián)消息ID,用于Confirm回調(diào)
        );
        // 4. 生產(chǎn)者Confirm回調(diào):監(jiān)聽消息是否成功送達Broker
        rabbitTemplate.setConfirmCallback((correlation, ack, cause) -> {
            // 獲取回調(diào)的消息ID
            String msgIdCallback = correlation.getId();
            if (ack) {
                // ack為true:消息成功送達Broker
                log.info("消息發(fā)送成功,msgId:{}", msgIdCallback);
            } else {
                // ack為false:消息發(fā)送失敗
                log.error("消息發(fā)送失敗,msgId:{},失敗原因:{}", msgIdCallback, cause);
                // 失敗處理:可重試發(fā)送(建議最多3次),或入庫定時重發(fā)(避免消息丟失)
                retrySendMsg(dto, msgIdCallback);
            }
        });
    }
    /**
     * 消息發(fā)送失敗重試(簡單重試邏輯,可根據(jù)業(yè)務(wù)優(yōu)化)
     */
    private void retrySendMsg(OrderCreateDTO dto, String msgId) {
        int retryCount = 3; // 重試3次
        for (int i = 0; i < retryCount; i++) {
            try {
                Thread.sleep(1000 * (i + 1)); // 指數(shù)退避重試(1s、2s、3s)
                CorrelationData correlationData = new CorrelationData(msgId);
                rabbitTemplate.convertAndSend(
                        ORDER_EXCHANGE,
                        ORDER_CREATE_ROUTING_KEY,
                        dto,
                        message -> {
                            message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                            message.getMessageProperties().setMessageId(msgId);
                            return message;
                        },
                        correlationData
                );
                log.info("消息重試發(fā)送成功,msgId:{},重試次數(shù):{}", msgId, i + 1);
                return;
            } catch (Exception e) {
                log.error("消息重試發(fā)送失敗,msgId:{},重試次數(shù):{}", msgId, i + 1, e);
                if (i == retryCount - 1) {
                    // 重試3次仍失敗,入庫定時重發(fā)(此處省略入庫邏輯,可結(jié)合定時任務(wù)實現(xiàn))
                    log.error("消息重試3次失敗,msgId:{},已入庫待定時重發(fā)", msgId);
                }
            }
        }
    }
}

4. 消費者:手動 ACK + Redis 冪等防重

核心邏輯:

  • 手動 ACK:業(yè)務(wù)處理成功后,調(diào)用 basicAck 確認消息;異常則調(diào)用 basicNack 重新入隊
  • Redis 冪等:用 setIfAbsent 存儲 msgId,已消費則直接 ACK,避免重復(fù)消費
  • 業(yè)務(wù)兜底:結(jié)合數(shù)據(jù)庫唯一索引,防止 Redis 掛了導致的重復(fù)消費
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import com.rabbitmq.client.Channel;
import java.util.concurrent.TimeUnit;
/**
 * 訂單消息消費者(手動ACK + 冪等防重)
 */
@Component
@Slf4j
@RequiredArgsConstructor
public class OrderConsumer {
    // 注入Redis模板,用于冪等判重
    private final StringRedisTemplate redisTemplate;
    // 注入訂單服務(wù),處理核心業(yè)務(wù)邏輯
    private final OrderService orderService;
    // 隊列名稱(需與交換機、路由鍵綁定)
    private static final String ORDER_CREATE_QUEUE = "order.create.queue";
    // Redis 冪等鍵前綴(區(qū)分不同業(yè)務(wù)的消息)
    private static final String MQ_CONSUMED_KEY_PREFIX = "mq:consumed:order:";
    /**
     * 消費訂單創(chuàng)建消息
     * @param dto 消息內(nèi)容(自動反序列化)
     * @param message 消息對象,用于獲取msgId
     * @param channel 信道對象,用于手動ACK/NACK
     */
    @RabbitListener(queues = ORDER_CREATE_QUEUE) // 監(jiān)聽指定隊列
    public void consumeOrderMsg(OrderCreateDTO dto, Message message, Channel channel) throws Exception {
        // 1. 獲取消息ID(生產(chǎn)者存入的msgId)
        String msgId = message.getMessageProperties().getMessageId();
        // 2. 獲取消息投遞標簽(用于手動ACK/NACK)
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        // 3. Redis 冪等判重:setIfAbsent 原子操作,避免并發(fā)重復(fù)消費
        String redisKey = MQ_CONSUMED_KEY_PREFIX + msgId;
        // 存入Redis,有效期24小時(根據(jù)業(yè)務(wù)調(diào)整,確保消息消費完成后不會被重復(fù)判斷)
        Boolean consumeFlag = redisTemplate.opsForValue().setIfAbsent(redisKey, "1", 24, TimeUnit.HOURS);
        // 4. 已消費過:直接ACK,避免重復(fù)處理
        if (consumeFlag == null || !consumeFlag) {
            log.warn("消息已消費,無需重復(fù)處理,msgId:{}", msgId);
            // 手動ACK:deliveryTag為當前消息標簽,false表示不批量確認
            channel.basicAck(deliveryTag, false);
            return;
        }
        try {
            // 5. 處理核心業(yè)務(wù)邏輯:創(chuàng)建訂單、扣減庫存等
            orderService.createOrder(dto);
            // 6. 業(yè)務(wù)處理成功:手動ACK,通知RabbitMQ刪除消息
            channel.basicAck(deliveryTag, false);
            log.info("消息消費成功,msgId:{},訂單號:{}", msgId, dto.getOrderSn());
        } catch (Exception e) {
            log.error("消息消費異常,msgId:{},訂單號:{}", msgId, dto.getOrderSn(), e);
            // 7. 異常處理:根據(jù)異常類型決定是否重入隊
            // 可重試異常(如網(wǎng)絡(luò)波動、數(shù)據(jù)庫臨時不可用):重入隊(third參數(shù)為true)
            // 不可重試異常(如業(yè)務(wù)校驗失敗、參數(shù)錯誤):直接拒絕,不重入隊(third參數(shù)為false)
            if (isRetryException(e)) {
                log.info("消息消費異常(可重試),將重入隊,msgId:{}", msgId);
                channel.basicNack(deliveryTag, false, true);
            } else {
                log.info("消息消費異常(不可重試),直接拒絕,msgId:{}", msgId);
                // 不可重試異常:拒絕消息,不重入隊(可結(jié)合死信隊列處理)
                channel.basicReject(deliveryTag, false);
            }
        }
    }
    /**
     * 判斷是否為可重試異常(根據(jù)自身業(yè)務(wù)調(diào)整)
     */
    private boolean isRetryException(Exception e) {
        // 示例:網(wǎng)絡(luò)異常、數(shù)據(jù)庫異??芍卦嚕瑯I(yè)務(wù)異常不可重試
        return e instanceof RuntimeException
                && (e.getMessage().contains("網(wǎng)絡(luò)") || e.getMessage().contains("數(shù)據(jù)庫"));
    }
}

5. 業(yè)務(wù)層:Redis 冪等 + 數(shù)據(jù)庫兜底

**核心:**即使 Redis 掛了,通過數(shù)據(jù)庫唯一索引(order_sn)兜底,確保不會重復(fù)創(chuàng)建訂單、重復(fù)扣庫存。

import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
/**
 * 訂單服務(wù)(核心業(yè)務(wù)邏輯)
 */
@Service
@Slf4j
@RequiredArgsConstructor
public class OrderService {
    private final OrderMapper orderMapper;
    private final ProductMapper productMapper;
    /**
     * 創(chuàng)建訂單(帶冪等校驗)
     * @param dto 訂單創(chuàng)建DTO
     */
    @Transactional(rollbackFor = Exception.class) // 事務(wù)管理,異?;貪L
    public void createOrder(OrderCreateDTO dto) {
        String orderSn = dto.getOrderSn();
        // 1. 數(shù)據(jù)庫冪等兜底:查詢訂單是否已存在(訂單表order_sn字段建唯一索引)
        Integer orderCount = orderMapper.countByOrderSn(orderSn);
        if (orderCount > 0) {
            log.warn("訂單已存在,無需重復(fù)創(chuàng)建,訂單號:{}", orderSn);
            return;
        }
        // 2. 核心業(yè)務(wù)邏輯:創(chuàng)建訂單、扣減庫存(根據(jù)自身業(yè)務(wù)實現(xiàn))
        // ① 扣減商品庫存(需加鎖,避免超賣,此處省略分布式鎖邏輯)
        Product product = productMapper.selectById(dto.getProductId());
        if (product == null || product.getStock() < dto.getQuantity()) {
            throw new RuntimeException("商品不存在或庫存不足,訂單號:" + orderSn);
        }
        product.setStock(product.getStock() - dto.getQuantity());
        productMapper.updateById(product);
        // ② 插入訂單記錄
        Order order = new Order();
        order.setOrderSn(orderSn);
        order.setUserId(dto.getUserId());
        order.setOrderAmount(dto.getOrderAmount());
        order.setProductId(dto.getProductId());
        order.setQuantity(dto.getQuantity());
        order.setStatus(0); // 0:待支付
        orderMapper.insert(order);
        log.info("訂單創(chuàng)建成功,訂單號:{}", orderSn);
    }
}

6. 交換機、隊列綁定(可選,兩種方式)

方式1:代碼綁定(推薦,部署時自動創(chuàng)建,無需手動操作)

import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
 * RabbitMQ 交換機、隊列綁定配置
 */
@Configuration
public class RabbitMQConfig {
    // 交換機名稱(與生產(chǎn)者、消費者一致)
    public static final String ORDER_EXCHANGE = "order.exchange";
    // 隊列名稱(與消費者一致)
    public static final String ORDER_CREATE_QUEUE = "order.create.queue";
    // 路由鍵(與生產(chǎn)者一致)
    public static final String ORDER_CREATE_ROUTING_KEY = "order.create";
    // 1. 聲明交換機(direct類型,持久化)
    @Bean
    public DirectExchange orderExchange() {
        // durable=true:交換機持久化,重啟后不丟失
        return new DirectExchange(ORDER_EXCHANGE, true, false);
    }
    // 2. 聲明隊列(持久化)
    @Bean
    public Queue orderCreateQueue() {
        // durable=true:隊列持久化;exclusive=false:不排他;autoDelete=false:不自動刪除
        return new Queue(ORDER_CREATE_QUEUE, true, false, false);
    }
    // 3. 綁定交換機、隊列、路由鍵
    @Bean
    public Binding orderCreateBinding() {
        return BindingBuilder.bind(orderCreateQueue())
                .to(orderExchange())
                .with(ORDER_CREATE_ROUTING_KEY);
    }
}

方式2:RabbitMQ 管理界面手動綁定(適合測試環(huán)境,生產(chǎn)環(huán)境推薦代碼綁定)

  • 登錄 RabbitMQ 管理界面(默認地址:http://localhost:15672)
  • 創(chuàng)建交換機:類型 direct,名稱 order.exchange,勾選 durable
  • 創(chuàng)建隊列:名稱 order.create.queue,勾選 durable
  • 綁定:交換機 → 隊列,路由鍵填寫 order.create

四、關(guān)鍵知識點與避坑點(實戰(zhàn)重點)

消息不丟失的3個關(guān)鍵

  • 生產(chǎn)者:開啟 Confirm 機制,確保消息到達 Broker,失敗重試/落庫
  • 消息:設(shè)置持久化(MessageDeliveryMode.PERSISTENT),交換機、隊列也需持久化
  • 消費者:手動 ACK,業(yè)務(wù)成功后再確認,異常合理處理(重入隊/死信)

防重復(fù)消費的2層保障

  • 第一層:Redis setIfAbsent 原子操作(高效判重,適合高并發(fā))
  • 第二層:數(shù)據(jù)庫唯一索引(兜底,防止 Redis 掛了導致的重復(fù)消費)

常見坑及解決方案

  • 坑1:消息自動 ACK → 解決方案:配置 acknowledge-mode: manual,手動 ACK
  • 坑2:消息未持久化 → 解決方案:設(shè)置消息、交換機、隊列均為持久化(durable=true)
  • 坑3:Redis 掛了導致重復(fù)消費 → 解決方案:數(shù)據(jù)庫唯一索引兜底
  • 坑4:消費者并發(fā)過高沖垮數(shù)據(jù)庫 → 解決方案:配置 prefetch 限流,控制每次獲取的消息數(shù)
  • 坑5:消息發(fā)送失敗后不重試 → 解決方案:實現(xiàn) Confirm 回調(diào),失敗后指數(shù)退避重試,重試失敗入庫定時重發(fā)

五、測試驗證(快速驗證可用性)

  • 啟動 RabbitMQ 服務(wù)(本地可通過 Docker 快速部署)
  • 啟動 Spring Boot 項目,自動創(chuàng)建交換機、隊列并綁定
  • 編寫測試類,調(diào)用 OrderProducer 的 sendOrderMsg 方法發(fā)送消息
  • 查看日志:消息發(fā)送成功 → 消費成功 → 訂單創(chuàng)建成功
  • 測試異常場景:關(guān)閉數(shù)據(jù)庫,發(fā)送消息,查看是否重入隊;恢復(fù)數(shù)據(jù)庫后,查看是否正常消費
  • 測試重復(fù)消費:手動將消息重新入隊,查看是否會重復(fù)創(chuàng)建訂單(應(yīng)提示“訂單已存在”)

六、總結(jié)

本文提供的代碼的是生產(chǎn)環(huán)境真實落地版本,涵蓋了 RabbitMQ 消息可靠投遞和防重復(fù)消費的全流程,無需修改核心邏輯,只需根據(jù)自身業(yè)務(wù)調(diào)整實體類和業(yè)務(wù)方法,即可快速集成到項目中。

核心思路:生產(chǎn)者靠 Confirm 保送達,消息靠持久化保存活,消費者靠手動 ACK 保消費,冪等靠 Redis+數(shù)據(jù)庫保唯一,四者結(jié)合,徹底解決 RabbitMQ 消息丟失和重復(fù)消費的痛點。

后續(xù)可優(yōu)化方向:消息重試機制(結(jié)合定時任務(wù))、死信隊列(處理不可重試異常消息)、分布式鎖(防止庫存超賣),可根據(jù)業(yè)務(wù)復(fù)雜度逐步迭代

到此這篇關(guān)于SpringBoot+RabbitMQ實現(xiàn)消息可靠投遞+防重復(fù)消費(可直接落地)的文章就介紹到這了,更多相關(guān)SpringBoot RabbitMQ 消息可靠投遞+防重復(fù)消費內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

最新評論

开江县| 时尚| 汾阳市| 乡宁县| 临西县| 泰顺县| 托克托县| 翁牛特旗| 武义县| 黎城县| 辰溪县| 樟树市| 乃东县| 甘孜县| 临海市| 安新县| 渝中区| 苍梧县| 宁陵县| 荔浦县| 萨嘎县| 龙川县| 萍乡市| 鹤山市| 云南省| 游戏| 深圳市| 阳城县| 靖宇县| 罗甸县| 汝州市| 洛川县| 庆云县| 丽江市| 阳春市| 江川县| 象州县| 屏南县| 饶河县| 星座| 米易县|