詳解RabbitMQ中延遲隊(duì)列結(jié)合業(yè)務(wù)場(chǎng)景的使用
消息進(jìn)入隊(duì)列后不會(huì)立即被消費(fèi),只有到達(dá)指定時(shí)間后才會(huì)被消費(fèi)。業(yè)務(wù)場(chǎng)景就是支付時(shí)間內(nèi)未支付就清除訂單或者用戶(hù)注冊(cè)一段時(shí)間后發(fā)短信問(wèn)候。在這里想說(shuō)的是這只是一種思想,并不是真正的一種用法,這種思想所需要的用法就是用上消息TTL存活時(shí)間以及死信隊(duì)列來(lái)實(shí)現(xiàn)。
生產(chǎn)者端
目錄結(jié)構(gòu)

導(dǎo)入依賴(lài)
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<version>2.5.0</version>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.12</version>
<scope>test</scope>
</dependency>
</dependencies>修改yml
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
publisher-returns: true # 開(kāi)啟退回回調(diào)
#三個(gè)類(lèi)型:none默認(rèn)不開(kāi)啟確認(rèn)回調(diào) correlated開(kāi)啟確認(rèn)回調(diào)
#simple也會(huì)確認(rèn)回調(diào) 還會(huì)調(diào)用waitForConfirms()方法或waitForConfirmsOrDie()方法
publisher-confirm-type: correlated # 開(kāi)啟確認(rèn)回調(diào)
業(yè)務(wù)邏輯
@Component
public class RabbitMQConfig {
public static final String EXCHANGE_NAME = "order_exchange_name";
public static final String QUEUE_NAME = "order_queue_name";
public static final String DELAY_EXCHANGE_NAME = "delay_exchange_name";
public static final String DELAY_QUEUE_NAME = "delay_queue_name";
@Bean("orderExchange")
public Exchange testExchange(){
return ExchangeBuilder.topicExchange(EXCHANGE_NAME).durable(true).build();
}
@Bean("delayExchange")
public Exchange deadExchange(){
return ExchangeBuilder.topicExchange(DELAY_EXCHANGE_NAME).durable(true).build();
}
//訂單隊(duì)列綁定延遲交換機(jī)并且?guī)下酚涉I
@Bean("orderQueue")
public Queue testQueue(){
return QueueBuilder.durable(QUEUE_NAME).deadLetterExchange(DELAY_EXCHANGE_NAME)
.deadLetterRoutingKey("order.delay.user").build();
}
@Bean("delayQueue")
public Queue deadQueue(){
return QueueBuilder.durable(DELAY_QUEUE_NAME).build();
}
@Bean
public Binding link(@Qualifier("orderExchange") Exchange exchange,
@Qualifier("orderQueue") Queue queue){
return BindingBuilder.bind(queue).to(exchange).with("order.#").noargs();
}
@Bean
public Binding deadLink(@Qualifier("delayExchange") Exchange exchange,
@Qualifier("delayQueue") Queue queue){
return BindingBuilder.bind(queue).to(exchange).with("order.delay.#").noargs();
}
}@SpringBootTest
@RunWith(SpringRunner.class)
class RabbitmqProducerApplicationTests {
@Autowired
private RabbitTemplate rabbitTemplate;
@Test
void testProducer() throws InterruptedException {
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
@Override
public void confirm(CorrelationData correlationData, boolean b, String s) {
if(b) System.out.println("交換機(jī)成功接受到了消息");
else System.out.println("消息失敗原因" + s);
}
});
// 設(shè)置交換機(jī)處理失敗消息的模式
// true:消息到達(dá)不了隊(duì)列時(shí) 會(huì)將消息重新返回給生產(chǎn)者 false:消息到達(dá)不了隊(duì)列直接丟棄
rabbitTemplate.setMandatory(true);
rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
@Override
public void returnedMessage(Message message, int i, String s, String s1, String s2) {
System.out.println("隊(duì)列接受不到交換機(jī)的消息進(jìn)行了失敗回調(diào)");
}
});
// 以上代碼均是為了保證消息的可靠性傳遞
// 對(duì)消息進(jìn)行后置處理 設(shè)置其過(guò)期時(shí)間為10s
MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
message.getMessageProperties().setExpiration("10000");
return message;
}
};
// 下單成功發(fā)送消息
rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_NAME,"order.delay.user","ikun書(shū)籍", messagePostProcessor);
}
}測(cè)試結(jié)果

消費(fèi)者端
目錄結(jié)構(gòu)

導(dǎo)入依賴(lài)
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<version>2.5.0</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</dependency>
</dependencies>修改yml
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
listener:
simple:
acknowledge-mode: manual # 開(kāi)啟手動(dòng)確認(rèn)
業(yè)務(wù)邏輯
@Slf4j
@Component
public class OrderListener implements ChannelAwareMessageListener {
@RabbitListener(queues = "delay_queue_name") // 監(jiān)聽(tīng)的是死信隊(duì)列
@Override
public void onMessage(Message message, Channel channel) throws Exception {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
// 這里要注意在監(jiān)聽(tīng)這條消息前 肯定會(huì)有接口點(diǎn)擊支付去更改支付狀態(tài)的業(yè)務(wù)邏輯 此處我是做了10s的訂單業(yè)務(wù)而已
log.info("您在時(shí)間為:{},時(shí)有一條訂單為:{}", LocalDateTime.now().minusSeconds(10), new String(message.getBody()));
// 下面開(kāi)始接受訂單消息邏輯
log.info("將訂單id傳入數(shù)據(jù)庫(kù)查詢(xún)訂單支付字段");
log.info("字段為支付成功狀態(tài)就手動(dòng)確認(rèn)簽收");
log.info("字段為未支付狀態(tài)就取消訂單并且回滾事務(wù)");
channel.basicAck(deliveryTag,false);// 僅確認(rèn)本次消息
} catch (Exception e){
log.info("出現(xiàn)異常 拒絕簽收消息 并且不重回隊(duì)列");
channel.basicNack(deliveryTag,false,false);
}
}
}測(cè)試結(jié)果


到此這篇關(guān)于詳解RabbitMQ中延遲隊(duì)列結(jié)合業(yè)務(wù)場(chǎng)景的使用的文章就介紹到這了,更多相關(guān)RabbitMQ延遲隊(duì)列內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
實(shí)戰(zhàn)分布式醫(yī)療掛號(hào)系統(tǒng)之設(shè)置微服務(wù)接口開(kāi)發(fā)模塊
java 動(dòng)態(tài)增加定時(shí)任務(wù)示例
springmvc中進(jìn)行數(shù)據(jù)保存以及日期參數(shù)的保存過(guò)程解析
Java實(shí)現(xiàn)簡(jiǎn)單學(xué)生信息管理系統(tǒng)
springboot嵌套子類(lèi)使用方式—前端與后臺(tái)開(kāi)發(fā)的注意事項(xiàng)
Java實(shí)現(xiàn)多個(gè)wav文件合成一個(gè)的方法示例

