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

SpringBoot4.0整合RabbitMQ死信隊(duì)列詳解

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

為啥那么講解死信隊(duì)列,因?yàn)楹枚嗳瞬粫?huì)使用,不知道什么場(chǎng)景下使用,此案例是我在公司實(shí)現(xiàn)的一種方式,讓大家都可以學(xué)習(xí)到

一、死信隊(duì)列的好處

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

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

2.異常消息管理

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

3.靈活的重試機(jī)制

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

4.系統(tǒng)解耦

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

二、注解式配置說(shuō)明

1.主配置注解

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

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

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

三、詳細(xì)整合步驟

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
    # 開(kāi)啟消息返回機(jī)制
    publisher-returns: true
    # 開(kāi)啟確認(rèn)機(jī)制
    publisher-confirm-type: correlated
    listener:
      simple:
        # 手動(dòng)確認(rèn)
        acknowledge-mode: manual
        # 重試配置
        retry:
          enabled: true
          max-attempts: 3
          initial-interval: 1000

3.完整配置類

@Configuration
@Slf4j
public class RabbitMQFullConfig {
    
    // ========== 正常業(yè)務(wù)隊(duì)列配置 ==========
    @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);
        // 隊(duì)列最大長(zhǎng)度
        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");
    }
    
    // ========== 死信隊(duì)列配置 ==========
    @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");
    }
    
    // ========== 重試隊(duì)列(延時(shí)隊(duì)列替代方案)==========
    @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
        );
        
        // 確認(rèn)回調(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.消息消費(fèi)者(完整版)

@Component
@Slf4j
public class OrderMessageConsumer {
    
    private static final int MAX_RETRY_COUNT = 3;
    
    @Autowired
    private MessageProducer messageProducer;
    
    /**
     * 監(jiān)聽(tīng)訂單隊(duì)列
     */
    @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è)務(wù)處理
            boolean success = processOrderBusiness(order);
            
            if (success) {
                // 業(yè)務(wù)成功,確認(rèn)消息
                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ù)未超限,重新入隊(duì)
                    log.warn("訂單處理失敗,第{}次重試: {}", retryCount, order.getId());
                    
                    // 設(shè)置重試計(jì)數(shù)
                    headers.put("x-retry-count", retryCount);
                    
                    // 延遲重試
                    messageProducer.sendDelayMessage(order, 5000);
                    
                    // 確認(rèn)消息,避免重新投遞
                    channel.basicAck(deliveryTag, false);
                } else {
                    // 超過(guò)重試次數(shù),進(jìn)入死信隊(duì)列
                    log.error("訂單處理失敗次數(shù)超過(guò)上限,進(jìn)入死信隊(duì)列: {}", order.getId());
                    channel.basicNack(deliveryTag, false, false);
                }
            }
        } catch (Exception e) {
            log.error("處理訂單消息異常: {}", e.getMessage());
            try {
                // 拒絕消息,進(jìn)入死信隊(duì)列
                channel.basicNack(deliveryTag, false, false);
            } catch (IOException ex) {
                log.error("拒絕消息失敗: {}", ex.getMessage());
            }
        }
    }
    
    /**
     * 監(jiān)聽(tīng)死信隊(duì)列
     */
    @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);
        
        // 人工處理或其他補(bǔ)償措施
        manualProcess(order);
    }
    
    /**
     * 監(jiān)聽(tīng)延遲隊(duì)列
     */
    @RabbitListener(queues = "delay.queue")
    public void handleDelayMessage(@Payload OrderDTO order) {
        log.info("收到延遲消息,開(kāi)始重試: {}", order);
        
        // 重新發(fā)送到訂單隊(duì)列
        messageProducer.sendOrderMessage(order);
    }
    
    private boolean processOrderBusiness(OrderDTO order) {
        // 業(yè)務(wù)處理邏輯
        // 返回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());
    }
}

四、使用場(chǎng)景

1.訂單超時(shí)取消

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

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

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

3.庫(kù)存鎖定與釋放

// 庫(kù)存鎖定15分鐘后自動(dòng)釋放
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.消息重試機(jī)制

// 分級(jí)重試策略
public class RetryStrategy {
    // 第一次重試:5秒后
    // 第二次重試:30秒后
    // 第三次重試:5分鐘后
    // 超過(guò)3次進(jìn)入死信隊(duì)列
}

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

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

六、最佳實(shí)踐建議

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

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

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

相關(guān)文章

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

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

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

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

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

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

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

    Java中泛型學(xué)習(xí)之細(xì)節(jié)篇

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

    springboot整合sa-token中的redis報(bào)netty錯(cuò)誤問(wèn)題

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

    maven父工程relativepath標(biāo)簽使用解讀

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

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

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

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

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

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

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

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

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

最新評(píng)論

吴桥县| 永川市| 抚顺县| 东乡族自治县| 英吉沙县| 施秉县| 无锡市| 尤溪县| 辛集市| 临澧县| 双鸭山市| 合阳县| 凤山县| 枣阳市| 东明县| 嘉荫县| 西丰县| 韶关市| 开鲁县| 芜湖市| 泌阳县| 铜川市| 仙游县| 廉江市| 改则县| 湟源县| 班戈县| 岐山县| 石首市| 盖州市| 洛川县| 那坡县| 广州市| 通河县| 普兰县| 砀山县| 和平区| 德昌县| 庄河市| 新化县| 望江县|