单机限流到分布式限流的演进需求
API网关或微服务网关在多实例部署时,单机限流无法满足全局限流要求。例如业务约定某个API的QPS上限为1000,部署4个实例后每台限流250QPS并不能保证全局QPS不超过1000——某台实例流量突增时可能独占更多配额。分布式限流的核心目标是在多节点间公平分配限流配额,并保证全局限流阈值不被突破。
常见的分布式限流方案有三类:集中式计数器(Redis原子操作)、令牌桶集群同步(gossip协议或中心协调)、滑动窗口日志(Redis Sorted Set)。集中式方案最简洁,令牌桶集群同步在超大规模场景延迟更低。
基于Redis的Token Bucket实现
Redis Token Bucket利用Lua脚本保证令牌发放的原子性,避免多实例竞争导致的超卖:
-- redis_token_bucket.lua
local key = KEYS[1]
local max_tokens = tonumber(ARGV[1])
local refill_rate = tonumber(ARGV[2])
local requested = tonumber(ARGV[3])
local now = tonumber(ARGV[4])
local data = redis.call('HMGET', key, 'tokens', 'last_refill')
local tokens = tonumber(data[1]) or max_tokens
local last_refill = tonumber(data[2]) or now
local elapsed = now - last_refill
local refill = elapsed * refill_rate / 1000
tokens = math.min(max_tokens, tokens + refill)
local allowed = 0
if tokens >= requested then
tokens = tokens - requested
allowed = 1
end
redis.call('HMSET', key, 'tokens', tokens, 'last_refill', now)
redis.call('PEXPIRE', key, max_tokens / refill_rate * 1000 * 2)
return { allowed, tokens }
Go调用端:
func (l *RedisTokenBucket) Allow(ctx context.Context, key string, requested float64) (bool, float64, error) {
now := time.Now().UnixMilli()
result, err := l.luaScript.Run(ctx, l.redis,
[]string{"ratelimit:" + key},
l.maxTokens, l.refillRate, requested, now,
).Slice()
if err != nil {
return false, 0, err
}
allowed := result[0].(int64) == 1
remaining := result[1].(float64)
return allowed, remaining, nil
}
集群Token Bucket同步方案
当QPS达到数十万级时,Redis中心化方案的网络延迟成为瓶颈。集群Token Bucket方案将令牌池预分配给各节点,节点本地独立决策,周期性同步:
type ClusterTokenBucket struct {
localBucket *LocalTokenBucket
nodeID string
syncInterval time.Duration
redis *redis.Client
totalTokens float64
}
func (c *ClusterTokenBucket) Start(ctx context.Context) {
ticker := time.NewTicker(c.syncInterval)
for {
select {
case <-ticker.C:
c.syncTokens(ctx)
case <-ctx.Done():
return
}
}
}
func (c *ClusterTokenBucket) syncTokens(ctx context.Context) {
consumed := c.localBucket.GetConsumed()
c.redis.HIncrBy(ctx, "cluster:ratelimit:consumed", c.nodeID, int64(consumed))
totalConsumed, _ := c.redis.HGetAll(ctx, "cluster:ratelimit:consumed").Result()
globalUsed := sumValues(totalConsumed)
remaining := c.totalTokens - float64(globalUsed)
nodeCount := len(totalConsumed)
perNodeQuota := remaining / float64(nodeCount)
c.localBucket.SetQuota(perNodeQuota)
}
同步间隔越长,节点间配额偏差越大,但Redis访问压力越低。生产环境建议100ms-1s的同步间隔,在精确性和性能间取得平衡。
限流响应与降级策略
被限流的请求不应直接返回500,而应携带标准限流响应头:
func RateLimitMiddleware(bucket *RedisTokenBucket) gin.HandlerFunc {
return func(c *gin.Context) {
allowed, remaining, _ := bucket.Allow(c.Request.Context(), c.ClientIP(), 1)
c.Header("X-RateLimit-Remaining", fmt.Sprintf("%.0f", remaining))
if !allowed {
c.Header("Retry-After", "1")
c.JSON(http.StatusTooManyRequests, gin.H{
"error": "rate_limit_exceeded",
"message": "请求频率超过限制,请稍后重试",
})
c.Abort()
return
}
c.Next()
}
}
对于关键业务路径,被限流时可降级到缓存数据或只读模式,而非直接拒绝服务。降级策略应在限流配置中预先定义,与限流阈值一同管理。
限流指标监控与动态调整
分布式限流上线后必须持续监控关键指标:每秒限流拒绝数(rate_limit_rejected_total)、平均令牌剩余率(remaining/max_tokens)、Redis Lua脚本执行延迟P99。
动态调整方面,可基于历史流量模式实现弹性限流——工作日高峰期适当放宽阈值,低谷期收紧。调整逻辑通过配置中心下发,限流组件热加载,无需重启服务。
func (l *Limiter) WatchConfig(ctx context.Context, configKey string) {
ch := l.configClient.Watch(ctx, configKey)
for newCfg := range ch {
l.mu.Lock()
l.maxTokens = newCfg.MaxTokens
l.refillRate = newCfg.RefillRate
l.mu.Unlock()
log.Printf("限流配置已更新: max=%v, rate=%v",
newCfg.MaxTokens, newCfg.RefillRate)
}
}
分布式限流不是配置一次就结束的工作,持续监控和动态调整才能让限流策略既保护后端服务又不过度牺牲业务流量。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-fen-bu-shi-xian-liu-shi-xian-tokenbucket-ji-qun/