You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Guava RateLimiter QPS控制异常求助:需严格维持每秒5000条消息

解决Guava RateLimiter无法严格控制固定窗口QPS的问题

咱们先来拆解下你的问题:你需要的是严格的固定窗口限流——每一秒最多发5000条,哪怕上一秒只用了2000条,剩下的3000额度也不能留到下一秒。但Guava的RateLimiter用的是令牌桶算法,天生就会累积令牌,这正是你出现速率波动的核心原因。

为什么你的RateLimiter实现达不到预期?

Guava RateLimiter的核心逻辑是令牌桶:它会以每秒5000个的速度持续生成令牌,没被使用的令牌会暂时存起来(默认最多存1秒的量,也就是5000个)。如果某一秒你只发了1000条,桶里就会剩下4000个令牌;下一秒新生成5000个令牌后,总共有9000个可用,自然就会出现一下子发9000条的波动。

而且它的acquire()方法是用来控制长期平滑速率的,不是严格的每秒限额,完全不符合你“额度不累积”的需求。

解决方案:实现固定窗口限流

既然Guava的组件不适用,咱们可以自己写一个轻量的固定窗口限流工具,逻辑非常直观:

  1. 维护一个当前窗口的请求计数
  2. 每过1秒就重置计数
  3. 每次发送消息前先检查当前窗口的计数是否超过5000,没超过就允许发送,否则等待

代码实现

先写一个固定窗口限流的工具类:

import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;

public class FixedWindowRateLimiter {
    // 每秒最大允许的请求数
    private final int maxRequestsPerSecond;
    // 当前窗口的请求计数
    private final AtomicInteger currentCount = new AtomicInteger(0);
    // 当前窗口的开始时间(毫秒)
    private final AtomicLong windowStartTime = new AtomicLong(System.currentTimeMillis());

    public FixedWindowRateLimiter(int maxRequestsPerSecond) {
        this.maxRequestsPerSecond = maxRequestsPerSecond;
    }

    // 尝试获取发送许可,返回是否可以发送
    public boolean tryAcquire() {
        long now = System.currentTimeMillis();
        // 判断是否进入了新的窗口(间隔超过1秒)
        if (now - windowStartTime.get() > 1000) {
            // 用CAS确保只有一个线程能重置窗口,避免并发下的计数混乱
            if (windowStartTime.compareAndSet(windowStartTime.get(), now)) {
                currentCount.set(0);
            }
        }
        // 计数加1,判断是否超过限额
        return currentCount.incrementAndGet() <= maxRequestsPerSecond;
    }

    // 阻塞式获取许可,直到可以发送
    public void acquire() throws InterruptedException {
        // 循环尝试,直到获取到许可
        while (!tryAcquire()) {
            // 短暂休眠,避免自旋占用过多CPU
            Thread.sleep(1);
        }
    }
}

使用方式

替换掉原来的RateLimiter,用这个工具类即可:

FixedWindowRateLimiter limiter = new FixedWindowRateLimiter(5000);
for (int i = 0; i < msisdnSize; i++) {
    limiter.acquire();
    // 执行发送消息的逻辑
}

补充说明

这个实现用了原子类来处理并发场景,确保多线程下计数和窗口切换的正确性。如果你的发送逻辑是单线程的,甚至可以简化掉原子类,用普通变量就行,但用原子类会更稳妥。

如果之后需要更复杂的限流策略(比如滑动窗口、漏桶),可以考虑用Sentinel、Resilience4j这类专门的限流组件,但对于你现在的需求,这个轻量实现完全足够。

内容的提问来源于stack exchange,提问作者Guru Prasad

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 07:22:40