首页技术栈归档照片墙音乐日记随想收藏夹友链留言关于

Spring Boot使用消息队列

写作时间:2026-07-08

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";
}
avatar

yuanyourdomain

写代码,做研究,记录生活。

RECOMMENDED

MyBatis 动态 SQL

2026-07-08

Nginx 基础入门

2026-07-08

Maven 多模块与私服

2026-07-08