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

RabbitMQ Spring-AMQP 事务机制与消息限流配置详解

RabbitMQ 事务机制确保消息发布确认的原子性,通过设置 channelTransacted 启用事务模式,同时需禁用 publisher confirms 以避免冲突。消息限流通过手动确认模式配合 prefetch 参数控制消费者未确认消息数量,防止生产速度超过消费能力导致积压。配置示例展示了 Spring Boot 中 RabbitTemplate 的事务 Bean 定义及队列交换器绑定方式,实现可靠的消息传输与流量控制。

云朵棉花糖发布于 2026/3/23更新于 2026/9/2967 浏览
RabbitMQ Spring-AMQP 事务机制与消息限流配置详解

在这里插入图片描述

1. 事务

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

发送消息

@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 自动配置默认会确认模式,但 RabbitMQ 不允许同一个通道同时使用事务模式和确认模式,所以需要确保 publisher confirms 被禁用。

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

配置 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) 是 RabbitMQ 防止生产者发送消息速度超过消费者处理能力,导致消息积压和系统崩溃的保护机制。

发送消息

@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 "发送成功";
    }
}

配置消费者

@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);
        }
    }
}

限制每个消费者未确认的最大消息数

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual # 消费者确认机制
        prefetch: 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");
    }
}

在这里插入图片描述

目录

  1. 1. 事务
  2. 发送消息
  3. 配置 RabbitTemplate 和事务管理器
  4. 配置队列
  5. 2. 消息限流
  6. 发送消息
  7. 配置消费者
  8. 限制每个消费者未确认的最大消息数
  9. 声明和配置交换器、队列和绑定关系

更多推荐文章

查看全部
  • 在 Linux Ubuntu 上安装 Qt 5 详细教程
  • Fooocus 部署实战:从本地环境搭建到云端快速启动
  • Java Web 开发基础:Spring Web MVC 核心解析
  • 天工 AI 辅助产品经理工作流程与多模态功能体验
  • Python FastAPI 入门教程:新手快速构建 RESTful API
  • SpringBoot 自动配置原理深度解析与实战
  • 实现一行或多行文本溢出省略效果的常用方法
  • 二分查找实战:山峰数组峰顶索引与寻找峰值
  • 从 vw/vh 到 clamp(),前端响应式设计的痛点与进化
  • 大模型时代企业 AI 发展趋势分析
  • Spatial Joy 2025 全球 AR&AI 赛事:开发者要的资源、玩法、避坑攻略都在这
  • 量子计算驱动 Python 医疗诊断:变分量子分类器实战
  • OpenClaw 集成 Telegram 机器人开发指南
  • FastAPI:Python 高性能 Web 框架
  • 大型模型科普指南
  • Elasticsearch 进阶实战:JavaRestClient 操作索引与文档及海量数据批处理指南
  • OpenClaw self-improving-agent 技能详解:让 AI 从错误中自主学习
  • RabbitMQ/Spring-AMQP 事务与消息限流高级特性详解
  • 在 PyTorch 镜像中部署 Whisper.cpp 语音识别模型
  • Seedance 2.0 双分支扩散变换器架构解析与工程实现

相关免费在线工具

  • 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