Spring Boot使用消息队列
使用@RabbitListener进行简单的消息监听
依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
将hello-java-queue绑定一个交换机hello-java-exchange指定他的Routing key为hello2.java

application.yaml
#mq的基本配置
spring:
rabbitmq:
host: 120.25.213.140
port: 5672
virtual-host: /
username: admin
password: admin
生产者
@Test
public void sendMessageTest() {
OrderReturnReasonEntity reasonEntity = new OrderReturnReasonEntity();
reasonEntity.setId(1L);
reasonEntity.setCreateTime(new Date());
reasonEntity.setName("reason");
reasonEntity.setStatus(1);
reasonEntity.setSort(2);
String msg = "Hello World";
//1、发送消息,如果发送的消息是个对象,会使用序列化机制,将对象写出去,对象必须实现Serializable接口
//2、发送的对象类型的消息,可以是一个json
//分别对应为交换机,routingkey,发送消息的类型,和一个唯一id
rabbitTemplate.convertAndSend("hello-java-exchange","hello2.java",
reasonEntity,new CorrelationData(UUID.randomUUID().toString()));
log.info("消息发送完成:{}",reasonEntity);
}
消费者
@RabbitListener(queues = {"hello-java-queue"}) //指定接收消息的队列名称hello-java-queue1
public void reciveMessage(Message message, //message 原生的详细信息。头+体
OrderReturnReasonEntity content, // T<发送的消息的类型> OrderReturnReasonEntity content 就是发送消息时间的自己传的类型这里是实体类
Channel channel){ //channel 当前消息的传输通道
System.out.println("接收到的消息:"+message.getBody());
}
在启动类加上@EnableRabbit注解
使用@RabbitHandler接收多种类型的消息
业务层(消费者)
@RabbitListener(queues = {"hello-java-queue"}) //指定接收消息的队列名称hello-java-queue1
@Service("orderItemService")
public class OrderItemServiceImpl extends ServiceImpl<OrderItemDao, OrderItemEntity> implements OrderItemService {
@Resource
private RabbitTemplate rabbitTemplate;
@RabbitHandler //接收指定的消息
public void reciveMessage(Message message, //message 原生的详细信息。头+体
OrderReturnReasonEntity content, // T<发送的消息的类型> OrderReturnReasonEntity content 就是发送消息时间的自己传的类型这里是实体类
Channel channel){ //channel 当前消息的传输通道
System.out.println("接收到的消息:"+content);
}
@RabbitHandler //接收指定的消息
public void orderMessage(OrderEntity content){ //指定接收消息的类型
System.out.println("接收到的消息:"+content);
}
}
视图层(生产者)
@RestController
@Slf4j
public class RabbitmqController {
@Resource
RabbitTemplate rabbitTemplate;
@GetMapping("rabbitmq")
public String rabbitmq(@RequestParam(value = "num",defaultValue = "10",required = false) Integer num) {
for (int i = 0; i < num; i++) {
if (i % 2 == 0) {
//1、发送消息,如果发送的消息是个对象,会使用序列化机制,将对象写出去,对象必须实现Serializable接口
OrderReturnReasonEntity reasonEntity = new OrderReturnReasonEntity();
reasonEntity.setId(1L);
reasonEntity.setCreateTime(new Date());
reasonEntity.setName("reason");
reasonEntity.setStatus(1);
reasonEntity.setSort(2);
//2、发送的对象类型的消息,可以是一个json
//分别对应为交换机,routingkey,发送消息的类型,和一个唯一id
rabbitTemplate.convertAndSend("hello-java-exchange", "hello2.java",
reasonEntity,new CorrelationData(UUID.randomUUID().toString()));
log.info("消息发送完成:{}", reasonEntity);
} else {
OrderEntity orderEntity = new OrderEntity();
orderEntity.setOrderSn(UUID.randomUUID().toString());
rabbitTemplate.convertAndSend("hello-java-exchange", "hello2.java",
orderEntity,new CorrelationData(UUID.randomUUID().toString()));
}
}
return "OK";
}
}
RabbitMq可靠投递机制(手动应答)
application.properties
#开启发布认
spring.rabbitmq.publisher-confirm-type=correlated
#开启发送消息抵达队列的确认
spring.rabbitmq.publisher-returns=true
#只有抵达队列以异步的方式优先回调
spring.rabbitmq.template.mandatory=true
#手动ack消息,手动应答
spring.rabbitmq.listener.simple.acknowledge-mode=manual
MyRabbitConfig 配置类
/**
* mq配置类
*/
@Configuration
public class MyRabbitConfig {
@Resource
private RabbitTemplate rabbitTemplate;
/**RabbitMQ以JSON进行转换*/
// @Bean
// public MessageConverter messageConverter() {
// return new Jackson2JsonMessageConverter();
// }
/**
* 定制RabbitTemplate
* 1、服务收到消息就会回调
* 1、spring.rabbitmq.publisher-confirm-type=correlated
* 2、设置确认回调
* 2、消息正确抵达队列就会进行回调
* 1、spring.rabbitmq.publisher-returns: true
* spring.rabbitmq.template.mandatory: true
* 2、设置确认回调ReturnCallback
*
* 3、消费端确认(保证每个消息都被正确消费,此时才可以broker删除这个消息)
* 1、默认是自动确认的,只要消息接收到,客户端会自动确认,服务端就会移除这个消息
* 间题:
* 我们收到很多消息,自动回复给服务器ack,只有一个消息处理成功,宕机了。发生消息丢失;
* 手动确认模式。只要我们没有明确告诉,货物被签收。没有Ack。消息就一直是unached状态。即使Consumer宕机。
* 消息不会丢失,会重新变为Ready,下一次有新的consumer连接进来就发给他
* 2、如何签收:
* channel.basicAcl( deliveryTag,false);签收;业务成功完成就应该签收
* channeL.basicNack(deliveryTag,false,true);拒签;业务失败,拒签
*/
@PostConstruct //MyRabbitConfig对象创建完成以后,执行这个方法
public void initRabbitTemplate() {
/**
* 消息成功回调
* 1、只要消息抵达Broker就ack=true
* correlationData:当前消息的唯一关联数据(这个是消息的唯一id)
* ack:消息是否成功收到
* cause:失败的原因
*/
//设置确认回调
rabbitTemplate.setConfirmCallback((correlationData,ack,cause) -> {
System.out.println("回调成功的消息confirm...correlationData["+correlationData+"]==>ack:["+ack+"]==>cause:["+cause+"]");
});
/**
* 消息失败回调
* 只要消息没有投递给指定的队列,就触发这个失败回调
* message:投递失败的消息详细信息
* replyCode:回复的状态码
* replyText:回复的文本内容
* exchange:当时这个消息发给哪个交换机
* routingKey:当时这个消息用哪个路邮键
*/
rabbitTemplate.setReturnCallback((message,replyCode,replyText,exchange,routingKey) -> {
System.out.println("回调失败的消息"+"Fail Message["+message+"]==>replyCode["+replyCode+"]" +
"==>replyText["+replyText+"]==>exchange["+exchange+"]==>routingKey["+routingKey+"]");
});
}
}
server(消费者)
@RabbitListener(queues = {"hello-java-queue"}) //指定接收消息的队列名称hello-java-queue1
@Service("orderItemService")
public class OrderItemServiceImpl extends ServiceImpl<OrderItemDao, OrderItemEntity> implements OrderItemService {
@Resource
private RabbitTemplate rabbitTemplate;
// @RabbitListener(queues = {"hello-java-queue"}) //指定接收消息的队列名称hello-java-queue1
@RabbitHandler //接收指定的消息
public void reciveMessage(Message message, //message 原生的详细信息。头+体
OrderReturnReasonEntity content, // T<发送的消息的类型> OrderReturnReasonEntity content 就是发送消息时间的自己传的类型这里是实体类
Channel channel){ //channel 当前消息的传输通道
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
if (deliveryTag%2==0) {
//签收货物,false单个确认
channel.basicAck(deliveryTag, false);// 手动ack模式
System.out.println("签收了货物:"+deliveryTag);
}else {
//退货
//参数 ,退货的变化,是否批量退货false单个当前的true全部退货,是否重新加入队列true发回服务器,flase直接丢掉
channel.basicNack(deliveryTag,false,true);
//消息拒绝以后重新放到队列里面,让别人继续进行消费解锁
//channel.basicReject(message.getMessageProperties().getDeliveryTag(), true);
System.out.println("没有签收了货物:"+deliveryTag);
}
} catch (Exception e) {
e.printStackTrace();
}
System.out.println("接收到的消息:"+content);
}
@RabbitHandler //接收指定的消息
public void orderMessage(OrderEntity content){ //指定接收消息的类型
System.out.println("接收到的消息:"+content);
}
}
RabbitMQ模拟业务使用延迟队列(基于死性)
配置类
package com.yuan.gulimall.order.config;
import com.rabbitmq.client.Channel;
import com.yuan.gulimall.order.entity.OrderEntity;
import org.springframework.amqp.core.*;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.io.IOException;
import java.util.HashMap;
@Configuration
public class MyMQConfig {
/**
* 监听过期的消息
* @param entity
*/
@RabbitListener(queues = "order.release.order.queue")
public void listen(OrderEntity entity, Channel channel, Message message) throws IOException {
System.out.println("1分钟收到过期的订单信息:准备关闭订单"+entity.getOrderSn());
channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
}
/* @Bean作用容器中的Queue、Exchange、Binding 会自动创建(在RabbitMQ)不存在的情况下 */
@Bean
public Queue orderDelayQueue(){ //死信队列
/*
Queue(String name, 队列名字
boolean durable, 是否持久化
boolean exclusive, 是否排他
boolean autoDelete, 是否自动删除
Map<String, Object> arguments) 属性
*/
HashMap<String, Object> arguments = new HashMap<>();
arguments.put("x-dead-letter-exchange", "order-event-exchange"); //指定死信路由交换机
arguments.put("x-dead-letter-routing-key", "order.release.order");//指定死性路由key
arguments.put("x-message-ttl", 60000); // 消息过期时间 1分钟
Queue queue = new Queue("order.delay.queue", true, false, false, arguments);
return queue;
}
/**
* 普通队列
* @return
*/
@Bean
public Queue orderReleaseOrderQueue(){
Queue queue = new Queue("order.release.order.queue", true, false, false);
return queue;
}
/**
* 交换机
* @return
*/
@Bean
public Exchange orderEventExchange(){
/*
* String name, 交换机名称
* boolean durable, 是否持久化
* boolean autoDelete, 是否自动删除
* Map<String, Object> arguments 属性
* */
return new TopicExchange("order-event-exchange", true, false);
}
/**
* 绑定关系
* @return
*/
@Bean
public Binding orderCreateOrderBingding(){
/*
* String destination, 目的地(队列名或者交换机名字)
* DestinationType destinationType, 目的地类型(Queue、Exhcange)
* String exchange,
* String routingKey,
* Map<String, Object> arguments
* */
return new Binding("order.delay.queue", //指的是和哪个队列进行绑定
Binding.DestinationType.QUEUE, //绑定的类型为一个队列
"order-event-exchange", // 指定哪个交换机和这个目的地进行绑定
"order.create.order", //指定路由key
null); //绑定的熟悉
}
@Bean
public Binding orderReleaseOrderBingding(){
return new Binding("order.release.order.queue",
Binding.DestinationType.QUEUE,
"order-event-exchange",
"order.release.order",
null);
}
}
测试类
@ResponseBody
@GetMapping(value = "/test/createOrder")
public String createOrderTest() {
//订单下单成功
OrderEntity orderEntity = new OrderEntity();
orderEntity.setOrderSn(UUID.randomUUID().toString());
orderEntity.setModifyTime(new Date());
//给MQ发送消息将在一分钟之后接收到 order-event-exchange 指定绑定的交换机 order.create.order 指定绑定的路由key
rabbitTemplate.convertAndSend("order-event-exchange","order.create.order",orderEntity);
return "ok";
}