Go语言分布式限流器设计与Token Bucket算法工程实战

限流算法对比与Token Bucket选型理由

常见的限流算法有固定窗口计数器、滑动窗口日志、滑动窗口计数器和令牌桶四种。固定窗口在边界处存在两倍流量的突刺问题;滑动窗口日志精度高但内存消耗与请求量线性增长;滑动窗口计数器在窗口边界有少量误差但实现简单。

令牌桶(Token Bucket)是工业场景的首选算法。它以恒定速率向桶中放入令牌,请求到达时从桶中取走令牌,桶满则丢弃多余令牌,桶空则拒绝请求。令牌桶允许一定程度的突发流量——桶中积累的令牌可以在瞬间被消耗,适合有流量波动的API场景。

Go语言的time.Ticker和channel天然适配令牌桶模型:Ticker控制令牌注入速率,channel的缓冲区容量就是桶大小。

单机令牌桶限流器实现

package ratelimit

import (
    "context"
    "sync"
    "time"
)

type TokenBucket struct {
    mu       sync.Mutex
    rate     float64   // 每秒令牌注入速率
    burst    int       // 桶容量(最大突发量)
    tokens   float64   // 当前令牌数
    lastTime time.Time // 上次更新时间
}

func NewTokenBucket(rate float64, burst int) *TokenBucket {
    return &TokenBucket{
        rate:     rate,
        burst:    burst,
        tokens:   float64(burst),
        lastTime: time.Now(),
    }
}

func (tb *TokenBucket) Allow() bool {
    tb.mu.Lock()
    defer tb.mu.Unlock()
    now := time.Now()
    elapsed := now.Sub(tb.lastTime).Seconds()
    tb.lastTime = now
    tb.tokens += elapsed * tb.rate
    if tb.tokens > float64(tb.burst) {
        tb.tokens = float64(tb.burst)
    }
    if tb.tokens >= 1 {
        tb.tokens--
        return true
    }
    return false
}

func (tb *TokenBucket) Wait(ctx context.Context) error {
    for {
        if tb.Allow() {
            return nil
        }
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-time.After(time.Millisecond * 50):
        }
    }
}

这种惰性计算(lazy refill)方式不需要后台协程持续注入令牌,只在每次Allow调用时根据时间差补充,实现简单且零额外开销。

分布式限流:Redis + Lua原子脚本

单机限流在多实例部署时无法协调配额。分布式限流需要一个共享存储来维护全局令牌状态。Redis的原子Lua脚本是最常用的方案:

var tokenBucketScript = redis.NewScript(`
    local key = KEYS[1]
    local rate = tonumber(ARGV[1])
    local burst = tonumber(ARGV[2])
    local now = tonumber(ARGV[3])
    local requested = tonumber(ARGV[4])
    
    local info = redis.call("HMGET", key, "tokens", "last_time")
    local tokens = tonumber(info[1])
    local last_time = tonumber(info[2])
    
    if tokens == nil then
        tokens = burst
        last_time = now
    end
    
    local elapsed = now - last_time
    tokens = math.min(burst, tokens + elapsed * rate)
    last_time = now
    
    local allowed = 0
    if tokens >= requested then
        tokens = tokens - requested
        allowed = 1
    end
    
    redis.call("HMSET", key, "tokens", tokens, "last_time", last_time)
    redis.call("EXPIRE", key, math.ceil(burst / rate) + 10)
    
    return { allowed, tokens }
`)

Lua脚本在Redis中原子执行,不会出现并发竞争。HMSET更新令牌数和时间戳,EXPIRE设置过期时间自动清理不活跃的key。

Go端调用:

func (rl *RedisRateLimiter) Allow(ctx context.Context, key string) (bool, error) {
    now := time.Now().UnixMilli() / 1000
    result, err := tokenBucketScript.Run(ctx, rl.client,
        []string{key}, rl.rate, rl.burst, now, 1).Slice()
    if err != nil {
        return false, err
    }
    return result[0].(int64) == 1, nil
}

滑动窗口限流与分布式配额分配

令牌桶允许突发,某些场景需要严格限制窗口内总请求数。滑动窗口限流用Sorted Set记录每个请求的时间戳,窗口滑动时淘汰过期记录,剩余数量即为当前窗口已用配额:

var slidingWindowScript = redis.NewScript(`
    local key = KEYS[1]
    local limit = tonumber(ARGV[1])
    local window = tonumber(ARGV[2])
    local now = tonumber(ARGV[3])
    
    redis.call("ZREMRANGEBYSCORE", key, 0, now - window)
    local count = redis.call("ZCARD", key)
    
    if count < limit then
        redis.call("ZADD", key, now, now .. ":" .. math.random(1000000))
        redis.call("EXPIRE", key, window / 1000 + 1)
        return 1
    end
    return 0
`)

Sorted Set中member用时间戳+随机数拼接避免重复。ZREMRANGEBYSCORE按分数范围删除过期成员,时间复杂度O(logN + M)。当QPS很高时Sorted Set内存开销较大,可以结合计数器分片优化。

多级限流架构与熔断联动

生产环境中限流与熔断配合使用。当限流拒绝率超过阈值时触发熔断器打开,直接快速失败不再尝试回源,避免下游压力持续累积:

type CircuitBreaker struct {
    mu           sync.Mutex
    state        State // Closed, Open, HalfOpen
    failures     int
    threshold    int
    resetTimeout time.Duration
    lastOpen     time.Time
}

func (cb *CircuitBreaker) Execute(fn func() error) error {
    cb.mu.Lock()
    if cb.state == Open {
        if time.Since(cb.lastOpen) > cb.resetTimeout {
            cb.state = HalfOpen
            cb.failures = 0
        } else {
            cb.mu.Unlock()
            return ErrCircuitOpen
        }
    }
    cb.mu.Unlock()
    
    err := fn()
    cb.mu.Lock()
    defer cb.mu.Unlock()
    
    if err != nil {
        cb.failures++
        if cb.failures >= cb.threshold {
            cb.state = Open
            cb.lastOpen = time.Now()
        }
        return err
    }
    
    if cb.state == HalfOpen {
        cb.state = Closed
    }
    cb.failures = 0
    return nil
}

限流器与熔断器的组合形成多级防护:限流器控制正常流量在配额内,熔断器在下游异常时快速切断请求流。这种分层防御策略在微服务网关中已是标配,Spring Cloud Gateway和Istio都内置了类似机制。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-fen-bu-shi-xian-liu-qi-she-ji-yu-tokenbucket-suan/

(0)
小编小编
上一篇 5小时前
下一篇 5小时前

相关推荐