SpringBoot集成RabbitMQ收发消息

wen java案例 1

本文目录导读:

SpringBoot集成RabbitMQ收发消息

  1. 添加依赖
  2. 配置RabbitMQ
  3. 消息队列配置
  4. 消息发送
  5. 消息接收
  6. 回调配置
  7. 使用示例
  8. 重要配置说明

我来详细介绍SpringBoot集成RabbitMQ实现消息收发。

添加依赖

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

配置RabbitMQ

application.yml配置

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    virtual-host: /
    # 发送确认
    publisher-confirm-type: correlated
    # 发送回退
    publisher-returns: true
    # 手动ACK
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 1

消息队列配置

交换机、队列、绑定配置

@Configuration
public class RabbitConfig {
    // 1. 直连交换机
    public static final String DIRECT_EXCHANGE = "direct.exchange";
    public static final String DIRECT_QUEUE = "direct.queue";
    public static final String DIRECT_ROUTING_KEY = "direct.routing.key";
    // 2. 主题交换机
    public static final String TOPIC_EXCHANGE = "topic.exchange";
    public static final String TOPIC_QUEUE = "topic.queue";
    public static final String TOPIC_ROUTING_KEY = "topic.#";
    // 3. 广播交换机
    public static final String FANOUT_EXCHANGE = "fanout.exchange";
    public static final String FANOUT_QUEUE_A = "fanout.queue.a";
    public static final String FANOUT_QUEUE_B = "fanout.queue.b";
    // 创建直连交换机
    @Bean
    public DirectExchange directExchange() {
        return new DirectExchange(DIRECT_EXCHANGE, true, false);
    }
    // 创建主题交换机
    @Bean
    public TopicExchange topicExchange() {
        return new TopicExchange(TOPIC_EXCHANGE, true, false);
    }
    // 创建广播交换机
    @Bean
    public FanoutExchange fanoutExchange() {
        return new FanoutExchange(FANOUT_EXCHANGE, true, false);
    }
    // 创建队列
    @Bean
    public Queue directQueue() {
        return QueueBuilder.durable(DIRECT_QUEUE).build();
    }
    @Bean
    public Queue topicQueue() {
        return QueueBuilder.durable(TOPIC_QUEUE).build();
    }
    @Bean
    public Queue fanoutQueueA() {
        return QueueBuilder.durable(FANOUT_QUEUE_A).build();
    }
    @Bean
    public Queue fanoutQueueB() {
        return QueueBuilder.durable(FANOUT_QUEUE_B).build();
    }
    // 绑定关系
    @Bean
    public Binding directBinding() {
        return BindingBuilder.bind(directQueue())
                .to(directExchange())
                .with(DIRECT_ROUTING_KEY);
    }
    @Bean
    public Binding topicBinding() {
        return BindingBuilder.bind(topicQueue())
                .to(topicExchange())
                .with(TOPIC_ROUTING_KEY);
    }
    @Bean
    public Binding fanoutBindingA() {
        return BindingBuilder.bind(fanoutQueueA())
                .to(fanoutExchange());
    }
    @Bean
    public Binding fanoutBindingB() {
        return BindingBuilder.bind(fanoutQueueB())
                .to(fanoutExchange());
    }
}

消息发送

消息实体类

@Data
@NoArgsConstructor
@AllArgsConstructor
public class OrderMessage implements Serializable {
    private static final long serialVersionUID = 1L;
    private String orderId;
    private String userId;
    private BigDecimal amount;
    private LocalDateTime createTime;
}

消息发送者

@Component
@Slf4j
public class MessageSender {
    @Autowired
    private RabbitTemplate rabbitTemplate;
    /**
     * 发送直连交换机消息
     */
    public void sendDirectMessage(OrderMessage orderMessage) {
        CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
        rabbitTemplate.convertAndSend(
            RabbitConfig.DIRECT_EXCHANGE,
            RabbitConfig.DIRECT_ROUTING_KEY,
            orderMessage,
            correlationData
        );
        log.info("直接交换机消息发送成功: {}", orderMessage);
    }
    /**
     * 发送主题交换机消息
     */
    public void sendTopicMessage(String routingKey, Object message) {
        rabbitTemplate.convertAndSend(
            RabbitConfig.TOPIC_EXCHANGE,
            routingKey,
            message
        );
        log.info("主题交换机消息发送成功: routingKey={}, message={}", routingKey, message);
    }
    /**
     * 发送广播交换机消息
     */
    public void sendFanoutMessage(Object message) {
        rabbitTemplate.convertAndSend(
            RabbitConfig.FANOUT_EXCHANGE,
            "",
            message
        );
        log.info("广播交换机消息发送成功: {}", message);
    }
    /**
     * 延迟消息发送(需要安装延迟插件)
     */
    public void sendDelayMessage(Object message, long delayTime) {
        rabbitTemplate.convertAndSend(
            "delay.exchange",
            "delay.routing.key",
            message,
            msg -> {
                msg.getMessageProperties().setDelay(Long.valueOf(delayTime).intValue());
                return msg;
            }
        );
        log.info("延迟消息发送成功: {}", message);
    }
}

消息接收

手动ACK消息消费者

@Component
@Slf4j
public class MessageConsumer {
    /**
     * 直连队列消费者
     */
    @RabbitListener(queues = RabbitConfig.DIRECT_QUEUE)
    public void handleDirectMessage(OrderMessage message, 
                                   Message mqMessage, 
                                   Channel channel) throws IOException {
        long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
        try {
            log.info("接收到直连消息: {}", message);
            // 业务处理
            processOrder(message);
            // 手动ACK
            channel.basicAck(deliveryTag, false);
            log.info("消息处理成功, deliveryTag: {}", deliveryTag);
        } catch (Exception e) {
            log.error("消息处理失败: {}", e.getMessage());
            // 判断是否重试
            boolean isRetry = mqMessage.getMessageProperties()
                .getHeader("x-retry-count") != null && 
                (Integer) mqMessage.getMessageProperties().getHeader("x-retry-count") < 3;
            if (isRetry) {
                // 重试,重新入队
                channel.basicNack(deliveryTag, false, true);
            } else {
                // 丢弃或进入死信队列
                channel.basicNack(deliveryTag, false, false);
                log.warn("消息超过最大重试次数,已丢弃: {}", message);
            }
        }
    }
    /**
     * 主题队列消费者
     */
    @RabbitListener(queues = RabbitConfig.TOPIC_QUEUE)
    public void handleTopicMessage(Object message, 
                                  Message mqMessage, 
                                  Channel channel) throws IOException {
        long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
        try {
            log.info("接收到主题消息: {}, routingKey={}", 
                message, mqMessage.getMessageProperties().getReceivedRoutingKey());
            // 业务处理
            channel.basicAck(deliveryTag, false);
        } catch (Exception e) {
            channel.basicNack(deliveryTag, false, true);
        }
    }
    /**
     * 广播队列A消费者
     */
    @RabbitListener(queues = RabbitConfig.FANOUT_QUEUE_A)
    public void handleFanoutMessageA(Object message, 
                                    Message mqMessage, 
                                    Channel channel) throws IOException {
        long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
        try {
            log.info("广播队列A接收到消息: {}", message);
            channel.basicAck(deliveryTag, false);
        } catch (Exception e) {
            channel.basicNack(deliveryTag, false, true);
        }
    }
    /**
     * 广播队列B消费者
     */
    @RabbitListener(queues = RabbitConfig.FANOUT_QUEUE_B)
    public void handleFanoutMessageB(Object message, 
                                    Message mqMessage, 
                                    Channel channel) throws IOException {
        long deliveryTag = mqMessage.getMessageProperties().getDeliveryTag();
        try {
            log.info("广播队列B接收到消息: {}", message);
            channel.basicAck(deliveryTag, false);
        } catch (Exception e) {
            channel.basicNack(deliveryTag, false, true);
        }
    }
    private void processOrder(OrderMessage orderMessage) {
        // 业务处理逻辑
        log.info("处理订单: {}", orderMessage.getOrderId());
    }
}

回调配置

@Component
@Slf4j
public class RabbitTemplateConfig implements RabbitTemplate.ConfirmCallback, 
                                            RabbitTemplate.ReturnsCallback {
    @PostConstruct
    public void init() {
        rabbitTemplate.setConfirmCallback(this);
        rabbitTemplate.setReturnsCallback(this);
        rabbitTemplate.setMandatory(true);
    }
    @Autowired
    private RabbitTemplate rabbitTemplate;
    /**
     * 消息发送确认回调
     */
    @Override
    public void confirm(CorrelationData correlationData, 
                       boolean ack, 
                       String cause) {
        if (ack) {
            log.info("消息发送成功: correlationId={}", correlationData.getId());
        } else {
            log.error("消息发送失败: correlationId={}, cause={}", 
                correlationData.getId(), cause);
            // 可以进行重试处理
        }
    }
    /**
     * 消息回退回调(消息无法路由到队列时触发)
     */
    @Override
    public void returnedMessage(ReturnedMessage returnedMessage) {
        log.error("消息无法路由: exchange={}, routingKey={}, replyCode={}, replyText={}", 
            returnedMessage.getExchange(),
            returnedMessage.getRoutingKey(),
            returnedMessage.getReplyCode(),
            returnedMessage.getReplyText());
        // 可以将消息保存到数据库,进行后续处理
    }
}

使用示例

@RestController
@RequestMapping("/mq")
@Slf4j
public class RabbitController {
    @Autowired
    private MessageSender messageSender;
    @PostMapping("/send/direct")
    public String sendDirect(@RequestBody OrderMessage orderMessage) {
        messageSender.sendDirectMessage(orderMessage);
        return "发送成功";
    }
    @PostMapping("/send/topic")
    public String sendTopic(@RequestParam String routingKey, 
                           @RequestBody Object message) {
        messageSender.sendTopicMessage(routingKey, message);
        return "发送成功";
    }
    @PostMapping("/send/fanout")
    public String sendFanout(@RequestBody Object message) {
        messageSender.sendFanoutMessage(message);
        return "发送成功";
    }
}

重要配置说明

  1. 手动ACK:确保消息被正确处理后才确认
  2. 消息持久化:队列和消息需要设置durable=true
  3. 重试机制:设置合理重试次数,避免死循环
  4. 回调机制:实现ConfirmCallback和ReturnsCallback监控消息状态
  5. 死信队列:处理失败消息,避免消息丢失

这样配置,就可以在SpringBoot中完整地使用RabbitMQ进行消息收发。

抱歉,评论功能暂时关闭!