NOTE
RateLimiter
1. What it is 2. Usage 3. Source analysis 4. Summary 5. References
This is a historical learning note and may contain outdated or incomplete understanding.
1. What It Is
A rate-limiting tool implemented by Google based on the token-bucket algorithm.
2. Usage
RateLimiter limiter = RateLimiter.create(10);
for (int i = 0; i < 100; i++)
{
limiter.acquire();
System.out.println(i + ":get lock");
}
3. Source Analysis
3.1. Creating a RateLimiter
- RateLimiter
public static RateLimiter create(double permitsPerSecond) {
return create(SleepingStopwatch.createFromSystemTimer(), permitsPerSecond);
}
static RateLimiter create(SleepingStopwatch stopwatch, double permitsPerSecond) {
// SmoothBursty is used by default (token generation rate is constant).
RateLimiter rateLimiter = new SmoothBursty(stopwatch, 1.0 /* maxBurstSeconds */);
// Set the number of tokens generated per second.
rateLimiter.setRate(permitsPerSecond);
return rateLimiter;
}
public final void setRate(double permitsPerSecond) {
checkArgument(
permitsPerSecond > 0.0 && !Double.isNaN(permitsPerSecond), "rate must be positive");
// SmoothRateLimiter.doSetRate
synchronized (mutex()) {
doSetRate(permitsPerSecond, stopwatch.readMicros());
}
}
- SmoothBursty
/**
* Current number of stored permits.
*/
double storedPermits;
/**
* Maximum number of stored permits.
*/
double maxPermits;
/**
* Interval for adding permits.
*/
double stableIntervalMicros;
/**
* Earliest time when the next request can obtain permits.
* Because RateLimiter allows permit pre-consumption, after the previous request
* pre-consumes permits, the next request has to wait until nextFreeTicketMicros.
*/
private long nextFreeTicketMicros = 0L;
SmoothBursty(SleepingStopwatch stopwatch, double maxBurstSeconds) {
super(stopwatch);
this.maxBurstSeconds = maxBurstSeconds;
}
- SmoothRateLimiter
final void doSetRate(double permitsPerSecond, long nowMicros) {
resync(nowMicros);
double stableIntervalMicros = SECONDS.toMicros(1L) / permitsPerSecond;//// calculate the maximum stored permit count
this.stableIntervalMicros = stableIntervalMicros;
doSetRate(permitsPerSecond, stableIntervalMicros);
}
3.2. Acquiring a Permit
- RateLimiter
public double acquire() {
// Acquire one permit by default.
return acquire(1);
}
public double acquire(int permits) {
long microsToWait = reserve(permits);
// Wait for the corresponding amount of time.
stopwatch.sleepMicrosUninterruptibly(microsToWait);
return 1.0 * microsToWait / SECONDS.toMicros(1L);
}
final long reserve(int permits) {
// Check that permits is positive.
checkPermits(permits);
// Not thread-safe, so lock it.
synchronized (mutex()) {
// Acquire the requested number of permits.
return reserveAndGetWaitLength(permits, stopwatch.readMicros());
}
}
final long reserveAndGetWaitLength(int permits, long nowMicros) {
// SmoothRateLimiter.reserveEarliestAvailable
long momentAvailable = reserveEarliestAvailable(permits, nowMicros);
// Return the required wait time.
return max(momentAvailable - nowMicros, 0);
}
- SmoothRateLimiter
final long reserveEarliestAvailable(int requiredPermits, long nowMicros) {
// Generate permits only when permits are requested.
resync(nowMicros);
long returnValue = nextFreeTicketMicros;// Return the previously calculated nextFreeTicketMicros.
// The current request pays for the previous request's pre-consumption.
// This is also why RateLimiter can pre-consume permits to handle bursts.
// To disable pre-consumption, return the updated nextFreeTicketMicros here instead.
double storedPermitsToSpend = min(requiredPermits, this.storedPermits);// Number of stored permits that can be consumed.
double freshPermits = requiredPermits - storedPermitsToSpend;// Additional permits still needed.
long waitMicros = storedPermitsToWaitTime(this.storedPermits, storedPermitsToSpend)
+ (long) (freshPermits * stableIntervalMicros);// Calculate wait time from freshPermits.
this.nextFreeTicketMicros = nextFreeTicketMicros + waitMicros;// The nextFreeTicketMicros calculated this time is not returned.
this.storedPermits -= storedPermitsToSpend;// Subtract the permits being distributed from the current stored permit count.
return returnValue;
}
resync
This uses the token-bucket algorithm. Tokens are put into the bucket at a certain rate, and a thread can execute only after obtaining a token.
Who continuously generates tokens and puts them into the bucket?
- Start a scheduled task. Problem: this consumes too many resources.
- Lazy calculation. Before obtaining a token, call a function to check whether tokens need to be added to the bucket.
void resync(long nowMicros) {
// if nextFreeTicket is in the past, resync to now
// Current time > the next time at which permits are allowed to be generated.
// This means additional permits can now be obtained.
if (nowMicros > nextFreeTicketMicros) {
// Calculate the number of permits that can be generated during this time period.
double newPermits = (nowMicros - nextFreeTicketMicros) / coolDownIntervalMicros();
// Limit the maximum number of permits.
storedPermits = min(maxPermits, storedPermits + newPermits);
// Set the next allowed permit-generation time to the current time.
nextFreeTicketMicros = nowMicros;
}
}
4. Summary
- Lock.
- Generate permits.
- Calculate the number of permits from the elapsed time since the previous generation point and store them. The number cannot exceed the value configured by
create.
- Calculate the number of permits from the elapsed time since the previous generation point and store them. The number cannot exceed the value configured by
- Consume permits.
- Use the current number of stored permits and the requested number of permits to calculate how many additional permits are needed, then calculate the wait time and update it (the next request must wait for this duration).
- Wait for the duration calculated by the previous request.
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub