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

SpringBoot4.0整合RabbitMQ死信隊列詳解

 更新時間:2026年02月12日 11:01:33   作者:小壞說Java  
本文主要介紹了SpringBoot4.0整合RabbitMQ死信隊列詳解,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

為啥那么講解死信隊列,因為好多人不會使用,不知道什么場景下使用,此案例是我在公司實現(xiàn)的一種方式,讓大家都可以學習到

一、死信隊列的好處

1.提高系統(tǒng)可靠性

  • 避免消息丟失,確保處理失敗的消息有備份
  • 防止因消息處理異常導致的消息無限重試

2.異常消息管理

  • 將異常消息與正常消息分離
  • 便于監(jiān)控和排查問題消息

3.靈活的重試機制

  • 支持延遲重試
  • 可設(shè)置不同的重試策略

4.系統(tǒng)解耦

  • 業(yè)務邏輯與異常處理邏輯分離
  • 提高代碼的可維護性

二、注解式配置說明

1.主配置注解

@Configuration
public class RabbitMQConfig {
    
    // 主隊列
    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue")
            .deadLetterExchange("dlx.exchange")  // 死信交換器
            .deadLetterRoutingKey("dlx.routing.key")  // 死信路由鍵
            .ttl(10000)  // 消息10秒未消費進入死信
            .maxLength(1000)  // 隊列最大長度
            .build();
    }
    
    // 死信隊列
    @Bean
    public Queue deadLetterQueue() {
        return QueueBuilder.durable("dl.queue")
            .build();
    }
    
    // 死信交換器
    @Bean
    public DirectExchange deadLetterExchange() {
        return new DirectExchange("dlx.exchange");
    }
    
    // 綁定死信交換器和隊列
    @Bean
    public Binding deadLetterBinding() {
        return BindingBuilder.bind(deadLetterQueue())
            .to(deadLetterExchange())
            .with("dlx.routing.key");
    }
}

2.監(jiān)聽器注解

@Component
public class OrderMessageListener {
    
    // 監(jiān)聽正常隊列
    @RabbitListener(queues = "order.queue")
    public void processOrderMessage(OrderDTO order, 
                                   Channel channel, 
                                   @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
        try {
            // 業(yè)務處理邏輯
            if (processOrder(order)) {
                // 手動確認
                channel.basicAck(tag, false);
            } else {
                // 拒絕消息,進入死信隊列
                channel.basicNack(tag, false, false);
            }
        } catch (Exception e) {
            // 異常時拒絕
            channel.basicNack(tag, false, false);
        }
    }
    
    // 監(jiān)聽死信隊列
    @RabbitListener(queues = "dl.queue")
    public void processDeadLetter(OrderDTO order) {
        log.error("收到死信消息: {}", order);
        // 死信消息處理邏輯
        handleDeadLetter(order);
    }
}

三、詳細整合步驟

1.添加依賴

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

2.配置屬性

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    # 開啟消息返回機制
    publisher-returns: true
    # 開啟確認機制
    publisher-confirm-type: correlated
    listener:
      simple:
        # 手動確認
        acknowledge-mode: manual
        # 重試配置
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000

3.完整配置類

@Configuration
@Slf4j
public class RabbitMQFullConfig {
    
    // ========== 正常業(yè)務隊列配置 ==========
    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("order.exchange", true, false);
    }
    
    @Bean
    public Queue orderQueue() {
        Map<String, Object> args = new HashMap<>();
        // 死信交換器
        args.put("x-dead-letter-exchange", "order.dlx.exchange");
        // 死信路由鍵
        args.put("x-dead-letter-routing-key", "order.dlx.key");
        // 消息TTL(毫秒)
        args.put("x-message-ttl", 30000);
        // 隊列最大長度
        args.put("x-max-length", 10000);
        return QueueBuilder.durable("order.queue")
            .withArguments(args)
            .build();
    }
    
    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
            .to(orderExchange())
            .with("order.key");
    }
    
    // ========== 死信隊列配置 ==========
    @Bean
    public DirectExchange deadLetterExchange() {
        return new DirectExchange("order.dlx.exchange", true, false);
    }
    
    @Bean
    public Queue deadLetterQueue() {
        return QueueBuilder.durable("order.dl.queue")
            .build();
    }
    
    @Bean
    public Binding deadLetterBinding() {
        return BindingBuilder.bind(deadLetterQueue())
            .to(deadLetterExchange())
            .with("order.dlx.key");
    }
    
    // ========== 重試隊列(延時隊列替代方案)==========
    @Bean
    public CustomExchange delayExchange() {
        Map<String, Object> args = new HashMap<>();
        args.put("x-delayed-type", "direct");
        return new CustomExchange("delay.exchange", 
            "x-delayed-message", true, false, args);
    }
    
    @Bean
    public Queue delayQueue() {
        return QueueBuilder.durable("delay.queue")
            .build();
    }
    
    @Bean
    public Binding delayBinding() {
        return BindingBuilder.bind(delayQueue())
            .to(delayExchange())
            .with("delay.key")
            .noargs();
    }
}

4.消息生產(chǎn)者

@Component
@Slf4j
public class MessageProducer {
    
    @Autowired
    private RabbitTemplate rabbitTemplate;
    
    // 發(fā)送普通消息
    public void sendOrderMessage(OrderDTO order) {
        CorrelationData correlationData = new CorrelationData(order.getId());
        
        rabbitTemplate.convertAndSend(
            "order.exchange",
            "order.key",
            order,
            message -> {
                // 設(shè)置消息屬性
                message.getMessageProperties()
                    .setExpiration("30000")  // 消息TTL
                    .setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                return message;
            },
            correlationData
        );
        
        // 確認回調(diào)
        correlationData.getFuture().addCallback(
            result -> {
                if (result.isAck()) {
                    log.info("消息發(fā)送成功: {}", order.getId());
                }
            },
            ex -> log.error("消息發(fā)送失敗: {}", ex.getMessage())
        );
    }
    
    // 發(fā)送延遲消息
    public void sendDelayMessage(OrderDTO order, int delayTime) {
        rabbitTemplate.convertAndSend(
            "delay.exchange",
            "delay.key",
            order,
            message -> {
                message.getMessageProperties()
                    .setHeader("x-delay", delayTime);
                return message;
            }
        );
    }
}

5.消息消費者(完整版)

@Component
@Slf4j
public class OrderMessageConsumer {
    
    private static final int MAX_RETRY_COUNT = 3;
    
    @Autowired
    private MessageProducer messageProducer;
    
    /**
     * 監(jiān)聽訂單隊列
     */
    @RabbitListener(queues = "order.queue")
    public void handleOrderMessage(
            @Payload OrderDTO order,
            @Headers Map<String, Object> headers,
            Channel channel,
            @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
        
        try {
            log.info("收到訂單消息: {}", order);
            
            // 模擬業(yè)務處理
            boolean success = processOrderBusiness(order);
            
            if (success) {
                // 業(yè)務成功,確認消息
                channel.basicAck(deliveryTag, false);
                log.info("訂單處理成功: {}", order.getId());
            } else {
                // 獲取重試次數(shù)
                Integer retryCount = (Integer) headers.get("x-retry-count");
                retryCount = (retryCount == null) ? 1 : retryCount + 1;
                
                if (retryCount <= MAX_RETRY_COUNT) {
                    // 重試次數(shù)未超限,重新入隊
                    log.warn("訂單處理失敗,第{}次重試: {}", retryCount, order.getId());
                    
                    // 設(shè)置重試計數(shù)
                    headers.put("x-retry-count", retryCount);
                    
                    // 延遲重試
                    messageProducer.sendDelayMessage(order, 5000);
                    
                    // 確認消息,避免重新投遞
                    channel.basicAck(deliveryTag, false);
                } else {
                    // 超過重試次數(shù),進入死信隊列
                    log.error("訂單處理失敗次數(shù)超過上限,進入死信隊列: {}", order.getId());
                    channel.basicNack(deliveryTag, false, false);
                }
            }
        } catch (Exception e) {
            log.error("處理訂單消息異常: {}", e.getMessage());
            try {
                // 拒絕消息,進入死信隊列
                channel.basicNack(deliveryTag, false, false);
            } catch (IOException ex) {
                log.error("拒絕消息失敗: {}", ex.getMessage());
            }
        }
    }
    
    /**
     * 監(jiān)聽死信隊列
     */
    @RabbitListener(queues = "order.dl.queue")
    public void handleDeadLetterMessage(
            @Payload OrderDTO order,
            @Headers Map<String, Object> headers) {
        
        log.error("收到死信消息: {}", order);
        
        // 記錄死信消息
        logDeadLetter(order, headers);
        
        // 發(fā)送告警
        sendAlert(order);
        
        // 人工處理或其他補償措施
        manualProcess(order);
    }
    
    /**
     * 監(jiān)聽延遲隊列
     */
    @RabbitListener(queues = "delay.queue")
    public void handleDelayMessage(@Payload OrderDTO order) {
        log.info("收到延遲消息,開始重試: {}", order);
        
        // 重新發(fā)送到訂單隊列
        messageProducer.sendOrderMessage(order);
    }
    
    private boolean processOrderBusiness(OrderDTO order) {
        // 業(yè)務處理邏輯
        // 返回true表示成功,false表示失敗
        return new Random().nextBoolean();
    }
    
    private void logDeadLetter(OrderDTO order, Map<String, Object> headers) {
        // 記錄死信日志
        log.info("記錄死信: {}, headers: {}", order, headers);
    }
    
    private void sendAlert(OrderDTO order) {
        // 發(fā)送告警通知
        log.warn("發(fā)送告警: 訂單{}處理失敗", order.getId());
    }
    
    private void manualProcess(OrderDTO order) {
        // 人工處理邏輯
        log.info("等待人工處理訂單: {}", order.getId());
    }
}

四、使用場景

1.訂單超時取消

// 訂單創(chuàng)建時發(fā)送延遲消息
public void createOrder(OrderDTO order) {
    // 保存訂單
    orderService.save(order);
    
    // 發(fā)送30分鐘過期的消息
    rabbitTemplate.convertAndSend(
        "order.exchange",
        "order.key",
        order,
        message -> {
            message.getMessageProperties()
                .setExpiration("1800000");  // 30分鐘
            return message;
        }
    );
}

2.支付回調(diào)重試

// 支付回調(diào)失敗時進入死信隊列,人工處理
@RabbitListener(queues = "payment.callback.queue")
public void handlePaymentCallback(PaymentDTO payment) {
    if (!paymentService.processCallback(payment)) {
        throw new RuntimeException("支付回調(diào)處理失敗");
    }
}

3.庫存鎖定與釋放

// 庫存鎖定15分鐘后自動釋放
public void lockInventory(String orderId) {
    inventoryService.lock(orderId);
    
    // 發(fā)送15分鐘后到期的消息
    rabbitTemplate.convertAndSend(
        "inventory.exchange",
        "inventory.lock.key",
        orderId,
        message -> {
            message.getMessageProperties()
                .setExpiration("900000");  // 15分鐘
            return message;
        }
    );
}

4.消息重試機制

// 分級重試策略
public class RetryStrategy {
    // 第一次重試:5秒后
    // 第二次重試:30秒后
    // 第三次重試:5分鐘后
    // 超過3次進入死信隊列
}

五、優(yōu)點總結(jié)

  1. 可靠性:確保消息不丟失,即使處理失敗也有備份
  2. 靈活性:支持多種死信策略(超時、長度限制、拒絕等)
  3. 可維護性:異常處理與正常業(yè)務邏輯分離
  4. 監(jiān)控性:死信隊列便于監(jiān)控和統(tǒng)計異常消息
  5. 可擴展性:支持多種重試和補償機制

六、最佳實踐建議

  1. 合理設(shè)置TTL:根據(jù)業(yè)務需求設(shè)置合適的過期時間
  2. 監(jiān)控死信隊列:設(shè)置告警,及時處理死信消息
  3. 限制隊列大小:防止消息積壓
  4. 記錄詳細日志:便于問題排查
  5. 死信消息分析:定期分析死信原因,優(yōu)化系統(tǒng)

通過Spring Boot整合RabbitMQ死信隊列,可以構(gòu)建更加健壯、可靠的消息驅(qū)動系統(tǒng),有效處理各種異常場景,提高系統(tǒng)的整體穩(wěn)定性。

到此這篇關(guān)于SpringBoot4.0整合RabbitMQ死信隊列詳解的文章就介紹到這了,更多相關(guān)SpringBoot RabbitMQ死信隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • SpringSecurity自定義成功失敗處理器的示例代碼

    SpringSecurity自定義成功失敗處理器的示例代碼

    這篇文章主要介紹了SpringSecurity自定義成功失敗處理器,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2020-09-09
  • Spring Cloud出現(xiàn)Options Forbidden 403問題解決方法

    Spring Cloud出現(xiàn)Options Forbidden 403問題解決方法

    本篇文章主要介紹了Spring Cloud出現(xiàn)Options Forbidden 403問題解決方法,具有一定的參考價值,有興趣的可以了解一下
    2017-11-11
  • Java實現(xiàn)的自定義類加載器示例

    Java實現(xiàn)的自定義類加載器示例

    這篇文章主要介紹了Java實現(xiàn)的自定義類加載器,結(jié)合具體實例形式分析了java自定義類加載器的原理與具體實現(xiàn)技巧,需要的朋友可以參考下
    2019-07-07
  • Java中泛型學習之細節(jié)篇

    Java中泛型學習之細節(jié)篇

    泛型在java中有很重要的地位,在面向?qū)ο缶幊碳案鞣N設(shè)計模式中有非常廣泛的應用,下面這篇文章主要給大家介紹了關(guān)于Java中泛型細節(jié)的相關(guān)資料,文中通過實例代碼介紹的非常詳細,需要的朋友可以參考下
    2022-02-02
  • springboot整合sa-token中的redis報netty錯誤問題

    springboot整合sa-token中的redis報netty錯誤問題

    整合Spring Boot與sa-token-redis-jackson時遇到Netty版本沖突,通過將netty-common升級到與sa-token-redis-jackson兼容的版本4.1.79解決
    2024-11-11
  • maven父工程relativepath標簽使用解讀

    maven父工程relativepath標簽使用解讀

    文章主要介紹了在使用Maven構(gòu)建父子工程時如何通過設(shè)置父工程和子工程的pom文件來管理依賴和版本,當子工程是Spring Boot項目時,可以通過關(guān)閉`relativePath`標簽來繼承Spring Boot的父工程,同時在父工程中使用`dependencyManagement`標簽來統(tǒng)一管理Spring Boot的依賴版本
    2024-11-11
  • java線程池中Worker線程執(zhí)行流程原理解析

    java線程池中Worker線程執(zhí)行流程原理解析

    這篇文章主要為大家介紹了java線程池中Worker線程執(zhí)行流程原理解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2022-11-11
  • Spring使用注解解決線程并發(fā)問題的多種方案

    Spring使用注解解決線程并發(fā)問題的多種方案

    在?Spring?中,線程并發(fā)問題的核心是共享資源的競爭,注解本身不能直接解決并發(fā)問題,但?Spring?提供了配套注解結(jié)合并發(fā)編程規(guī)范,可以優(yōu)雅地實現(xiàn)線程安全控制,下面我會從最常用的場景出發(fā),講解?Spring?中如何通過注解解決并發(fā)問題,需要的朋友可以參考下
    2026-03-03
  • 一文搞懂Spring中的注解與反射

    一文搞懂Spring中的注解與反射

    這篇文章主要為大家介紹了Spring中的注解與反射的原理與實現(xiàn),文中的示例代碼講解詳細,對我們了解Spring有一定的幫助,需要的可以參考一下
    2022-06-06
  • 使用JAVA獲取nacos配置信息出現(xiàn)null,獲取不到的解決

    使用JAVA獲取nacos配置信息出現(xiàn)null,獲取不到的解決

    文章討論了在使用Java獲取Nacos配置信息時遇到的問題,特別是調(diào)用ConfigService獲取配置時出現(xiàn)null的情況,作者嘗試了多種解決方法,最終發(fā)現(xiàn)更換jar包版本(1.*)解決了問題,作者分享了個人經(jīng)驗,希望能對大家有所幫助
    2025-12-12

最新評論

江孜县| 瑞昌市| 永和县| 常德市| 霞浦县| 拉萨市| 东乌珠穆沁旗| 高阳县| 鹤峰县| 同心县| 彩票| 民丰县| 高阳县| 阿图什市| 安平县| 迭部县| 道孚县| 都匀市| 盈江县| 略阳县| 阿拉善右旗| 榆树市| 曲水县| 丰原市| 腾冲县| 滨海县| 南部县| 兴仁县| 和平区| 普兰店市| 武川县| 浙江省| 遵义县| 垣曲县| 桂东县| 桃园市| 株洲县| 营山县| 屯昌县| 湘阴县| 南昌市|