跳到主要内容
极客日志极客日志面向AI+效率的开发者社区
首页博客我的书AI学习GitHub 精选镜像AI 生图工具UI配色美学关于
搜索内容 / 工具 / 仓库 / 镜像...⌘K搜索
注册
博客列表
Javajava

RabbitMQ/Spring-AMQP 高级特性:事务机制与消息限流实战

RabbitMQ 事务机制确保消息发布的原子性,需禁用 Publisher Confirms 以避免模式冲突。消息限流通过 Prefetch Count 控制消费者未确认消息数,结合手动 ACK 机制,防止消息积压和系统崩溃。本文展示了如何在 Spring AMQP 中配置事务模板及限流策略,涵盖生产者发送、消费者接收及核心配置文件详解。

w795471发布于 2026/3/23更新于 2026/10/870 浏览
RabbitMQ/Spring-AMQP 高级特性:事务机制与消息限流实战

RabbitMQ/Spring-AMQP 高级特性:事务机制与消息限流实战

在分布式系统中,消息的可靠投递至关重要。Spring AMQP 提供了多种机制来保障这一点,其中事务(Transaction)和限流(Flow Control)是生产环境中最常用的两个手段。

1. 事务机制

AMQP 协议实现了事务机制,主要用于确保消息发布的原子性。简单来说,它允许你将多个操作(如发送多条消息、确认操作)绑定在一起,要么全部成功,要么全部失败。

实现示例

在 Spring Boot 中,我们可以利用 @Transactional 注解配合配置好的 RabbitTemplate 来实现事务发送。

@RestController
@RequestMapping("/producer")
public class ProducerController {
    @Resource(name = "transRabbitTemplate")
    private RabbitTemplate transRabbitTemplate;

    @Transactional
    @RequestMapping("/trans")
    public String trans() {
        // 第一条消息
        transRabbitTemplate.convertAndSend("", Constants.TRANS_QUEUE, "trans test ---> 1");
        
        // 模拟业务异常,导致事务回滚
        int num = 5 / 0;
        
        // 这条消息不会发送
        transRabbitTemplate.convertAndSend("", Constants.TRANS_QUEUE, "trans test ---> 2");
        return "发送成功";
    }
}

事务执行流程

注意: Spring Boot 的 RabbitMQ 自动配置默认开启 Publisher Confirms 模式,但 RabbitMQ 不允许同一个通道同时使用事务模式和 Confirm 模式。因此,必须显式禁用 Confirm 功能。

spring:
  rabbitmq:
    publisher-confirm-type: none
    publisher-returns: false

配置 Bean

我们需要手动创建支持事务的 RabbitTemplate 和事务管理器。

@Configuration
public class RabbitTemplateConfig {
    @Bean("transRabbitTemplate")
    public RabbitTemplate transRabbitTemplate(ConnectionFactory connectionFactory) {
        RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
        // 开启通道事务
        rabbitTemplate.setChannelTransacted(true);
        return rabbitTemplate;
    }

    @Bean
    public RabbitTransactionManager rabbitTransactionManager(ConnectionFactory connectionFactory) {
        return new RabbitTransactionManager(connectionFactory);
    }
}

队列配置保持标准即可:

@Configuration
public class RabbitMQConfig {
    @Bean("transQueue")
    public Queue transQueue() {
        return QueueBuilder.durable(Constants.TRANS_QUEUE).build();
    }
}

2. 消息限流

消息限流(Flow Control)是防止生产者发送速度超过消费者处理能力的关键机制。如果处理不及时,会导致消息积压甚至系统崩溃。

生产者代码

这里演示一个简单的批量发送场景:

@RestController
@RequestMapping("/producer")
public class ProducerController {
    @Resource(name = "rabbitTemplate")
    private RabbitTemplate rabbitTemplate;

    @RequestMapping("/qos")
    public String qos() {
        for (int i = 0; i < 20; i++) {
            rabbitTemplate.convertAndSend(Constants.QOS_EXCHANGE, "qos", "qos test:" + i);
        }
        return "发送成功";
    }
}

消费者配置与手动 ACK

为了实施限流,消费者端需要关闭自动确认,改为手动确认,并设置预取数量(Prefetch Count)。这意味着消费者每次只拉取少量消息,处理完一条再拉下一条。

@Component
@Slf4j
public class QosListener {
    @RabbitListener(queues = Constants.QOS_QUEUE)
    public void handMessage(Message message, Channel channel) throws IOException {
        long deliveryTag = message.getMessageProperties().getDeliveryTag();
        try {
            log.info("接收到消息:{}, deliveryTag:{}", 
                new String(message.getBody(), StandardCharsets.UTF_8), deliveryTag);
            log.info("处理成功");
            // 处理成功后手动确认
            // channel.basicAck(deliveryTag, true);
        } catch (Exception e) {
            // 处理失败,拒绝消息并重新入队
            channel.basicNack(deliveryTag, true, true);
        }
    }
}

关键配置

在 application.yml 中,我们需要将监听器模式设为 manual,并限制 prefetch 数量。

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual # 开启手动确认
        prefetch: 5              # 每个消费者最多未确认 5 条消息

最后,确保交换器、队列和绑定关系配置正确:

@Configuration
public class RabbitMQConfig {
    @Bean("qosQueue")
    public Queue qosQueue() {
        return QueueBuilder.durable(Constants.QOS_QUEUE).build();
    }

    @Bean("qosExchange")
    public DirectExchange qosExchange() {
        return ExchangeBuilder.directExchange(Constants.QOS_EXCHANGE).build();
    }

    @Bean("qosBinding")
    public Binding qosBinding(@Qualifier("qosExchange") DirectExchange directExchange,
                              @Qualifier("qosQueue") Queue queue) {
        return BindingBuilder.bind(queue).to(directExchange).with("qos");
    }
}

通过合理配置 Prefetch 和手动 ACK,我们可以在保证消息不丢失的前提下,有效控制消费端的压力,避免系统过载。

目录

  1. RabbitMQ/Spring-AMQP 高级特性:事务机制与消息限流实战
  2. 1. 事务机制
  3. 实现示例
  4. 配置 Bean
  5. 2. 消息限流
  6. 生产者代码
  7. 消费者配置与手动 ACK
  8. 关键配置

更多推荐文章

查看全部
  • 基于 Next.js 构建支持 TokenP 钱包登录的 DApp 前端实战
  • 机器人通讯架构选型:CAN/FD、RS485 与 EtherCAT 深度对比
  • 滑动窗口算法详解与经典例题实战
  • Java 文件操作与基础:流对象、数字码表及缓冲区
  • Python Web 开发框架选型指南
  • 如何成为一个懂 AI 的产品经理
  • Z-Image-Turbo 新手入门:从零上手 AI 绘画
  • Vitis AI 推理加速实战:从零实现 FPGA 部署
  • Claude Code 高级编程技巧与实战项目详解
  • Python 调用高德地图 MCP 服务查询天气示例
  • ThinkPad T480 安装 macOS 的实用配置笔记
  • 强化学习:演员评论家 Actor-Critic 算法详解与实现
  • C++26 CPU 资源隔离机制与性能优化实践
  • Windows + WSL + Ubuntu 安装 OpenClaw 及飞书机器人百炼模型配置流程
  • Node.js 环境搭建与 npm 配置实战指南
  • Linux 进程优先级与 O(1) 调度算法解析
  • 基于 Next.js 与 Wagmi 实现 TokenP 钱包登录的 DApp 前端开发
  • Python Flask 轻量级 Web 框架基础指南
  • OpenClaw 及其六大开源替代方案对比与选型指南
  • Flutter wallet_connect 鸿蒙适配:Web3 钱包连接与 DApp 授权

相关免费在线工具

  • Keycode 信息

    查找任何按下的键的javascript键代码、代码、位置和修饰符。 在线工具,Keycode 信息在线工具,online

  • Escape 与 Native 编解码

    JavaScript 字符串转义/反转义;Java 风格 \uXXXX(Native2Ascii)编码与解码。 在线工具,Escape 与 Native 编解码在线工具,online

  • JavaScript / HTML 格式化

    使用 Prettier 在浏览器内格式化 JavaScript 或 HTML 片段。 在线工具,JavaScript / HTML 格式化在线工具,online

  • JavaScript 压缩与混淆

    Terser 压缩、变量名混淆,或 javascript-obfuscator 高强度混淆(体积会增大)。 在线工具,JavaScript 压缩与混淆在线工具,online

  • Base64 字符串编码/解码

    将字符串编码和解码为其 Base64 格式表示形式即可。 在线工具,Base64 字符串编码/解码在线工具,online

  • Base64 文件转换器

    将字符串、文件或图像转换为其 Base64 表示形式。 在线工具,Base64 文件转换器在线工具,online