
做后端的人大部分都觉得 RabbitMQ 自带削峰填谷的能力生产者尽管往里发消息就行了。直到线上某次活动流量上来睡醒一看核心队列里积压了几百万条消息紧接着 RabbitMQ 内存告警、磁盘告警连接被 block下游消费者疯狂重试最后连数据库连接池都撑不住——这一整套雪崩链路我只亲身经历了一次就把“生产者限流”这件事刻进了肌肉记忆。这篇文章就是用 Java Spring Boot把 RabbitMQ 生产端的限流方案从信号量到令牌桶完整梳理一遍。文章里的代码可以直接复制到本地 Demo 跑起来也附带参数计算思路、监控手段和排障经验。适合正在用 Spring Boot 做消息服务的开发者也适合准备 Java 后端面试、或者负责消息中间件稳定性的人作为参考。限流不是消费者的事生产端同样需要一道闸门这篇文章就是讲清楚这道闸门怎么设计、怎么落地。1. 为什么生产者也需要限流1.1 生产端失控的雪崩链条很多人误以为 RabbitMQ 既然叫“消息队列”天然就能缓冲压力。这话对了一半队列确实能把生产端和消费端的速率解耦但队列本身的存储能力是有上限的。RabbitMQ 的消息默认先落内存再按策略刷盘一旦积压速度超过消费速度内存占用持续上涨触发memory_alarm此时 broker 会主动 block 所有生产连接。接下来发生的事情比较扎心生产者连接被 block发消息开始抛异常或者长时间阻塞业务代码里的重试逻辑被触发重试风暴进一步放大流量消费端还在努力消费但积压的消息带着过期时间堆积大量消息变成死信生产者所在的 JVM 线程被 MQ 客户端 IO 拖住Tomcat 线程池被打满接口大面积超时。这一整条链路里RabbitMQ 只是一个导火索真正的问题是生产端在无限制地把流量灌进一个容量有限的中间件。消费者就算再优化也需要时间中间件就算再稳定也有内存和磁盘的物理边界。所以生产端限流不是可选项是稳定性设计里的必需品。1.2 限流方案的选型逻辑生产端限流本质上是把上游流量“驯服”成一个下游和 broker 都能接受的速率。常见的限流手段有信号量、固定窗口计数器、滑动窗口、漏桶和令牌桶各自的适用场景差异很大方案控制维度是否支持突发实现成本适用场景信号量并发数支持但速率不受控最低保护单机线程资源固定窗口请求数/时间窗边界处容易突刺低简单接口限流滑动窗口请求数/精确时间窗平滑中网关层接口限流漏桶固定处理速率不支持突发中严格控制消费速率令牌桶平均速率 突发容量支持有限突发中生产者发消息场景这篇文章重点讲信号量和令牌桶是因为这两个方案在 RabbitMQ 生产者限流里最有代表性。信号量解决的是“同时有多少线程在发消息”的问题实现极简单是很多人做消息防崩的第一道闸门令牌桶解决的是“每秒钟能发多少条消息”的问题更接近流量速率控制的本质。两个方案可以独立用也可以组合用后面我会给出具体代码和调参经验。2. 信号量方案最朴素的并发闸门2.1 信号量的核心思路Java 的Semaphore是基于 AQS 实现的一个并发工具核心就两个操作acquire()获取许可release()释放许可。初始化时可以指定许可数量比如new Semaphore(32)表示最多允许 32 个线程同时执行后续代码段。拿它做生产者限流思路非常直白在发消息之前先acquire发完消息在finally里release。这相当于在 MQ 客户端前面加了一道闸门同一时间最多只有 N 个线程在调rabbitTemplate.convertAndSend()。为什么要控制并发数因为 RabbitMQ 的 Java client 底层会为每个 Channel 分配一个发送线程如果无限创建 Channel、无限并发发送首先遭殃的是 JVM 的线程数和堆内存其次才是 broker。信号量的优势是零依赖、开箱即用。不需要引入额外组件不需要考虑 Redis 可用性代码量也就十几行。对于单机部署、或者发送量不高的业务这就是最实用的第一道防线。2.2 Spring Boot 实战代码先建一个消息发送服务用信号量包住发送动作。这里我建议用tryAcquire加超时时间的方式而不是无参的acquire()因为无参等待可能让调用线程无限阻塞接口侧直接超时。Component public class OrderMessageProducer { private static final Logger log LoggerFactory.getLogger(OrderMessageProducer.class); private final RabbitTemplate rabbitTemplate; /** * 同一时间最多 32 个线程在发送消息 */ private final Semaphore sendSemaphore new Semaphore(32); public OrderMessageProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate rabbitTemplate; } public boolean sendOrderMessage(String message) { boolean acquired false; try { // 最多等 500ms拿不到许可就直接降级 acquired sendSemaphore.tryAcquire(500, TimeUnit.MILLISECONDS); if (!acquired) { log.warn(send queue is full, degrade order message); return false; } rabbitTemplate.convertAndSend(order.exchange, order.routing.key, message); return true; } catch (InterruptedException e) { Thread.currentThread().interrupt(); return false; } catch (Exception e) { log.error(send message failed, e); return false; } finally { if (acquired) { sendSemaphore.release(); } } } }核心要点是两个超时时间要结合业务接口的容忍度去设置一般 300ms 到 1s 之间比较合理太短会导致瞬时流量大量降级太长会拖垮上游调用方release()必须在finally里执行否则一旦发送抛异常许可就永久丢失很快所有线程都会卡死在acquire上。踩过这个坑的人都知道信号量泄漏造成的“假死”比接口本身超时更隐蔽排查半天才能对上号。2.3 信号量的局限与适用边界信号量也有明显短板。它限制的是并发线程数不是消息发送速率。假设 32 个线程同时发送每个线程每秒可以发 100 条总速率依然有 3200 TPS如果消费者端每秒只能处理 500 条积压照样发生只是时间上延缓了而已。另外信号量是进程内的多实例部署时每个节点的 32 个许可各自独立整体并发数可能是 32 乘以实例数单机限流并不等于集群限流。想实现全局限流就必须引入 Redis 或者别的分布式协调组件复杂度就上来了。信号量也不感知 RabbitMQ 的 block 状态。broker 已经内存告警、把所有生产连接都 block 了信号量依然“大方地”放行线程进入发送逻辑只是线程最终在 MQ 客户端那里阻塞或者报错。所以信号量适合做“在途并发控制”但解决不了“速率匹配”的问题。想从根本上让生产速率贴合消费速率还得看令牌桶。3. 令牌桶方案让流量像水龙头一样可控3.1 令牌桶原理与实际效果令牌桶的类比特别直观想象一个景区售票窗口旁边有一个桶每隔固定时间往桶里放一张令牌桶满了令牌就溢出来。游客必须先拿到令牌才能买票没令牌就必须在门口等。令牌桶本身有一个容量上限所以如果积累了一批令牌短时间内的确可以允许一波游客同时涌入但长期来看单位时间能进去的人数严格受制于放令牌的速率。对应到 RabbitMQ 生产者限流场景令牌就是“发一条消息的资格”放令牌的速率就是“每秒允许发送的消息数”。为什么要留容量上限因为生产者偶尔会有批量发送需求比如一次业务操作要发 10 条消息绝对均匀的限流会让这一次操作等 10 个时间片体验太差。令牌桶允许这 10 条消息一次性消耗攒下来的 10 个令牌既照顾了突发又不至于突破平均速率的天花板。与信号量对比令牌桶直接把“限流”定义从并发数变成了速率方向更准确。RabbitMQ 消费者处理消息的速率是可以压测出来的生产者限流速率只要设成消费者速率的 70% 到 85%队列积压就会稳定在可控范围内这比盲目控制线程数要科学得多。3.2 自研一个轻量令牌桶有不少团队为了不引入 Guava选择自研令牌桶。这里给一个基于时间戳计算的实现不需要额外定时任务靠“上次发放时间”和当前时间差来补令牌代码简单且运行稳定public class TokenBucket { private final long capacity; private final double refillRatePerSecond; private double tokens; private long lastRefillTimeMillis; private final ReentrantLock lock new ReentrantLock(); public TokenBucket(long capacity, double refillRatePerSecond) { this.capacity capacity; this.refillRatePerSecond refillRatePerSecond; this.tokens capacity; this.lastRefillTimeMillis System.currentTimeMillis(); } public boolean tryAcquire(int permits, long timeoutMillis) { lock.lock(); try { refill(); if (tokens permits) { tokens - permits; return true; } return false; } finally { lock.unlock(); } } private void refill() { long current System.currentTimeMillis(); long elapsed current - lastRefillTimeMillis; if (elapsed 0) { double newTokens elapsed / 1000.0 * refillRatePerSecond; tokens Math.min(capacity, tokens newTokens); lastRefillTimeMillis current; } } }这个实现有几个细节需要注意。refill()只有在tryAcquire被调用时才会执行这意味着长时间没有消息发送时令牌会一直累积到容量上限不会过期这是正确的令牌桶行为计算令牌数时用了double在高并发大流量下要注意精度损失但实际场景中误差可以忽略锁用的是ReentrantLock比synchronized更灵活也支持后续扩展 tryLock 超时逻辑。容量建议设成每秒放行数量的 1 到 2 倍例如速率是 160 TPS容量设 160相当于允许最多 1 秒的突发量不会造成瞬间冲击。3.3 使用 Guava RateLimiter 的实战写法如果项目里不排斥第三方依赖强烈建议直接用 Guava 的RateLimiter。它本身就是线程安全的可以直接注册成 Spring Bean 使用底层实现比大多数人自己写的要细腻很多。引入依赖dependency groupIdcom.google.guava/groupId artifactIdguava/artifactId version33.0.1-jre/version /dependency配置类里声明限流器Configuration public class RabbitRateLimitConfig { /** * 每秒放行 160 个令牌容量允许 160 个突发 */ Bean public RateLimiter producerRateLimiter() { return RateLimiter.create(160.0); } }发送消息的服务里注入这个 BeanComponent public class RateLimitedMessageSender { private static final Logger log LoggerFactory.getLogger(RateLimitedMessageSender.class); private final RabbitTemplate rabbitTemplate; private final RateLimiter producerRateLimiter; public RateLimitedMessageSender(RabbitTemplate rabbitTemplate, RateLimiter producerRateLimiter) { this.rabbitTemplate rabbitTemplate; this.producerRateLimiter producerRateLimiter; } public boolean sendWithRateLimit(String exchange, String routingKey, Object message) { // 最多等 300ms拿不到令牌就快速失败 if (!producerRateLimiter.tryAcquire(300, TimeUnit.MILLISECONDS)) { log.warn(rate limit triggered, message rejected); return false; } try { rabbitTemplate.convertAndSend(exchange, routingKey, message); return true; } catch (Exception e) { log.error(send message failed, e); throw new RuntimeException(send message failed, e); } } }tryAcquire(300, TimeUnit.MILLISECONDS)和信号量的超时逻辑思路一致没拿到令牌就快速失败把压力挡在业务入口而不是让线程一直挂在 MQ 客户端里。这里有个很多人容易搞错的点RateLimiter.acquire()是阻塞式的能保证请求一定能发出去但可能阻塞几十毫秒甚至几百毫秒对接口 RT 影响很大tryAcquire()是快速失败的实际业务里更常用。如果调用方接受短暂等待也可以做一层缓冲例如先tryAcquire拿不到再降级写本地表异步重试。3.4 令牌桶在实际项目中的几个坑第一Guava 的RateLimiter.create()默认创建的是SmoothBursty它允许一定程度的突发甚至在没有任何预热的情况下前几秒可能会放行超过实际速率的流量。对 MQ 发送场景来说这个突发窗口有可能在系统刚刚启动时造成短暂的冲击。更平滑的方式是RateLimiter.create(rate, warmupPeriod, TimeUnit.SECONDS)带预热时间让限流器从低速率逐步爬升到目标速率。如果业务场景里消费端冷启动也需要预热强烈建议用预热版本。第二多实例部署时每个实例各自持有一个RateLimiter全局限流不成立。假设两台实例每台限 160 TPS全局限流实际上是 320 TPS。要实现精确的全局限流需要把令牌桶挪到 Redis 侧用 Lua 脚本保证原子性或者直接使用 Redisson 的RRateLimiter。这个思路不复杂但要注意 Redis 本身的性能和可用性量力而行。第三令牌桶只能解决“发太快”的问题解决不了“发不出去”的问题。消息路由失败、broker 拒绝、序列化异常这些错误还是得靠 confirm 回调、重试和死信策略兜底。限流和可靠性是两套独立的机制别混为一谈。4. 与 RabbitMQ 原生机制的组合拳4.1 开启发布确认与路由回调限流只解决速率问题消息发送成功与否需要另一套机制确认。Spring Boot 里先把发布确认和路由回调打开配置如下spring: rabbitmq: host: 127.0.0.1 port: 5672 publisher-confirm-type: correlated publisher-returns: true template: mandatory: truepublisher-confirm-type: correlated表示发送消息时携带CorrelationDatabroker 异步回调确认消息是否落盘mandatory: true配合publisher-returns: true路由不到队列时会把消息退回来。然后在RabbitTemplate初始化时注册回调PostConstruct public void init() { rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { log.error(message confirm failed, cause: {}, cause); // 这里做补偿记录日志、落库、重试 } }); rabbitTemplate.setReturnsCallback(returned - { log.error(message routed failed, exchange: {}, routingKey: {}, replyText: {}, returned.getExchange(), returned.getRoutingKey(), returned.getReplyText()); }); }有两点经验分享。confirm 回调是异步线程池执行的不要在回调里做耗时操作更不要在主线程里同步等待 confirm 结果否则限流带来的吞吐提升会被等待消耗掉route 失败的消息说明业务上有路由配置错误这类问题不仅要记录日志最好能发送到独立的“失败消息表”或者备用队列里人工处理。4.2 连接被 block 时的监听与降级RabbitMQ 内存或磁盘达到阈值时会主动 block 所有生产连接。此时生产端不管怎么限流消息都发不出去继续发只会加重 broker 的负担。正确的做法是监听连接的 block 事件在 block 期间直接停掉发送动作unblock 之后再恢复。Spring AMQP 提供了连接监听机制可以拿到底层 com.rabbitmq.client.ConnectionConfiguration public class RabbitConnectionBlockListenerConfig { Bean public CachingConnectionFactory rabbitConnectionFactory() { CachingConnectionFactory factory new CachingConnectionFactory(); factory.setHost(127.0.0.1); factory.setPort(5672); factory.setPublisherConfirmType(CachingConnectionFactory.ConfirmType.CORRELATED); factory.addConnectionListener(new ConnectionListener() { Override public void onCreate(Connection connection) { com.rabbitmq.client.Connection delegate (com.rabbitmq.client.Connection) connection.getDelegate(); delegate.addBlockedListener(new BlockedListener() { Override public void handleBlocked(String reason) { // 通知业务层进入降级模式 RateLimitSwitchHolder.setBlocked(true); } Override public void handleUnblocked() { RateLimitSwitchHolder.setBlocked(false); } }); } }); return factory; } }RateLimitSwitchHolder可以设计成一个简单的开关类block 标记为 true 时生产者服务直接拒绝发消息返回降级结果解除 block 后恢复正常发送。这个机制配合令牌桶非常关键令牌桶控制正常情况下的发送速率block 监听控制极端情况下的发送开关两道防线一起用才算真正防崩。4.3 消费端 QoS 与生产端限流的配合生产端限流不是孤立动作消费端的basicQos设置跟它是配套的。basicQos可以限制每个消费者在本地缓存的最大未确认消息数例如prefetch10让消费者处理完一条并确认之后才从 broker 拉取下一条避免消费者本地堆积大量消息导致内存溢出。Spring Boot 里可以直接配置spring: rabbitmq: listener: simple: prefetch: 10 acknowledge-mode: manual concurrency: 5 max-concurrency: 10生产令牌桶把发送速率限制在消费者最大吞吐的 80% 左右消费端用 QoS 保证每条消息都被处理完再拉新消息Broker 里的队列积压就会保持在一个稳定区间不会出现“生产者慢下来了消费者也拉不完”的情况。实际调优时可以先定消费端 prefetch再倒推生产端速率先统计消费者单条消息平均处理耗时假设并发消费 5 个线程平均每条耗时 50ms理论最大吞吐是 100 TPS生产端限速就设在 80 TPS 左右。5. 常见问题与排查技巧实录5.1 高频问题速查表现象可能原因解决方案连接被 block发送消息长时间阻塞broker 内存或磁盘达到阈值开启连接 block 监听block 期间暂停发送同时优化消费速度报reply-code320, reply-textCONNECTION_FORCED管理端强制关闭连接或 broker 资源告警检查 RabbitMQ 管理界面连接列表看是否触发了memory_alarm信号量拿到许可但消息发送失败broker 连接已关闭或交换机/队列不存在发消息前检查连接状态配合 confirm 回调确认路由限流之后队列积压依然上涨消费者吞吐远低于预期或限流速率设置过高重新压测消费端真实吞吐把生产速率降到 70%多实例部署限流不生效每个实例各持一个令牌桶换成 Redis 分布式令牌桶Guava RateLimiter 刚启动时突发流量大默认 SmoothBursty 允许预支令牌改用create(rate, warmupPeriod, TimeUnit.SECONDS)预热模式确认回调一直没执行未开启publisher-confirm-type检查 yaml 配置Confirm 回调需要 correlated 模式消息路由失败但没日志mandatory未开启退回消息被丢弃开启mandatory: true并注册 ReturnsCallback5.2 实战中容易踩的细节信号量和令牌桶的代码本身都不复杂线上出问题往往出在细节上。第一个细节是信号量释放。只要线程在tryAcquire成功之后、release之前有异常路径许可就一定漏掉。常规的防御手段是在finally里释放并且用 boolean 标记是否真正拿到了许可否则在RuntimeException的场景下会重复释放许可直接把信号量允许的并发数抬到无穷大。这个问题的隐蔽性在于重复释放不会立即报错只是并发上限悄悄没了。第二个细节是令牌桶容量和速率的比例。很多人把速率设成 100容量也设成 100这没问题。但把容量设成 1000速率 100就相当于允许瞬间发出 10 秒的流量队列积压会立即冲高。容量代表“容忍突发”的边界不是拍脑袋填的建议容量小于等于每秒速率的两倍。第三个细节是监控指标。限流方案上线之后至少要看三个指标消息发送 TPS、队列积压深度、消费者处理 TPS。可以通过 RabbitMQ 管理接口拉取队列积压curl -s -u guest:guest http://localhost:15672/api/queues/%2F/order.queue | jq .messages_ready, .messages_unacknowledged把这个数据接到 Grafana 或者定时任务里能快速判断限流参数是否合理。如果messages_ready持续上涨优先怀疑生产速率过高如果messages_unacknowledged居高不下优先怀疑消费者卡住或者 QoS prefetch 设置过大。第四个细节是降级策略。限流触发时最稳妥的降级不是直接丢弃消息而是写本地表或者 Redis 缓存后台再慢慢重发。我经常建议业务侧在RateLimiter.tryAcquire失败后把消息序列化后写入一张mq_resend_record表状态标记为 pending再起一个定时任务按低速率重发。这样既保护了 MQ 和消费者又不丢数据代价只是多一点存储。5.3 从信号量到令牌桶的演进建议如果你的系统现在还没有做任何生产端限流我建议按这个节奏演进第一周先上信号量方案代码量最小改动最轻能立刻挡住“线程堆积”这类最明显的风险运行稳定后再引入令牌桶把限流维度从并发数升级到速率再往后如果遇到多实例部署或者流量波动剧烈再考虑 Redis 分布式令牌桶和动态限流。动态限流是个很实用的扩展方向用消费者处理速率作为反馈信号动态调整生产端令牌桶速率。比如每隔 10 秒采集一次消费者 TPS如果连续三个周期消费者 TPS 低于生产速率就把令牌桶速率下调 10%如果消费者长期有余量则上调 5%。这样限流参数就不再需要人工反复调系统自己会找到平衡点。这个做法的前提是监控数据准确预留一个人工干预的开关不推荐在业务初期就上。最后再分享一个个人的体会RabbitMQ 生产端限流这个事越早做越好。我刚工作的时候总觉得限流是网关和消费端的事生产端只管发就行了结果线上出过一次故障之后才明白消息链路是一个整体每一端的稳定性都需要主动设计。信号量和令牌桶都不复杂复杂的反而是你愿不愿意在最平静的时候把这条防线先建起来。希望这篇实战记录能帮你少走一些弯路让你的 RabbitMQ 在流量高峰时也能稳稳撑住。