RocketMQ 消费端一秒进来几万条消息,数据库连接池只有 20 个连接——不限流,DB 瞬间打满。Guava 的 RateLimiter 就是给消费端加一个水龙头:令牌桶出令牌,消费线程抢到令牌才能执行,抢不到就等着。这篇文章从令牌桶的原理讲起,到两种落地写法(硬编码 + Nacos 动态配置 + @RefreshScope),全部用真实代码串下来。


没有限流的时候,问题长什么样

正常流量下一切正常。秒杀、大促、热点事件一来,消息队列瞬间堆积几十万条——消费端多线程全力消费,每秒几百条 SQL 打到 DB 上。

正常:    100 条/秒 → MySQL 稳如老狗
高峰:  10000 条/秒 → MySQL 连接池耗尽 → 请求超时 → 雪崩
MySQL Consumer RocketMQ MySQL Consumer RocketMQ 连接池满了! 堆积越来越严重 推送消息(大批量) INSERT × 100... INSERT × 100... INSERT × 100... INSERT(等连接……超时) 消费失败,消息重试

不是 MySQL 不行,是没有控制消费速度。解决办法就是在消费端加个限流器——令牌桶就是干这个的。


令牌桶是什么

每秒存 N 个

有令牌,拿走一个

有令牌,拿走一个

没令牌了,等着

令牌工厂
每秒生产 N 个令牌

令牌桶
最多存 M 个令牌

处理请求 1

处理请求 2

处理请求 3

原理就四句话

  1. 系统以固定速率往桶里放令牌(比如每秒 5000 个)
  2. 每个请求处理前必须从桶里拿到一个令牌
  3. 桶满了就不再放令牌
  4. 没令牌了就等着新令牌产生

Guava 的 RateLimiter 封装好了这一切。你只需要两行代码:

// 创建:每秒生成 5000 个令牌
private RateLimiter rateLimiter = RateLimiter.create(5000);

// 消费前拿令牌:有就直接返回,没有就阻塞等待
rateLimiter.acquire();

基础写法——硬编码在 Consumer 里

@Component
@RocketMQMessageListener(consumerGroup = "group_count_following_2_db",
        topic = MQConstants.TOPIC_COUNT_FOLLOWING_2_DB)
@Slf4j
public class CountFollowing2DBConsumer implements RocketMQListener<String> {

    @Resource
    private UserCountDOMapper userCountDOMapper;

    // 每秒 5000 个令牌——写死在代码里
    private RateLimiter rateLimiter = RateLimiter.create(5000);

    @Override
    public void onMessage(String body) {
        // 拿令牌,没有就阻塞等
        rateLimiter.acquire();

        // 消息体 JSON → DTO
        CountFollowUnfollowMqDTO countDTO = JsonUtils.parseObject(body, CountFollowUnfollowMqDTO.class);

        int count = countDTO.getType().equals(FollowUnfollowTypeEnum.FOLLOW.getCode()) ? 1 : -1;
        userCountDOMapper.insertOrUpdateFollowingTotalByUserId(count, countDTO.getUserId());
    }
}

这个写法下,RateLimiter 是这个 Consumer Bean 的一个成员变量,当前实例内共享。部署多实例的话,每个实例的 Consumer 各自持有一个独立的 RateLimiter——实例 A 的限流器管不了实例 B 的消费速度。所以你要是在 Nacos 里配了 rate-limit: 5000,部署了 3 个实例,打到 DB 的 QPS 上限实际上是 3 × 5000 = 15000


进阶写法——Nacos 动态配置,改限流阈值不重启

硬编码 5000,要是秒杀活动流量翻倍,想临时调成 10000——改代码、打包、部署、重启,半小时过去了。

把阈值抽到 Nacos 配置里,加上 @RefreshScope,控制台改完实时生效:

@Configuration
@RefreshScope
public class FollowUnfollowMqConsumerRateLimitConfig {

    @Value("${mq-consumer.follow-unfollow.rate-limit}")
    private double rateLimit;

    @Bean
    @RefreshScope
    public RateLimiter rateLimiter() {
        return RateLimiter.create(rateLimit);
    }
}
# application-dev.yml——阈值写本地
mq-consumer:
  follow-unfollow:
    rate-limit: 5000

或者放进 Nacos 配置文件里(动态刷新):

# Nacos 里的 your-service-dev.yaml
mq-consumer:
  follow-unfollow:
    rate-limit: 5000

改了 Nacos 里的值,@RefreshScope 销毁旧的 RateLimiter Bean,用新的阈值重建一个。

Consumer 端注入这个 Bean 即可:

@Component
@RocketMQMessageListener(consumerGroup = "group_follow_unfollow",
        topic = MQConstants.TOPIC_FOLLOW_OR_UNFOLLOW)
@Slf4j
public class FollowUnfollowConsumer implements RocketMQListener<Message> {

    @Resource
    private RateLimiter rateLimiter;  // 注入配置类里创建的 Bean

    @Override
    public void onMessage(Message message) {
        rateLimiter.acquire();  // 拿令牌
        // 消费逻辑...
    }
}

硬编码 vs Nacos 动态配置:

方式 优点 缺点
硬编码 简单,不需要额外配置 改阈值要重启
Nacos + @RefreshScope 控制台改完实时生效,不用重启 多了一个配置类

生产环境建议后者,原因有两个。

一是应急——线上 DB 突然扛不住了,Nacos 里把 5000 改成 2000,所有实例秒级生效,不用重启。

二是配合扩缩容,维持全局 QPS 不变——这个更重要:

DB 能扛住的写入上限:10000/s

部署 1 台实例:    rate-limit = 10000    → 总 QPS = 10000 ✅
扩容到 3 台实例:  不改配置的话 → 总 QPS = 30000 💥 DB 打挂
                  Nacos 里改成 3300 →   总 QPS ≈ 10000 ✅

实例数量变了,不改代码、不重启,Nacos 控制台改一个数字,所有实例的总 QPS 重新对齐到 DB 能承受的值。


RateLimiter 的核心方法

RateLimiter limiter = RateLimiter.create(5000);  // 每秒 5000 个令牌

// acquire():阻塞获取,没令牌就一直等
limiter.acquire();       // 拿 1 个令牌,阻塞
limiter.acquire(10);     // 拿 10 个令牌,阻塞

// tryAcquire():非阻塞获取,没令牌立刻返回 false
if (limiter.tryAcquire()) {
    // 拿到了,执行
} else {
    // 没拿到,直接跳过(适用于可丢弃的场景)
}

// tryAcquire(timeout):限定时间内的阻塞获取
if (limiter.tryAcquire(500, TimeUnit.MILLISECONDS)) {
    // 500ms 内拿到了
} else {
    // 超时没拿到
}

// getRate() / setRate():查询和动态修改速率
double currentRate = limiter.getRate();   // 当前速率
limiter.setRate(10000);                    // 动态改成 10000/s

MQ 消费端几乎只用 acquire()——等他一会儿比跳过这条消息重要得多。


放在 consumeMessage 的哪一行

很多人会把 acquire() 放在 onMessage 的最后几行——觉得"先处理业务,再限流"。错了。

acquire() 必须放在 onMessage 的第一行。 在消费逻辑之前就把令牌拿了,不让多余的请求进到 DB 操作里。

@Override
public void onMessage(Message message) {
    // ✅ 第一行就拿令牌——DB 操作之前就被限流了
    rateLimiter.acquire();

    // ❌ 不要放在这里——SQL 都执行完了才限流,没意义
    String bodyJsonStr = new String(message.getBody());
    // ...DB 操作...
}

你的各个消费者都设多少阈值

不是所有 Consumer 都需要限流。判断要不要加:

  1. Consumer 有没有 DB 操作?没有(比如只写 Redis)→ 不加,Redis 扛得住
  2. DB 操作是单条还是批量?批量(一次写多条)→ 阈值设低一点
  3. 顺带有没有其他外部调用?有 → 阈值设低一点,外部接口慢
// 批量写入型 Consumer——阈值设低一点,比如 1000/s
// 因为每次消费处理的不是一条 SQL,是一批
private RateLimiter rateLimiter = RateLimiter.create(1000);

// 单条写入型 Consumer——阈值可以高一些,比如 5000/s
private RateLimiter rateLimiter = RateLimiter.create(5000);

// 无 DB 操作,只操作 Redis——不需要限流
// Redis 单机 10w QPS,消费端这点量打不挂

注意:RateLimiter 是 JVM 级别的,限的是单个实例

RateLimiter 是 Guava 的本地限流器——它只是 JVM 内存里的一个计数器,不跨网络,不跨实例。

这意味着两件事:

  1. 一个 Consumer 里的 RateLimiter 对这个实例内所有消费线程生效。 ConsumeMode.ORDERLY 下,一个 Consumer 可能被分配了 4 个 MessageQueue,4 个线程各自消费。这 4 个线程共享同一个 RateLimiter——加起来每秒最多 5000 个令牌,一个线程 acquire() 阻塞了,其他线程也能拿到令牌。

  2. 部署了 3 个服务实例,打到 DB 的总 QPS 是 3 × 5000 = 15000 每个实例各自一个 RateLimiter,互不干扰。

所以如果你希望全局限制 DB 写入 QPS 不超过 5000(不管多少实例),Guava RateLimiter 做不到——它管不了别的 JVM。这种情况需要用 Redis 做分布式限流。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐