基于令牌桶算法的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
相关产品推荐
相关产品推荐

