<em>Mac</em>Book项目 2009年学校开始实施<em>Mac</em>Book项目,所有师生配备一本<em>Mac</em>Book,并同步更新了校园无线网络。学校每周进行电脑技术更新,每月发送技术支持资料,极大改变了教学及学习方式。因此2011
2021-06-01 09:32:01
使用springboot,實現以下功能,有兩個佇列1、2,往裡面傳送訊息,如果處理失敗發生異常,可以重試3次,重試3次均失敗,那麼就將訊息傳送到死信佇列進行統一處理,例如記錄資料庫、報警等
完整demo專案程式碼https://gitee.com/daenmax/rabbit-mq-demo
Windows10,IDEA,otp_win64_25.0,rabbitmq-server-3.10.4
1.雙擊C:Program FilesRabbitMQ Serverrabbitmq_server-3.10.4sbinrabbitmq-server.bat啟動MQ服務
2.然後存取http://localhost:15672/,預設賬號密碼均為guest,
3.手動新增一個虛擬主機為admin_host,手動建立一個使用者賬號密碼均為admin
pom.xml
<!-- RabbitMQ --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> <version>2.7.0</version> </dependency>
spring: rabbitmq: host: 127.0.0.1 port: 5672 username: admin password: admin virtual-host: admin_host publisher-confirm-type: correlated publisher-returns: true listener: simple: acknowledge-mode: manual retry: enabled: true #開啟失敗重試 max-attempts: 3 #最大重試次數 initial-interval: 1000 #重試間隔時間 毫秒
RabbitConfig
package com.example.rabitmqdemo.mydemo.config; import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.core.*; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.stereotype.Component; import java.util.HashMap; import java.util.Map; /** * Broker:它提供一種傳輸服務,它的角色就是維護一條從生產者到消費者的路線,保證資料能按照指定的方式進行傳輸, * Exchange:訊息交換機,它指定訊息按什麼規則,路由到哪個佇列。 * Queue:訊息的載體,每個訊息都會被投到一個或多個佇列。 * Binding:繫結,它的作用就是把exchange和queue按照路由規則繫結起來. * Routing Key:路由關鍵字,exchange根據這個關鍵字進行訊息投遞。 * vhost:虛擬主機,一個broker裡可以有多個vhost,用作不同使用者的許可權分離。 * Producer:訊息生產者,就是投遞訊息的程式. * Consumer:訊息消費者,就是接受訊息的程式. * Channel:訊息通道,在使用者端的每個連線裡,可建立多個channel. */ @Slf4j @Component public class RabbitConfig { //業務交換機 public static final String EXCHANGE_PHCP = "phcp"; //業務佇列1 public static final String QUEUE_COMPANY = "company"; //業務佇列1的key public static final String ROUTINGKEY_COMPANY = "companyKey"; //業務佇列2 public static final String QUEUE_PROJECT = "project"; //業務佇列2的key public static final String ROUTINGKEY_PROJECT = "projectKey"; //死信交換機 public static final String EXCHANGE_PHCP_DEAD = "phcp_dead"; //死信佇列1 public static final String QUEUE_COMPANY_DEAD = "company_dead"; //死信佇列2 public static final String QUEUE_PROJECT_DEAD = "project_dead"; //死信佇列1的key public static final String ROUTINGKEY_COMPANY_DEAD = "companyKey_dead"; //死信佇列2的key public static final String ROUTINGKEY_PROJECT_DEAD = "projectKey_dead"; // /** // * 解決重複確認報錯問題,如果沒有報錯的話,就不用啟用這個 // * // * @param connectionFactory // * @return // */ // @Bean // public RabbitListenerContainerFactory<?> rabbitListenerContainerFactory(ConnectionFactory connectionFactory) { // SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); // factory.setConnectionFactory(connectionFactory); // factory.setMessageConverter(new Jackson2JsonMessageConverter()); // factory.setAcknowledgeMode(AcknowledgeMode.MANUAL); // return factory; // } /** * 宣告業務交換機 * 1. 設定交換機型別 * 2. 將佇列繫結到交換機 * FanoutExchange: 將訊息分發到所有的繫結佇列,無routingkey的概念 * HeadersExchange :通過新增屬性key-value匹配 * DirectExchange:按照routingkey分發到指定佇列 * TopicExchange:多關鍵字匹配 */ @Bean("exchangePhcp") public DirectExchange exchangePhcp() { return new DirectExchange(EXCHANGE_PHCP); } * 宣告死信交換機 @Bean("exchangePhcpDead") public DirectExchange exchangePhcpDead() { return new DirectExchange(EXCHANGE_PHCP_DEAD); * 宣告業務佇列1 * * @return @Bean("queueCompany") public Queue queueCompany() { Map<String,Object> arguments = new HashMap<>(2); arguments.put("x-dead-letter-exchange",EXCHANGE_PHCP_DEAD); //繫結該佇列到死信交換機的佇列1 arguments.put("x-dead-letter-routing-key",ROUTINGKEY_COMPANY_DEAD); return QueueBuilder.durable(QUEUE_COMPANY).withArguments(arguments).build(); * 宣告業務佇列2 @Bean("queueProject") public Queue queueProject() { //繫結該佇列到死信交換機的佇列2 arguments.put("x-dead-letter-routing-key",ROUTINGKEY_PROJECT_DEAD); return QueueBuilder.durable(QUEUE_PROJECT).withArguments(arguments).build(); * 宣告死信佇列1 @Bean("queueCompanyDead") public Queue queueCompanyDead() { return new Queue(QUEUE_COMPANY_DEAD); * 宣告死信佇列2 @Bean("queueProjectDead") public Queue queueProjectDead() { return new Queue(QUEUE_PROJECT_DEAD); * 繫結業務佇列1和業務交換機 * @param queue * @param directExchange @Bean public Binding bindingQueueCompany(@Qualifier("queueCompany") Queue queue, @Qualifier("exchangePhcp") DirectExchange directExchange) { return BindingBuilder.bind(queue).to(directExchange).with(RabbitConfig.ROUTINGKEY_COMPANY); * 繫結業務佇列2和業務交換機 public Binding bindingQueueProject(@Qualifier("queueProject") Queue queue, @Qualifier("exchangePhcp") DirectExchange directExchange) { return BindingBuilder.bind(queue).to(directExchange).with(RabbitConfig.ROUTINGKEY_PROJECT); * 繫結死信佇列1和死信交換機 public Binding bindingQueueCompanyDead(@Qualifier("queueCompanyDead") Queue queue, @Qualifier("exchangePhcpDead") DirectExchange directExchange) { return BindingBuilder.bind(queue).to(directExchange).with(RabbitConfig.ROUTINGKEY_COMPANY_DEAD); * 繫結死信佇列2和死信交換機 public Binding bindingQueueProjectDead(@Qualifier("queueProjectDead") Queue queue, @Qualifier("exchangePhcpDead") DirectExchange directExchange) { return BindingBuilder.bind(queue).to(directExchange).with(RabbitConfig.ROUTINGKEY_PROJECT_DEAD); }
生產者
RabbltProducer
package com.example.rabitmqdemo.mydemo.producer; import com.example.rabitmqdemo.mydemo.config.RabbitConfig; import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.core.*; import org.springframework.amqp.rabbit.connection.CorrelationData; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import javax.annotation.Resource; import java.nio.charset.StandardCharsets; import java.util.UUID; @Component @Slf4j public class RabbltProducer implements RabbitTemplate.ConfirmCallback, RabbitTemplate.ReturnsCallback{ @Resource private RabbitTemplate rabbitTemplate; /** * 初始化訊息確認函數 */ @PostConstruct public void init() { rabbitTemplate.setConfirmCallback(this); rabbitTemplate.setReturnsCallback(this); rabbitTemplate.setMandatory(true); } /** * 傳送訊息伺服器確認函數 * @param correlationData * @param ack * @param cause */ @Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { if (ack) { System.out.println("訊息傳送成功" + correlationData); } else { System.out.println("訊息傳送失敗:" + cause); } } /** * 訊息傳送失敗,訊息回撥函數 * @param returnedMessage */ @Override public void returnedMessage(ReturnedMessage returnedMessage) { String str = new String(returnedMessage.getMessage().getBody()); System.out.println("訊息傳送失敗:" + str); } /** * 處理訊息傳送到佇列1 * @param str */ public void sendCompany(String str){ MessageProperties messageProperties = new MessageProperties(); messageProperties.setContentType("application/json"); Message message = new Message(str.getBytes(StandardCharsets.UTF_8),messageProperties); CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString()); this.rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_PHCP,RabbitConfig.ROUTINGKEY_COMPANY,message,correlationData); //也可以用下面的方式 //CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString()); //this.rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_PHCP,RabbitConfig.ROUTINGKEY_COMPANY,str,correlationData); } /** * 處理訊息傳送到佇列2 * @param str */ public void sendProject(String str){ MessageProperties messageProperties = new MessageProperties(); messageProperties.setContentType("application/json"); Message message = new Message(str.getBytes(StandardCharsets.UTF_8),messageProperties); CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString()); this.rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_PHCP,RabbitConfig.ROUTINGKEY_PROJECT,message,correlationData); //也可以用下面的方式 //CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString()); //this.rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_PHCP,RabbitConfig.ROUTINGKEY_PROJECT,str,correlationData); } }
RabbitConsumer
package com.example.rabitmqdemo.mydemo.consumer; import com.rabbitmq.client.Channel; import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; /** * 監聽業務交換機 * @author JeWang */ @Component @Slf4j public class RabbitConsumer { /** * 監聽業務佇列1 * @param message * @param channel * @throws IOException */ @RabbitListener(queues = "company") public void company(Message message, Channel channel) throws IOException { try{ System.out.println("次數" + message.getMessageProperties().getDeliveryTag()); channel.basicQos(1); Thread.sleep(2000); String s = new String(message.getBody()); log.info("處理訊息"+s); //下面兩行是嘗試手動丟擲異常,用來測試重試次數和傳送到死信交換機 //String str = null; //str.split("1"); //處理成功,確認應答 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); }catch (Exception e){ log.error("處理訊息時發生異常:"+e.getMessage()); Boolean redelivered = message.getMessageProperties().getRedelivered(); if(redelivered){ log.error("異常重試次數已到達設定次數,將傳送到死信交換機"); channel.basicReject(message.getMessageProperties().getDeliveryTag(),false); }else { log.error("訊息即將返回佇列處理重試"); channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); } } } /** * 監聽業務佇列2 * @param message * @param channel * @throws IOException */ @RabbitListener(queues = "project") public void project(Message message, Channel channel) throws IOException { try{ System.out.println("次數" + message.getMessageProperties().getDeliveryTag()); channel.basicQos(1); Thread.sleep(2000); String s = new String(message.getBody()); log.info("處理訊息"+s); //下面兩行是嘗試手動丟擲異常,用來測試重試次數和傳送到死信交換機 //String str = null; //str.split("1"); //處理成功,確認應答 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); }catch (Exception e){ log.error("處理訊息時發生異常:"+e.getMessage()); Boolean redelivered = message.getMessageProperties().getRedelivered(); if(redelivered){ log.error("異常重試次數已到達設定次數,將傳送到死信交換機"); channel.basicReject(message.getMessageProperties().getDeliveryTag(),false); }else { log.error("訊息即將返回佇列處理重試"); channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); } } } }
RabbitConsumer
package com.example.rabitmqdemo.mydemo.consumer; import com.rabbitmq.client.Channel; import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; /** * 監聽死信交換機 * @author JeWang */ @Component @Slf4j public class RabbitConsumerDead { /** * 處理死信佇列1 * @param message * @param channel * @throws IOException */ @RabbitListener(queues = "company_dead") public void company_dead(Message message, Channel channel) throws IOException { try{ channel.basicQos(1); String s = new String(message.getBody()); log.info("處理死信"+s); //在此處記錄到資料庫、報警之類的操作 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); }catch (Exception e){ log.error("接收異常:"+e.getMessage()); } } /** * 處理死信佇列2 * @param message * @param channel * @throws IOException */ @RabbitListener(queues = "project_dead") public void project_dead(Message message, Channel channel) throws IOException { try{ channel.basicQos(1); String s = new String(message.getBody()); log.info("處理死信"+s); //在此處記錄到資料庫、報警之類的操作 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); }catch (Exception e){ log.error("接收異常:"+e.getMessage()); } } }
MqController
package com.example.rabitmqdemo.mydemo.controller; import com.example.rabitmqdemo.mydemo.producer.RabbltProducer; import lombok.extern.slf4j.Slf4j; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; @RequestMapping("/def") @RestController @Slf4j public class MsgController { @Resource private RabbltProducer rabbltProducer; @RequestMapping("/handleCompany") public void handleCompany(@RequestBody String jsonStr){ rabbltProducer.sendCompany(jsonStr); } }
到此這篇關於SpringBoot整合RabbitMQ實戰附加死信交換機的文章就介紹到這了,更多相關SpringBoot整合RabbitMQ死信交換機內容請搜尋it145.com以前的文章或繼續瀏覽下面的相關文章希望大家以後多多支援it145.com!
相關文章
<em>Mac</em>Book项目 2009年学校开始实施<em>Mac</em>Book项目,所有师生配备一本<em>Mac</em>Book,并同步更新了校园无线网络。学校每周进行电脑技术更新,每月发送技术支持资料,极大改变了教学及学习方式。因此2011
2021-06-01 09:32:01
综合看Anker超能充系列的性价比很高,并且与不仅和iPhone12/苹果<em>Mac</em>Book很配,而且适合多设备充电需求的日常使用或差旅场景,不管是安卓还是Switch同样也能用得上它,希望这次分享能给准备购入充电器的小伙伴们有所
2021-06-01 09:31:42
除了L4WUDU与吴亦凡已经多次共事,成为了明面上的厂牌成员,吴亦凡还曾带领20XXCLUB全队参加2020年的一场音乐节,这也是20XXCLUB首次全员合照,王嗣尧Turbo、陈彦希Regi、<em>Mac</em> Ova Seas、林渝植等人全部出场。然而让
2021-06-01 09:31:34
目前应用IPFS的机构:1 谷歌<em>浏览器</em>支持IPFS分布式协议 2 万维网 (历史档案博物馆)数据库 3 火狐<em>浏览器</em>支持 IPFS分布式协议 4 EOS 等数字货币数据存储 5 美国国会图书馆,历史资料永久保存在 IPFS 6 加
2021-06-01 09:31:24
开拓者的车机是兼容苹果和<em>安卓</em>,虽然我不怎么用,但确实兼顾了我家人的很多需求:副驾的门板还配有解锁开关,有的时候老婆开车,下车的时候偶尔会忘记解锁,我在副驾驶可以自己开门:第二排设计很好,不仅配置了一个很大的
2021-06-01 09:30:48
不仅是<em>安卓</em>手机,苹果手机的降价力度也是前所未有了,iPhone12也“跳水价”了,发布价是6799元,如今已经跌至5308元,降价幅度超过1400元,最新定价确认了。iPhone12是苹果首款5G手机,同时也是全球首款5nm芯片的智能机,它
2021-06-01 09:30:45