Rate Limiter β Token Bucket Implementation
Implementasi rate limiter token bucket di Go: thread-safe, memory-efficient, dengan cleanup goroutine. Lengkap dengan penjelasan algoritma dan trade-off.
Rate Limiter β Token Bucket
Konsep
Token bucket adalah algoritma rate limiting yang paling banyak dipakai karena mengizinkan burst (lonjakan traffic) sambil tetap membatasi average rate.
Analogi: Ember bocor.
- Ember punya kapasitas maksimum (max tokens)
- Token ditambahkan dengan rate konstan (refill rate)
- Setiap request butuh 1 token
- Kalau token habis β request ditolak
- Kalau gak ada request β token menumpuk sampai kapasitas max
Kenapa Token Bucket?
| Algoritma | Burst Support | Memory | Implementasi |
|---|---|---|---|
| Fixed Window | β (reset di batas window) | O(1) | Termudah |
| Sliding Window Log | β | O(n) | Berat |
| Sliding Window Counter | β | O(1) | Medium |
| Token Bucket | β | O(1) | Medium |
| Leaky Bucket | β (smooth output) | O(1) | Medium |
Token bucket menang karena: burst support + O(1) memory + implementasi gampang.
Implementasi Go
package ratelimit
import (
"sync"
"time"
)
// TokenBucket implements token bucket rate limiting.
// Thread-safe. Suitable for per-user or per-IP rate limiting.
type TokenBucket struct {
mu sync.Mutex
rate float64 // tokens per second
burst float64 // max tokens (bucket capacity)
tokens float64 // current token count
lastRefill time.Time // last refill timestamp
}
// NewTokenBucket creates a new token bucket.
// rate: tokens per second (e.g., 10 = 10 requests/sec)
// burst: max burst size (e.g., 20 = allow 20 requests at once)
func NewTokenBucket(rate, burst float64) *TokenBucket {
return &TokenBucket{
rate: rate,
burst: burst,
tokens: burst, // start with full bucket
lastRefill: time.Now(),
}
}
// Allow checks if one request is allowed.
// Returns true if allowed, false if rate limited.
func (tb *TokenBucket) Allow() bool {
return tb.AllowN(1)
}
// AllowN checks if n requests are allowed.
func (tb *TokenBucket) AllowN(n float64) bool {
tb.mu.Lock()
defer tb.mu.Unlock()
// Refill tokens based on elapsed time
now := time.Now()
elapsed := now.Sub(tb.lastRefill).Seconds()
tb.tokens += elapsed * tb.rate
// Cap at burst
if tb.tokens > tb.burst {
tb.tokens = tb.burst
}
tb.lastRefill = now
// Check if enough tokens
if tb.tokens >= n {
tb.tokens -= n
return true
}
return false
}
// Tokens returns current token count (for monitoring).
func (tb *TokenBucket) Tokens() float64 {
tb.mu.Lock()
defer tb.mu.Unlock()
return tb.tokens
}Usage di HTTP Middleware
type RateLimiter struct {
buckets sync.Map // map[string]*TokenBucket β per-user/IP buckets
rate float64
burst float64
}
func NewRateLimiter(rate, burst float64) *RateLimiter {
return &RateLimiter{rate: rate, burst: burst}
}
func (rl *RateLimiter) Middleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
key := r.Header.Get("X-API-Key") // or r.RemoteAddr
if key == "" {
key = r.RemoteAddr
}
bucket, _ := rl.buckets.LoadOrStore(key, NewTokenBucket(rl.rate, rl.burst))
tb := bucket.(*TokenBucket)
if !tb.Allow() {
w.Header().Set("Retry-After", "1")
w.Header().Set("X-RateLimit-Limit", fmt.Sprintf("%.0f", rl.burst))
http.Error(w, `{"error":"rate limit exceeded"}`, http.StatusTooManyRequests)
return
}
next.ServeHTTP(w, r)
})
}Cleanup β Mencegah Memory Leak
Tanpa cleanup, sync.Map bakal terus tumbuh karena setiap IP baru bikin bucket baru yang gak pernah dihapus.
func (rl *RateLimiter) StartCleanup(ctx context.Context, interval time.Duration) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
rl.cleanup()
case <-ctx.Done():
return
}
}
}
func (rl *RateLimiter) cleanup() {
rl.buckets.Range(func(key, value interface{}) bool {
tb := value.(*TokenBucket)
// Remove buckets that haven't been used in 1 hour
tb.mu.Lock()
idle := time.Since(tb.lastRefill)
tb.mu.Unlock()
if idle > 1*time.Hour {
rl.buckets.Delete(key)
}
return true
})
}Test
func TestTokenBucket_Allow(t *testing.T) {
tb := NewTokenBucket(10, 20) // 10 req/s, burst 20
// Burst: all 20 should pass
for i := 0; i < 20; i++ {
assert.True(t, tb.Allow(), "burst request %d", i)
}
// Rate limited
assert.False(t, tb.Allow(), "should be rate limited")
// Wait for refill
time.Sleep(200 * time.Millisecond) // 10 * 0.2 = 2 tokens
assert.True(t, tb.Allow())
assert.True(t, tb.Allow())
assert.False(t, tb.Allow())
}
func TestTokenBucket_Concurrent(t *testing.T) {
tb := NewTokenBucket(1000, 1000)
var wg sync.WaitGroup
allowed := atomic.Int64{}
for i := 0; i < 1000; i++ {
wg.Add(1)
go func() {
defer wg.Done()
if tb.Allow() {
allowed.Add(1)
}
}()
}
wg.Wait()
assert.Equal(t, int64(1000), allowed.Load())
}Benchmark
func BenchmarkTokenBucket_Allow(b *testing.B) {
tb := NewTokenBucket(1e6, 1e6) // basically unlimited
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
tb.Allow()
}
})
}
// Result: ~50ns/op, ~20M ops/sec β mutex overhead minimalInterview Talking Points
"Kenapa token bucket, bukan fixed window?" Fixed window punya masalah di batas window:
- Window 1 (detik 0.0-1.0): 10 request di 0.9s β OK
- Window 2 (detik 1.0-2.0): 10 request di 1.1s β OK
- Total: 20 request dalam 0.2 detik! (0.9s - 1.1s)
Ini bukan rate limiting yang efektif. Token bucket mencegah ini karena token diisi secara kontinyu, bukan per-window.
"Bagaimana scaling rate limiter di multi-server?"
- Sticky session: request user yang sama selalu ke server yang sama β local rate limiter cukup
- Redis: centralized counter β Redis
INCR+EXPIRE, bisa pakai Lua script buat atomic - Envoy/L7 proxy: rate limiting di edge, sebelum request masuk ke service
Last updated on