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

基于令牌桶算法的Rate Limiter多线程限流请求消耗值异常排查

问题描述

我用Java实现了一个基于令牌桶算法的限流器,分别用单线程和Executor Service多线程测试:

  • 单线程测试:限流速率20请求/秒,运行10秒后总消耗请求数200,符合预期;
  • 多线程测试:相同配置下,即使加了synchronized关键字,总消耗请求数却为220,不符合预期。

期望单线程和多线程测试10秒内请求消耗数均为200,请问要修改哪些地方处理并发问题,或者我遗漏了什么关键点?

BucketToken.java

package algo.interviewquestions;

public class BucketToken {
    private Integer maxBucketSize;
    private Integer tokensAvailable;
    private Long refillBucketRate;
    private Long nextRefillTime;
    private Long totalRequestsReceived;
    private Long totalRequestsConsumed;
    private Integer bucketRefillCount;

    public BucketToken(Integer maxBucketSize, Long refillBucketRate) {
        this.maxBucketSize = maxBucketSize;
        this.refillBucketRate = refillBucketRate;
        this.nextRefillTime = System.currentTimeMillis() + refillBucketRate;
        this.tokensAvailable = maxBucketSize;
        this.totalRequestsReceived = 0L;
        this.totalRequestsConsumed = 0L;
        this.bucketRefillCount = 0;
        refill();
    }

    public Long getTotalRequestsConsumed() {
        return totalRequestsConsumed;
    }

    public Integer getTokensAvailable() {
        return tokensAvailable;
    }

    public Integer getBucketRefillCount() {
        return bucketRefillCount;
    }

    public Long getTotalRequestsReceived() {
        return totalRequestsReceived;
    }

     public synchronized boolean tryConsume() {
        totalRequestsReceived++;
        refill();
        if (this.tokensAvailable > 0) {
            this.tokensAvailable --;
            this.totalRequestsConsumed++;
            return true;
        }
        return false;
    }

    private synchronized void refill() {
        if (System.currentTimeMillis() < this.nextRefillTime) {
            return;
        }
        this.nextRefillTime = System.currentTimeMillis() + this.refillBucketRate;
        this.tokensAvailable = Math.max(this.maxBucketSize, this.tokensAvailable);
        this.bucketRefillCount++;
    }

}

BucketTokenTest.java

package algo.interviewquestions;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.IntStream;

public class BucketTokenTest {
    public static void main(String[] args) throws InterruptedException {
        BucketToken bucketToken = new BucketToken(20, 1000L);
        AtomicInteger totalRequestsConsumed1 = runInMultiThreadedApproach(bucketToken);

        BucketToken bucketToken2 = new BucketToken(20, 1000L);
        AtomicInteger totalRequestsConsumed2 = runWithSingleThread(bucketToken2);

        printRequestsStats(bucketToken, totalRequestsConsumed1);
        printRequestsStats(bucketToken2, totalRequestsConsumed2);
    }

    private static AtomicInteger runWithSingleThread(BucketToken bucketToken) {
        AtomicInteger totalRequestsConsumed = new AtomicInteger(0);
        long startTime = System.currentTimeMillis();
        while (System.currentTimeMillis() - startTime < 10000L) {
            if (bucketToken.tryConsume()) {
                totalRequestsConsumed.incrementAndGet();
                System.out.println(Thread.currentThread().getName() + " - Request Accepted");
                continue;
            }
            System.out.println(Thread.currentThread().getName() + " - Request Declined");
        }
        return totalRequestsConsumed;
    }

    private static void printRequestsStats(BucketToken bucketToken, AtomicInteger totalRequestsConsumed) {
        System.out.println("Total Requests Received = " + bucketToken.getTotalRequestsReceived());
        System.out.println("Total Requests Accepted = " + totalRequestsConsumed.get());
        System.out.println("Total Requests Accepted = " + bucketToken.getTotalRequestsConsumed());
        System.out.println("Bucket refill count = " + bucketToken.getBucketRefillCount());
    }

    private static AtomicInteger runInMultiThreadedApproach(BucketToken bucketToken) throws InterruptedException {
        int threadCount = 5;
        ExecutorService executorService = Executors.newFixedThreadPool(threadCount);
        AtomicInteger totalRequestsConsumed = new AtomicInteger(0);
        long startTime = System.currentTimeMillis();
        CountDownLatch latch = new CountDownLatch(threadCount);
        IntStream.rangeClosed(1,threadCount).boxed().forEach(t -> executorService.submit(() -> {
            while (System.currentTimeMillis() - startTime < 10000L) {
                if (bucketToken.tryConsume()) {
                    totalRequestsConsumed.incrementAndGet();
                    System.out.println(Thread.currentThread().getName() + " - Request Accepted");
                    continue;
                }
                System.out.println(Thread.currentThread().getName() + " - Request Declined");
            }
            latch.countDown();
        }));
        long startTime2 = System.currentTimeMillis();
        latch.await();
        System.out.println("Completed in " + (System.currentTimeMillis() - startTime2) + "ms");
        executorService.shutdown();
        return totalRequestsConsumed;
    }
}
问题分析与解决方案

你的核心问题出在令牌桶的补充逻辑和多线程下的时间窗口控制上,具体原因和修复方案如下:

1. 核心问题点

(1)令牌补充逻辑错误

当前refill()方法直接将tokensAvailable设为maxBucketSize,但正确的令牌桶算法应该是按时间间隔补充固定数量的令牌,而非直接填满桶。你设置的20请求/秒意味着每秒补充20个令牌,直接填满桶会导致剩余令牌被覆盖后超发。

(2)多线程下的时间触发误差

nextRefillTime每次设置为当前时间 + refillBucketRate,多线程环境下可能因系统时间微小差异,导致10秒周期内多触发一次补充,额外产生20个令牌(220-200的差值正好是单次补充量)。

(3)测试循环的时间边界问题

多线程测试中,每个线程的System.currentTimeMillis() - startTime < 10000L判断,可能因线程调度延迟,导致部分线程在10秒后仍执行一小段时间,消耗额外令牌。

2. 具体修复步骤

(1)修正令牌补充逻辑

改为按时间差计算应补充的令牌数,避免直接填满桶:

private synchronized void refill() {
    long now = System.currentTimeMillis();
    if (now < nextRefillTime) {
        return;
    }
    // 计算流逝的时间间隔数
    long intervalCount = (now - nextRefillTime) / refillBucketRate;
    if (intervalCount <= 0) {
        return;
    }
    // 按间隔数补充令牌,不超过桶上限
    int tokensToAdd = (int) (intervalCount * maxBucketSize);
    tokensAvailable = Math.min(maxBucketSize, tokensAvailable + tokensToAdd);
    // 更新下次补充时间,适配多间隔情况
    nextRefillTime += intervalCount * refillBucketRate;
    bucketRefillCount += intervalCount;
}

(2)统一测试时间边界

将循环判断改为基于固定结束时间,避免线程调度导致超时执行:

// 单线程测试修改
private static AtomicInteger runWithSingleThread(BucketToken bucketToken) {
    AtomicInteger totalRequestsConsumed = new AtomicInteger(0);
    long endTime = System.currentTimeMillis() + 10000L;
    while (System.currentTimeMillis() < endTime) {
        if (bucketToken.tryConsume()) {
            totalRequestsConsumed.incrementAndGet();
            System.out.println(Thread.currentThread().getName() + " - Request Accepted");
        } else {
            System.out.println(Thread.currentThread().getName() + " - Request Declined");
        }
    }
    return totalRequestsConsumed;
}

// 多线程测试修改
private static AtomicInteger runInMultiThreadedApproach(BucketToken bucketToken) throws InterruptedException {
    int threadCount = 5;
    ExecutorService executorService = Executors.newFixedThreadPool(threadCount);
    AtomicInteger totalRequestsConsumed = new AtomicInteger(0);
    long endTime = System.currentTimeMillis() + 10000L;
    CountDownLatch latch = new CountDownLatch(threadCount);
    IntStream.rangeClosed(1,threadCount).boxed().forEach(t -> executorService.submit(() -> {
        while (System.currentTimeMillis() < endTime) {
            if (bucketToken.tryConsume()) {
                totalRequestsConsumed.incrementAndGet();
                System.out.println(Thread.currentThread().getName() + " - Request Accepted");
            } else {
                System.out.println(Thread.currentThread().getName() + " - Request Declined");
            }
        }
        latch.countDown();
    }));
    long startTime2 = System.currentTimeMillis();
    latch.await();
    System.out.println("Completed in " + (System.currentTimeMillis() - startTime2) + "ms");
    executorService.shutdown();
    return totalRequestsConsumed;
}

(3)移除构造函数的多余调用

初始化时令牌已经是maxBucketSize,无需立即调用refill():

public BucketToken(Integer maxBucketSize, Long refillBucketRate) {
    this.maxBucketSize = maxBucketSize;
    this.refillBucketRate = refillBucketRate;
    this.nextRefillTime = System.currentTimeMillis() + refillBucketRate;
    this.tokensAvailable = maxBucketSize;
    this.totalRequestsReceived = 0L;
    this.totalRequestsConsumed = 0L;
    this.bucketRefillCount = 0;
    // 移除此处的refill()调用
}

3. 修复后的验证

修改后,单线程和多线程测试在10秒内的请求消耗数会稳定在200左右(系统时间精度导致的±1误差属于正常情况),符合限流预期。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:05:02