springboot實現rabbitmq消息確認的示例代碼
概述
RabbitMQ的消息確認有兩種。 一種是消息發(fā)送確認。這種是用來確認生產者將消息發(fā)送給交換器,交換器傳遞給隊列的過程中,消息是否成功投遞。發(fā)送確認分為兩步,一是確認是否到達交換器,二是確認是否到達隊列。 第二種是消費接收確認。這種是確認消費者是否成功消費了隊列中的消息。
一、運行效果

二、實現過程
①、引入rabbitmq包
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>②、修改application.properties配置
spring.rabbitmq.host=127.0.0.1 spring.rabbitmq.port=5672 spring.rabbitmq.username=guest spring.rabbitmq.password=guest # 發(fā)送者開啟 confirm 確認機制 spring.rabbitmq.publisher-confirms=true # 發(fā)送者開啟 return 確認機制 spring.rabbitmq.publisher-returns=true #################################################### # 設置消費端手動 ack spring.rabbitmq.listener.simple.acknowledge-mode=manual # 是否支持重試 spring.rabbitmq.listener.simple.retry.enabled=true
③、定義exchange和queue,并將queue綁定在exchange上
package com.mm.springbootrabbitmqconfirmdemo.config;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.FanoutExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Configuration
public class RabbitMQConfig {
@Bean(name = "confirmQueue")
public Queue confirmQueue(){
return new Queue("confirmQueue",true,false,false);
}
@Bean(name = "confirmExchange")
public FanoutExchange confirmExchange(){
return new FanoutExchange("confirmExchange");
}
@Bean
public Binding confirmFanoutExchangeAndQueue(@Qualifier("confirmExchange") FanoutExchange confirmExchange,
@Qualifier("confirmQueue") Queue confirmQueue){
return BindingBuilder.bind(confirmQueue).to(confirmExchange);
}
}④、消息發(fā)送確認
發(fā)送消息確認:用來確認生產者 producer 將消息發(fā)送到 broker ,broker 上的交換機 exchange 再投遞給隊列 queue的過程中,消息是否成功投遞。
消息從 producer 到 rabbitmq broker有一個 confirmCallback 確認模式。
消息從 exchange 到 queue 投遞失敗有一個 returnCallback 退回模式。
我們可以利用這兩個Callback來確保消息的100%送達。
1、 ConfirmCallback確認模式
消息只要被 rabbitmq broker 接收到就會觸發(fā) confirmCallback 回調 。
package com.mm.springbootrabbitmqconfirmdemo.service;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class ConfirmCallbackService implements RabbitTemplate.ConfirmCallback {
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause){
if (!ack) {
log.error("消息發(fā)送異常!");
} else {
log.info("發(fā)送者爸爸已經收到確認,correlationData={} ,ack={}, cause={}", correlationData.getId(), ack, cause);
}
}
}實現接口 ConfirmCallback ,重寫其confirm()方法,方法內有三個參數correlationData、ack、cause。
correlationData:對象內部只有一個id屬性,用來表示當前消息的唯一性。ack:消息投遞到broker的狀態(tài),true表示成功。cause:表示投遞失敗的原因。
但消息被 broker 接收到只能表示已經到達 MQ服務器,并不能保證消息一定會被投遞到目標 queue 里。所以接下來需要用到 returnCallback 。
2、 ReturnCallback 退回模式
如果消息未能投遞到目標 queue 里將觸發(fā)回調 returnCallback ,一旦向 queue 投遞消息未成功,這里一般會記錄下當前消息的詳細投遞數據,方便后續(xù)做重發(fā)或者補償等操作。
com.mm.springbootrabbitmqconfirmdemo.service; lombok.extern.slf4j.; org.springframework.amqp.core.Message; org.springframework.amqp.rabbit.core.RabbitTemplate; org.springframework.stereotype.; ReturnCallbackService RabbitTemplate.ReturnCallback returnedMessageMessage message, replyCode, String replyText, String exchange, String routingKey.info, replyCode, replyText, exchange, routingKey;
實現接口ReturnCallback,重寫 returnedMessage() 方法,方法有五個參數message(消息體)、replyCode(響應code)、replyText(響應內容)、exchange(交換機)、routingKey(隊列)。
下邊是具體的消息發(fā)送,在rabbitTemplate中設置 Confirm 和 Return 回調,我們通過setDeliveryMode()對消息做持久化處理,為了后續(xù)測試創(chuàng)建一個 CorrelationData對象,添加一個id 為10000000000。
⑤、消息發(fā)送確認
消息接收確認要比消息發(fā)送確認簡單一點,因為只有一個消息回執(zhí)(ack)的過程。使用@RabbitHandler注解標注的方法要增加 channel(信道)、message 兩個參數。
@Slf4j
@Component
@RabbitListener(queues = "confirm_test_queue")
public class ReceiverMessage1 {
@RabbitHandler
public void processHandler(String msg, Channel channel, Message message) throws IOException {
try {
log.info("小富收到消息:{}", msg);
//TODO 具體業(yè)務
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
if (message.getMessageProperties().getRedelivered()) {
log.error("消息已重復處理失敗,拒絕再次接收...");
channel.basicReject(message.getMessageProperties().getDeliveryTag(), false); // 拒絕消息
} else {
log.error("消息即將再次返回隊列處理...");
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
}
}
}
}消費消息有三種回執(zhí)方法,我們來分析一下每種方法的含義。
1、basicAck
basicAck:表示成功確認,使用此回執(zhí)方法后,消息會被rabbitmq broker 刪除。
void?basicAck(long?deliveryTag,?boolean?multiple)
deliveryTag:表示消息投遞序號,每次消費消息或者消息重新投遞后,deliveryTag都會增加。手動消息確認模式下,我們可以對指定deliveryTag的消息進行ack、nack、reject等操作。
multiple:是否批量確認,值為 true 則會一次性 ack所有小于當前消息 deliveryTag 的消息。
舉個栗子: 假設我先發(fā)送三條消息deliveryTag分別是5、6、7,可它們都沒有被確認,當我發(fā)第四條消息此時deliveryTag為8,multiple設置為 true,會將5、6、7、8的消息全部進行確認。
2、basicNack
basicNack :表示失敗確認,一般在消費消息業(yè)務異常時用到此方法,可以將消息重新投遞入隊列。
void?basicNack(long?deliveryTag,?boolean?multiple,?boolean?requeue)
deliveryTag:表示消息投遞序號。
multiple:是否批量確認。
requeue:值為 true 消息將重新入隊列。
3、basicReject
basicReject:拒絕消息,與basicNack區(qū)別在于不能進行批量操作,其他用法很相似。
void?basicReject(long?deliveryTag,?boolean?requeue)
deliveryTag:表示消息投遞序號。
requeue:值為 true 消息將重新入隊列。
三、項目結構圖

四、補充
1、別忘確認消息
這是一個非常沒技術含量的坑,但卻是非常容易犯錯的地方。
開啟消息確認機制,消費消息別忘了channel.basicAck,否則消息會一直存在,導致重復消費。
2、消息無限投遞
在我最開始接觸消息確認機制的時候,消費端代碼就像下邊這樣寫的,思路很簡單:處理完業(yè)務邏輯后確認消息, int a = 1 / 0 發(fā)生異常后將消息重新投入隊列。
@RabbitHandler
public void processHandler(String msg, Channel channel, Message message) throws IOException {
try {
log.info("消費者 2 號收到:{}", msg);
int a = 1 / 0;
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
}
}3、重復消費
如何保證 MQ 的消費是冪等性,這個需要根據具體業(yè)務而定,可以借助MySQL、或者redis將消息持久化,通過再消息中的唯一性屬性校驗。
可以看到使用了 RabbitMQ 以后,我們的業(yè)務鏈路明顯變長了,雖然做到了系統(tǒng)間的解耦,但可能造成消息丟失的場景也增加了。例如:
消息生產者 - > rabbitmq服務器(消息發(fā)送失?。?/p>
rabbitmq服務器自身故障導致消息丟失
消息消費者 - > rabbitmq服務(消費消息失?。?/p>
到此這篇關于springboot實現rabbitmq消息確認的示例代碼的文章就介紹到這了,更多相關springboot rabbitmq消息確認內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
IDEA Java win10環(huán)境配置的圖文教程
這篇文章主要介紹了IDEA Java win10環(huán)境配置,本文通過圖文并茂的形式給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2020-07-07
利用Spring MVC+Mybatis實現Mysql分頁數據查詢的過程詳解
這篇文章主要給大家介紹了關于利用Spring MVC+Mybatis實現Mysql分頁數據查詢的相關資料,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面跟著小編來一起學習學習吧。2017-08-08
SpringBoot?使用?Sa-Token?完成權限認證的操作方法
Sa-Token是一款輕量級Java權限認證框架,適用于快速搭建權限系統(tǒng),它提供了豐富的功能,下面給大家介紹SpringBoot?使用?Sa-Token?完成權限認證的相關操作,感興趣的朋友跟隨小編一起看看吧2025-02-02
Mybatis #foreach中相同的變量名導致值覆蓋的問題解決
本文主要介紹了Mybatis #foreach中相同的變量名導致值覆蓋的問題解決,具有一定的參考價值,感興趣的小伙伴們可以參考一下2021-07-07

