RabbitMQ 延时插件实现消息延时发送

起因:在实际项目开发过程中,需要使用RabbitMQ来实现消息队列的功能,但仅仅实现功能之后并不能对自己满足,既然学一次,就要更深的了解她,吃一吃架构方面的相关内容,提升自己。


1. 上下游数据的消费是否互相影响

上游同时接收到的数据消费后不影响下游发送过过去的信息

2. 消息延迟发送机制的实现

delay的应用场景:用户下了个单,30分钟后如果没有支付,我就将订单关闭,5分钟一轮询(29分钟+5)

2.1. 通过死信队列来实现

在消息发送的时候设置消息的TTL,并将该消息发送到一个没有人消费的队列上,将这个没有人消费的队列配置成死信触发队列:x-dead-letter-exchange、x-dead-letter-routing-key 当消息超过TTL后转发发给一个具体的执行队列,这个执行队列的消息需要监听和消费,当消息一进来就消费掉,这个消息的TTL就是delay的时长

2.2. 通过延时插件实现消息延时发送

2.2.1. 延时插件的配置

# 这个插件默认是不带的,需要下载,需要确保大版本和rabbitmq一致
wget https://dl.bintray.com/rabbitmq/community-plugins/3.6.x/rabbitmq_delayed_message_exchange/rabbitmq_delayed_message_exchange-20171215-3.6.x.zip
unzip rabbitmq_delayed_message_exchange-20171215-3.6.x.zip
mv rabbitmq_delayed_message_exchange-20171215-3.6.x.ez /usr/lib/rabbitmq/lib/rabbitmq_server-3.6.5/plugins/
# 解压并移动到plugins目录后先看一下是否成功,插件是热加载的,不用停服务
rabbitmq-plugins list
rabbitmq-plugins enable rabbitmq_delayed_message_exchange

应用完毕后看是否加载成功,去控制台看add exchange,是否有延迟交换机了

注意:如果是集群需要所有机器都加载这个插件

2.2.2. 创建延时交换机

出现的提示意思是需要你通过arguments来指定延迟交换机的type匹配类型

2.2.3. 创建延时消息

20秒后数据是能进入队列的,为什么提示not routed?

因为是延时队列,是发送到exchange成功了,但还没有到时间的时候exchange没有进行routing操作

3. springboot实现延时信息的收发

接收方

import com.icoding.basic.po.OrderInfo;
import com.rabbitmq.client.Channel;
import org.springframework.amqp.rabbit.annotation.*;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Headers;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.stereotype.Component;
import java.util.Map;
​
@Component
public class OrderReceiver {
​
 int flag = 0;
​
 @RabbitListener(bindings = @QueueBinding(
 value = @Queue(value = "delay-queue-other",durable = "true",autoDelete = "false"),
 exchange = @Exchange(value = "delay-exchange-other",durable = "true",type = "x-delayed-message",arguments = {
 @Argument(name = "x-delayed-type",value = "topic")
 }),
 key = "delay.#"
 )
 )
 @RabbitHandler
 public void onOrderMessage(@Payload OrderInfo orderInfo, @Headers Map<String,Object> headers, Channel channel) throws Exception{
 System.out.println("************消息接收开始***********");
 System.out.println("Order Name: "+orderInfo.getOrder_name());
 Long deliverTag = (Long)headers.get(AmqpHeaders.DELIVERY_TAG);
 //ACK进行签收,第一个参数是标识,第二个参数是批量接收为fasle
 //channel.basicAck(deliverTag,false);
 //前两个参数和上面一样,第三个参数是否重回队列
 channel.basicAck(deliverTag,false);
 }
}

发送方

import com.icoding.basic.po.OrderInfo;
import org.springframework.amqp.AmqpException;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessagePostProcessor;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
​
@Component
public class OrderSender {
​
 @Autowired
 RabbitTemplate rabbitTemplate;
​
 public void sendOrder(OrderInfo orderInfo) throws Exception{
 /**
 * exchange: 交换机名字
 * routingkey: 队列关联的key
 * object: 要传输的消息对象
 * correlationData: 消息的唯一id
 */
 CorrelationData correlationData = new CorrelationData();
 correlationData.setId(orderInfo.getMessage_id());
 MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
 @Override
 public Message postProcessMessage(Message message) throws AmqpException {
 message.getMessageProperties().getHeaders().put("x-delay",20000);
 return message;
 }
 };
 rabbitTemplate.convertAndSend("delay-exchange-other","delay.key",orderInfo,messagePostProcessor,correlationData);
 }
}

不要以为每天把功能完成了就行了,这种思想是要不得的,互勉~!

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。