起因:在实际项目开发过程中,需要使用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);
}
}
不要以为每天把功能完成了就行了,这种思想是要不得的,互勉~!