当前位置:首页 > 技术 > 正文内容

基于 Spring Boot 的 RabbitMQ 核心机制与生产实践

访客 技术 2026年9月4日 1
1. 队列与交换机的基本认知

Queue(队列)是 RabbitMQ 中保存消息的容器,消息必须投递到队列里,消费者再从队列中拉取。RabbitMQ 与 Kafka 的最大区别在于:RabbitMQ 以队列为核心存储单元,而 Kafka 以 Topic 为逻辑层面、以分区文件为物理存储。

当多个消费者订阅同一个队列时,RabbitMQ 会采用轮询(Round-Robin)方式把消息分摊给这些消费者,因此一个消息不会被所有消费者同时消费。需要注意的是,RabbitMQ 原生不支持"队列级广播",如果需要广播,应使用 Fanout 交换机绑定多个队列来实现,而不是在队列层面做二次开发。

交换机(Exchange)常见类型有四种:

  • Direct:精确匹配路由键。
  • Fanout:广播到所有绑定队列,忽略路由键。
  • Topic:支持通配符 *# 匹配路由键。
  • Headers:根据消息头属性匹配。

2. Spring Boot 基础配置

下面给出一个完整的 application.yml 示例,涵盖集群地址、发布确认、Return 回调、Mandatory 设置以及消费者手动确认等关键项。

server:
  port: 11000

spring:
  application:
    name: rabbitmq-demo
  rabbitmq:
    addresses: node1:5672,node2:5672
    username: admin
    password: admin
    virtual-host: /demo
    publisher-returns: true
    publisher-confirm-type: correlated
    connection-timeout: 15000
    template:
      mandatory: true
    default-exchange: demo.direct
    listener:
      simple:
        acknowledge-mode: manual
        concurrency: 5
        max-concurrency: 10
        prefetch: 1
        retry:
          enabled: true

关键配置说明:

  • publisher-returns / publisher-confirm-type:开启 Return 与 Confirm 两种回调。
  • template.mandatory:交换机无法路由到队列时,消息是否返回生产者,必须设置为 true 才能触发 ReturnCallback。
  • acknowledge-mode: manual:消费端手动确认,业务处理完成后再 ACK。
  • prefetch:限制每个消费者同时只能处理指定条数未确认消息,常用于削峰。

3. Java 配置类

通过 @Configuration 声明交换机、队列、绑定关系,并创建两个 RabbitTemplate 示例:一个使用自定义连接工厂,一个直接注入 Spring 管理的连接工厂。

@Configuration
public class AmqpConfig {

    @Value("${spring.rabbitmq.addresses}")
    private String addresses;

    @Value("${spring.rabbitmq.username}")
    private String user;
    @Value("${spring.rabbitmq.password}")
    private String pwd;
    @Value("${spring.rabbitmq.virtual-host}")
    private String vhost;
    @Value("${spring.rabbitmq.default-exchange}")
    private String defaultExchange;

    @Autowired
    private ConfirmCallbackService confirmCallback;
    @Autowired
    private ReturnCallbackService returnCallback;

    @Bean
    public ConnectionFactory amqpConnectionFactory() {
        CachingConnectionFactory factory = new CachingConnectionFactory();
        factory.setAddresses(addresses);
        factory.setUsername(user);
        factory.setPassword(pwd);
        factory.setVirtualHost(vhost);
        return factory;
    }

    @Bean("customRabbitTemplate")
    public RabbitTemplate customRabbitTemplate(ConnectionFactory amqpConnectionFactory) {
        RabbitTemplate template = new RabbitTemplate(amqpConnectionFactory);
        template.setMessageConverter(jsonMessageConverter());
        template.setExchange(defaultExchange);
        template.setMandatory(true);
        template.setConfirmCallback(confirmCallback);
        template.setReturnsCallback(returnCallback);
        template.setReplyTimeout(20000);
        template.setReceiveTimeout(20000);
        return template;
    }

    @Bean
    public Jackson2JsonMessageConverter jsonMessageConverter() {
        return new Jackson2JsonMessageConverter();
    }

    @Bean
    public DirectExchange orderDirectExchange() {
        return new DirectExchange("order.direct", true, false);
    }

    @Bean
    public FanoutExchange logFanoutExchange() {
        return new FanoutExchange("log.fanout", true, false);
    }

    @Bean
    public TopicExchange noticeTopicExchange() {
        return new TopicExchange("notice.topic", true, false);
    }

    @Bean
    public Queue orderQueue() {
        return new Queue("order.queue", true);
    }

    @Bean
    public Queue smsQueue() {
        return new Queue("sms.queue", true);
    }

    @Bean
    public Queue emailQueue() {
        return new Queue("email.queue", true);
    }

    @Bean
    public Binding orderBinding() {
        return BindingBuilder.bind(orderQueue())
                .to(orderDirectExchange())
                .with("order.created");
    }

    @Bean
    public Binding smsBinding() {
        return BindingBuilder.bind(smsQueue()).to(logFanoutExchange());
    }

    @Bean
    public Binding emailBinding() {
        return BindingBuilder.bind(emailQueue())
                .to(noticeTopicExchange())
                .with("notice.#");
    }
}

声明交换机时主要关注三个属性:durable(是否持久化)、autoDelete(没有队列绑定时是否自动删除)、以及自定义扩展参数 arguments

4. 消息转换器

RabbitTemplate 通过 MessageConverter 完成 Java 对象与 Message 之间的转换。可以自定义转换器,也可以直接使用 Jackson2JsonMessageConverter

public class PlainTextConverter implements MessageConverter {

    @Override
    public Message toMessage(Object object, MessageProperties properties) throws MessageConversionException {
        return new Message(object.toString().getBytes(StandardCharsets.UTF_8), properties);
    }

    @Override
    public Object fromMessage(Message message) throws MessageConversionException {
        return new String(message.getBody(), StandardCharsets.UTF_8);
    }
}

使用 Jackson2JsonMessageConverter 时,发送端与接收端的类型必须一致,包括类名、字段以及包路径,否则反序列化会失败。如果业务中需要多态或松耦合,建议自定义转换器或在消息头里存放类型信息。

5. 生产者与消费者示例

下面展示一个统一封装的发送服务,通过 MessagePostProcessor 在发送前设置消息 ID、Correlation ID 以及持久化投递模式。

@Service
@Slf4j
public class MsgPublisher {

    @Resource(name = "customRabbitTemplate")
    private RabbitTemplate rabbitTemplate;

    public void sendObject(String routingKey, Object payload) {
        rabbitTemplate.convertAndSend(routingKey, payload, this::decorateMessage);
    }

    public void sendObject(String exchange, String routingKey, Object payload) {
        rabbitTemplate.convertAndSend(exchange, routingKey, payload,
                this::decorateMessage,
                new CorrelationData(UUID.randomUUID().toString()));
    }

    private Message decorateMessage(Message message) {
        message.getMessageProperties().setMessageId(UUID.randomUUID().toString());
        message.getMessageProperties().setCorrelationId(UUID.randomUUID().toString());
        message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
        return message;
    }
}

消费端可以使用 @RabbitListener 标注在类上,再搭配多个 @RabbitHandler 方法,Spring 会根据 MessageConverter 转换后的实际类型分派到对应方法。

@Component
@Slf4j
@RabbitListener(queues = "notice.queue")
public class NoticeConsumer {

    @RabbitHandler
    public void onString(@Payload String text) {
        log.info("收到文本消息:{}", text);
    }

    @RabbitHandler
    public void onOrder(@Payload Order order) {
        log.info("收到订单消息:{}", order);
    }

    @RabbitHandler
    public void onMap(@Payload Map map) {
        map.forEach((k, v) -> log.info("{}={}", k, v));
    }
}

需要注意:如果多个 @RabbitHandler 方法都接受 Map 类型,Spring 无法确定调用哪个方法,会抛出异常。因此不建议把 Map 作为通用接收类型。

6. 消息确认与手动 ACK

生产端需要实现 RabbitTemplate.ConfirmCallbackRabbitTemplate.ReturnsCallback,分别处理"是否送达交换机"和"是否能路由到队列"。

@Component
@Slf4j
public class OrderConfirmCallback implements RabbitTemplate.ConfirmCallback {
    @Override
    public void confirm(CorrelationData data, boolean ack, String cause) {
        if (ack) {
            log.info("消息已到达交换机,id={}", data != null ? data.getId() : "");
        } else {
            log.error("消息未到达交换机,id={},原因:{}", data != null ? data.getId() : "", cause);
        }
    }
}
@Component
@Slf4j
public class OrderReturnCallback implements RabbitTemplate.ReturnsCallback {
    @Override
    public void returnedMessage(ReturnedMessage returned) {
        log.error("消息被退回,exchange={},routingKey={},replyCode={},replyText={}",
                returned.getExchange(),
                returned.getRoutingKey(),
                returned.getReplyCode(),
                returned.getReplyText());
    }
}

消费端开启手动确认后,需要在业务成功后主动 ACK;失败时可以选择重新入队或拒绝。

@RabbitHandler
public void handleOrder(@Payload Order order, Channel channel, Message message) throws IOException {
    long tag = message.getMessageProperties().getDeliveryTag();
    try {
        // 业务处理
        channel.basicAck(tag, false);
    } catch (Exception e) {
        if (message.getMessageProperties().getRedelivered()) {
            channel.basicReject(tag, false);
        } else {
            channel.basicNack(tag, false, true);
        }
    }
}

三种回执方法的区别如下:

  • basicAck(deliveryTag, multiple):确认成功,multiple=true 时批量确认小于当前 deliveryTag 的所有消息。
  • basicNack(deliveryTag, multiple, requeue):确认失败,requeue=true 时重新入队。
  • basicReject(deliveryTag, requeue):单条拒绝,不能批量。

四种投递场景对应的回调触发情况:

  • 找不到交换机:ConfirmCallback 返回 false。
  • 找到交换机但找不到队列:ConfirmCallback 返回 true,同时触发 ReturnCallback
  • 交换机和队列都找不到:同第一种,ConfirmCallback 返回 false。
  • 投递成功:ConfirmCallback 返回 true。
]]>
标签: RabbitMQ

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Dom\HTML_NO_DEFAULT_NS 的副作用:自动加闭合标签

在使用Dom\HTMLDocument时,Dom\HTML_NO_DEFAULT_NS 将禁止在解析过程中设置元素的命名空间, 此设置是为了与DOMDocument向后兼容而存在的。当使用它时,已知的一个副作用就是:自动加闭合标签例如 </img> 为什么会这样?当你使用:Dom\HTML_NO_DEFAULT_NS文档会变成 无命名空间模式,此时内部更接近 XML...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。