RabbitMQ队列延迟
1. 场景:
“订单下单成功后,15分钟未支付自动取消”
1.传统处理超时订单
采取定时任务轮训数据库订单,并且批量处理。其弊端也是显而易见的;对服务器、数据库性会有很大的要求,
并且当处理大量订单起来会很力不从心,而且实时性也不是特别好。当然传统的手法还可以再优化一下,
即存入订单的时候就算出订单的过期时间插入数据库,设置定时任务查询数据库的时候就只需要查询过期了的订单,
然后再做其他的业务操作
2.rabbitMQ延时队列方案
一台普通的rabbitmq服务器单队列容纳千万级别的消息还是没什么压力的,而且rabbitmq集群扩展支持的也是非常好的,
并且队列中的消息是可以进行持久化,即使我们重启或者宕机也能保证数据不丢失
2. TTL和DLX
rabbitMQ中是没有延时队列的,也没有属性可以设置,只能通过死信交换器(DLX)和设置过期时间(TTL)结合起来实现延迟队列
1.TTL
TTL是Time To Live的缩写, 也就是生存时间。
RabbitMq支持对消息和队列设置TTL,对消息这设置是在发送的时候指定,对队列设置是从消息入队列开始计算, 只要超过了队列的超时时间配置, 那么消息会自动清除。
如果两种方式一起使用消息对TTL和队列的TTL之间较小的为准,也就是消息5s过期,队列是10s,那么5s的生效。
默认是没有过期时间的,表示消息没有过期时间;如果设置为0,表示消息在投递到消费者的时候直接被消费,否则丢弃。
设置消息的过期时间用 x-message-ttl 参数实现,单位毫秒。
设置队列的过期时间用 x-expires 参数,单位毫秒,注意,不能设置为0。
2.DLX和死信队列
DLX即Dead-Letter-Exchange(死信交换机),它其实就是一个正常的交换机,能够与任何队列绑定。
死信队列是指队列(正常)上的消息(过期)变成死信后,能够后发送到另外一个交换机(DLX),然后被路由到一个队列上,
这个队列,就是死信队列
成为死信一般有以下几种情况:
消息被拒绝(basic.reject or basic.nack)且带requeue=false参数
消息的TTL-存活时间已经过期
队列长度限制被超越(队列满)
注1:如果队列上存在死信, RabbitMq会将死信消息投递到设置的DLX上去 ,
注2:通过在队列里设置x-dead-letter-exchange参数来声明DLX,如果当前DLX是direct类型还要声明
x-dead-letter-routing-key参数来指定路由键,如果没有指定,则使用原队列的路由键
3. 延迟队列
通过DLX和TTL模拟出延迟队列的功能,即,消息发送以后,不让消费者拿到,而是等待过期时间,变成死信后,发送给死信交换机再路由到死信队列进行消费
注1:延迟队列(即死信队列)产生流程见“images/01 死信队列产生流程.png”
4. 开发步骤
1.生产者创建一个正常消息,并添加消息过期时间/死信交换机/死信路由键这3个参数
关键代码1
new Queue(name, durable, exclusive, autoDelete, arguments);
new Queue(NORMAL_QUEUE, true, false, false, map)
参数说明:
name:队列名字
durable:true则持久队列
exclusive:如果我们声明一个排他队列(该队列将仅由声明者的连接使用),则为true
autoDelete:服务器不再使用时应删除队列,则为true
arguments:用于声明队列的参数
map.put("x-message-ttl", 10000);//message在该队列queue的存活时间最大为10秒
map.put("x-dead-letter-exchange", DELAY_EXCHANGE); //x-dead-letter-exchange参数是设置该队列的死信交换器(DLX)
map.put("x-dead-letter-routing-key", DELAY_ROUTING_KEY);//x-dead-letter-routing-key参数是给这个DLX指定路由键
关键代码2
new DirectExchange(NORMAL_EXCHANGE, true, false);
2.消费者A
正常情况下,由消费者A去消费队列“normal-queue”中的消息,但实际上没有,而是等消息过期
3.消费者B
消息过期后,变成死信,根据配置会被投递到DLX,然后根据死信路由键投到死信队列(即延时队列)中
5. 子模块间共享Model
1.创建公共子模块common
添加公共的JavaBean对象,并使用lombok简化代码
@Data:会为类的所有属性自动生成setter/getter、equals、canEqual、hashCode、toString方法
@NoArgsConstructor:无参构造器
@AllArgsConstructor:全参构造器
2.主模块
<packaging>pom</packaging>
<!-- 2.添加子模块 --> <modules> <module>rabbitmq-provider</module> <module>rabbitmq-consumer</module> <module>common</module> </modules> 3.各子模块 <!-- 1.packaging模式改为jar --> <packaging>jar</packaging>4.配置公共common模块
在主模块的POM的<dependencies>中添加公共子模块common
看代码
创建一个工程rabbitmq03 ,普通maven项目
pom.xml
1 1 <?xml version="1.0" encoding="UTF-8"?> 2 2 <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" 3 3 xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> 4 4 <modelVersion>4.0.0</modelVersion> 5 5 <parent> 6 6 <groupId>org.springframework.boot</groupId> 7 7 <artifactId>spring-boot-starter-parent</artifactId> 8 8 <version>2.2.2.RELEASE</version> 9 9 <relativePath/> <!-- lookup parent from repository --> 1010 </parent> 1111 <groupId>com.yuan</groupId> 1212 <artifactId>rabbitmq03</artifactId> 1313 <version>0.0.1-SNAPSHOT</version> 1414 <name>rabbitmq03</name> 1515 <packaging>pom</packaging> 1616 <description>Demo project for Spring Boot</description> 1717 1818 <properties> 1919 <java.version>1.8</java.version> 2020 </properties> 2121 2222 <modules> 2323 <module>rabbitmq-provider</module> 2424 <module>rabbitmq-consumer</module> 2525 </modules> 2626 2727 <dependencies> 2828 <dependency> 2929 <groupId>org.springframework.boot</groupId> 3030 <artifactId>spring-boot-starter-amqp</artifactId> 3131 </dependency> 3232 <dependency> 3333 <groupId>junit</groupId> 3434 <artifactId>junit</artifactId> 3535 <scope>test</scope> 3636 </dependency> 3737 <dependency> 3838 <groupId>org.springframework.boot</groupId> 3939 <artifactId>spring-boot-starter-web</artifactId> 4040 </dependency> 4141 4242 <dependency> 4343 <groupId>org.projectlombok</groupId> 4444 <artifactId>lombok</artifactId> 4545 <version>1.18.10</version> 4646 <scope>provided</scope> 4747 </dependency> 4848 4949 </dependencies> 5050 5151 <build> 5252 <plugins> 5353 <plugin> 5454 <groupId>org.springframework.boot</groupId> 5555 <artifactId>spring-boot-maven-plugin</artifactId> 5656 </plugin> 5757 </plugins> 5858 </build> 5959 6060 </project>
创建生产者模块rabbitmq-provider

pom.xml
1 1 <?xml version="1.0" encoding="UTF-8"?> 2 2 <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" 3 3 xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> 4 4 <modelVersion>4.0.0</modelVersion> 5 5 6 6 <parent> 7 7 <groupId>com.yuan</groupId> 8 8 <artifactId>rabbitmq03</artifactId> 9 9 <version>0.0.1-SNAPSHOT</version> 1010 </parent> 1111 <artifactId>rabbitmq-provider</artifactId> 1212 <version>0.0.1-SNAPSHOT</version> 1313 <name>rabbitmq-provider</name> 1414 <description>子模块-生产者</description> 1515 <packaging>jar</packaging> 1616 </project>
QueueDelayConfig
1 1 package com.yuan.rabbitmqprovider.rabbitmq; 2 2 3 3 4 4 import org.springframework.amqp.core.Binding; 5 5 import org.springframework.amqp.core.BindingBuilder; 6 6 import org.springframework.amqp.core.DirectExchange; 7 7 import org.springframework.amqp.core.Queue; 8 8 import org.springframework.context.annotation.Bean; 9 9 import org.springframework.context.annotation.Configuration; 1010 1111 import javax.lang.model.element.NestingKind; 1212 import java.util.HashMap; 1313 import java.util.Map; 1414 1515 @Configuration 1616 public class QueueDelayConfig { 1717 1818 /** 1919 * 定义正常的队列、交换机、路由键 2020 */ 2121 public static final String NORMAL_QUEUE="normal-queue"; 2222 public static final String NORMAL_EXCHANGE="normal-exchange"; 2323 public static final String NORMAL_ROUTINGKEY="normal-routingkey"; 2424 2525 /** 2626 * 定义死信的队列、交换机、路由键 2727 */ 2828 public static final String DELAY_QUEUE="delay-queue"; 2929 public static final String DELAY_EXCHANGE="delay-exchange"; 3030 public static final String DELAY_ROUTINGKEY="delay-routingkey"; 3131 3232 3333 /** 3434 * 定义正常队列 3535 * @return 3636 */ 3737 @Bean 3838 public Queue normalQueue(){ 3939 //设定消息过期时间/死信交换机/死信路由键3个参数 4040 Map<String, Object> map = new HashMap<String, Object>(); 4141 map.put("x-message-ttl", 15000);//message在该队列queue的存活时间最大为15秒 4242 map.put("x-dead-letter-exchange", DELAY_EXCHANGE); //x-dead-letter-exchange参数是设置该队列的死信交换器(DLX) 4343 map.put("x-dead-letter-routing-key", DELAY_ROUTINGKEY);//x-dead-letter-routing-key参数是给这个DLX指定路由键 4444 4545 return new Queue(NORMAL_QUEUE, true, false, false, map); 4646 } 4747 4848 @Bean 4949 public DirectExchange normalExchange(){ 5050 return new DirectExchange(NORMAL_EXCHANGE, true, false); 5151 } 5252 5353 @Bean 5454 public Binding normalRoutingkey(){ 5555 return BindingBuilder.bind(normalQueue()) 5656 .to(normalExchange()) 5757 .with(NORMAL_ROUTINGKEY); 5858 } 5959 6060 6161 /** 6262 * 定义死信队列 6363 */ 6464 @Bean 6565 public Queue delayQueue(){ 6666 return new Queue(DELAY_QUEUE, true); 6767 } 6868 6969 @Bean 7070 public DirectExchange delayExchange(){ 7171 return new DirectExchange(DELAY_EXCHANGE); 7272 } 7373 7474 @Bean 7575 public Binding delayRoutingkey(){ 7676 return BindingBuilder.bind(delayQueue()) 7777 .to(delayExchange()) 7878 .with(DELAY_ROUTINGKEY); 7979 } 8080 8181 8282 8383 8484 }
SendController
1 1 package com.yuan.rabbitmqprovider.controller; 2 2 3 3 4 4 import com.yuan.rabbitmqprovider.rabbitmq.QueueDelayConfig; 5 5 import lombok.extern.slf4j.Slf4j; 6 6 import org.springframework.amqp.rabbit.core.RabbitTemplate; 7 7 import org.springframework.beans.factory.annotation.Autowired; 8 8 import org.springframework.web.bind.annotation.RequestMapping; 9 9 import org.springframework.web.bind.annotation.RestController; 1010 1111 import java.time.LocalDateTime; 1212 import java.time.format.DateTimeFormatter; 1313 import java.util.HashMap; 1414 import java.util.Map; 1515 1616 @RestController 1717 @Slf4j 1818 public class SendController { 1919 2020 @Autowired 2121 private RabbitTemplate rabbitTemplate; 2222 2323 @RequestMapping("/sender") 2424 public Map<String, Object> sender(){ 2525 Map<String, Object> data = this.createData(); 2626 2727 rabbitTemplate.convertAndSend(QueueDelayConfig.NORMAL_EXCHANGE, 2828 QueueDelayConfig.NORMAL_ROUTINGKEY,data); 2929 Map<String, Object> result = new HashMap<String, Object>(); 3030 result.put("msg","OK"); 3131 result.put("code","1"); 3232 return result; 3333 } 3434 3535 3636 3737 private Map<String, Object> createData(){ 3838 Map<String, Object> map = new HashMap<String, Object>(); 3939 4040 String date = LocalDateTime.now().format(DateTimeFormatter.BASIC_ISO_DATE. 4141 ofPattern("yyyy-MM-dd HH:mm:ss")); 4242 map.put("msg","hello rabbitmq!!"); 4343 map.put("success",true); 4444 map.put("createdate", date); 4545 4646 4747 return map; 4848 } 4949 5050 5151 5252 }
最后配置一下yml文件
1 1 server: 2 2 port: 8081 3 3 servlet: 4 4 context-path: /rabbitmq-provider 5 5 spring: 6 6 rabbitmq: 7 7 virtual-host: / 8 8 username: guest 9 9 password: guest 1010 host: 192.168.238.129 1111 port: 5672
创建消费者模块rabbitmq-consumer
pom.xml
1 1 <?xml version="1.0" encoding="UTF-8"?> 2 2 <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" 3 3 xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> 4 4 <modelVersion>4.0.0</modelVersion> 5 5 <parent> 6 6 <groupId>com.yuan</groupId> 7 7 <artifactId>rabbitmq03</artifactId> 8 8 <version>0.0.1-SNAPSHOT</version> 9 9 </parent> 1010 <artifactId>rabbitmq-consumer</artifactId> 1111 <version>0.0.1-SNAPSHOT</version> 1212 <name>rabbitmq-consumer</name> 1313 <description>子模块-消费者</description> 1414 <packaging>jar</packaging> 1515 </project>
QueueRecevier
1 1 package com.yuan.rabbitmqconsumer.controller; 2 2 3 3 import lombok.extern.slf4j.Slf4j; 4 4 import org.springframework.amqp.rabbit.annotation.RabbitHandler; 5 5 import org.springframework.amqp.rabbit.annotation.RabbitListener; 6 6 import org.springframework.stereotype.Component; 7 7 8 8 import java.util.Map; 9 9 1010 @Component 1111 @Slf4j 1212 @RabbitListener(queues = {"delay-queue"}) //消费端监听队列,如果delay-queue死信队列中有消息过来就会被消费掉 1313 public class QueueRecevier { 1414 1515 @RabbitHandler 1616 public void handlerMessage(Map<String, Object> data){ 1717 log.info("QueueRecevier.handlerMessage,data={}",data); 1818 } 1919 2020 2121 2222 2323 }
标红处的log使用需要下载一个插件Lombok

直接右边install, 然后重启idea
yml文件配置
1 1 server: 2 2 port: 8082 3 3 servlet: 4 4 context-path: /rabbitmq-consumer 5 5 spring: 6 6 rabbitmq: 7 7 virtual-host: / 8 8 username: guest 9 9 password: guest 1010 host: 192.168.238.129 1111 port: 5672
启动生产者,访问http://localhost:8081/rabbitmq-provider/sender 发送请求。

生产端推送消息到正常队列等待被消费,我们设定的过期时间是15秒,,,


启动消费端,消费端会根据我们设定的监听去监听队列中是否有消息有则会被消费掉。。

6. json转换
1.生产者
@Bean
public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory, Jackson2JsonMessageConverter jackson2JsonMessageConverter) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
rabbitTemplate.setMessageConverter(jackson2JsonMessageConverter);//指定json转换器
return rabbitTemplate;
}
2.消费者
@Bean
public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(ConnectionFactory connectionFactory, Jackson2JsonMessageConverter jackson2JsonMessageConverter) {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory);
factory.setMessageConverter(jackson2JsonMessageConverter);
return factory; }
创建公共子模块common-vo
1 1 <?xml version="1.0" encoding="UTF-8"?> 2 2 <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" 3 3 xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> 4 4 <modelVersion>4.0.0</modelVersion> 5 5 <parent> 6 6 <groupId>com.yuan</groupId> 7 7 <artifactId>rabbitmq03</artifactId> 8 8 <version>0.0.1-SNAPSHOT</version> 9 9 </parent> 1010 <artifactId>common-vo</artifactId> 1111 <version>0.0.1-SNAPSHOT</version> 1212 <name>common-vo</name> 1313 <packaging>jar</packaging> 1414 <description>公共子模块</description> 1515 1616 1717 1818 </project>
创建一个model的Package,创建一个Order
1package com.yuan.commonvo.model; 2 3import lombok.Data; 4 5import java.lang.reflect.ParameterizedType; 6import java.util.Date; 7 8 9@Data 10public class Order { 11 12 private long orderId; 13 private String orderNo; 14 private Date createdate; 15 16 17}
vo包下创建一个OrderVo
1package com.yuan.commonvo.vo; 2 3import com.yuan.commonvo.model.Order; 4 5 6public class OrderVo extends Order { 7 8}
完了之后在父模块中添加common-vo子模块的一个pom依赖
1<modules> 2 <module>rabbitmq-provider</module> 3 <module>rabbitmq-consumer</module> 4 <module>common-vo</module> 5 </modules> 6 7<dependency> <groupId>com.yuan</groupId> <artifactId>common-vo</artifactId> <version>0.0.1-SNAPSHOT</version></dependency>
修改生产者SendController
1@RequestMapping("/sender") 2 public Map<String, Object> sender(){ 3// Map<String, Object> data = this.createData(); 4 5 OrderVo orderVo = new OrderVo(); 6 orderVo.setOrderId(1); 7 orderVo.setOrderNo("P001"); 8 9 rabbitTemplate.convertAndSend(QueueDelayConfig.NORMAL_EXCHANGE, 10 QueueDelayConfig.NORMAL_ROUTINGKEY,orderVo); 11 Map<String, Object> result = new HashMap<String, Object>(); 12 result.put("msg","OK"); 13 result.put("code","1"); 14 return result; 15 }
添加QueueProviderMessageConvert
package com.yuan.rabbitmqprovider.rabbitmq;import org.springframework.amqp.rabbit.connection.ConnectionFactory;import org.springframework.amqp.rabbit.core.RabbitTemplate;import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;import org.springframework.context.annotation.Bean;import org.springframework.context.annotation.Configuration;@Configurationpublic class QueueProviderMessageConvert { @Bean public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory){ RabbitTemplate rabbitTemplate=new RabbitTemplate(); rabbitTemplate.setConnectionFactory(connectionFactory); rabbitTemplate.setMessageConverter(jackson2JsonMessageConverter()); return rabbitTemplate; } @Bean public Jackson2JsonMessageConverter jackson2JsonMessageConverter(){ return new Jackson2JsonMessageConverter(); }}
修改消费端QueueRecevier
1package com.yuan.rabbitmqconsumer.controller; 2 3 4import com.yuan.commonvo.vo.OrderVo; 5import lombok.extern.slf4j.Slf4j; 6import org.springframework.amqp.rabbit.annotation.RabbitHandler; 7import org.springframework.amqp.rabbit.annotation.RabbitListener; 8import org.springframework.stereotype.Component; 9 10@Component 11@Slf4j 12@RabbitListener(queues = {"delay-queue"}) //消费端监听队列,如果delay-queue死信队列中有消息过来就会被消费掉 13public class QueueRecevier { 14 15 @RabbitHandler 16 public void handlerMessage(OrderVo orderVo){ 17 log.info("QueueRecevier.handlerMessage,data={}",orderVo); 18 } 19 20 21 22 23}
添加消费端QueueRecevierMessageConvert
package com.yuan.rabbitmqconsumer.controller;import org.springframework.amqp.rabbit.connection.ConnectionFactory;import org.springframework.amqp.rabbit.core.RabbitTemplate;import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;import org.springframework.context.annotation.Bean;import org.springframework.context.annotation.Configuration;@Configurationpublic class QueueRecevierMessageConvert { @Bean public RabbitTemplate createRabbitTemplate(ConnectionFactory connectionFactory){ RabbitTemplate rabbitTemplate=new RabbitTemplate(); rabbitTemplate.setConnectionFactory(connectionFactory); rabbitTemplate.setMessageConverter(jackson2JsonMessageConverter()); return rabbitTemplate; } @Bean public Jackson2JsonMessageConverter jackson2JsonMessageConverter(){ return new Jackson2JsonMessageConverter(); }}
测试:
