限流是高并发系统中保护服务稳定性的第一道防线。令牌桶算法相比固定窗口和滑动窗口,支持突发流量且实现简洁,是API接口规范中最常用的限流策略。本文使用Go语言 + Redis + Lua脚本实现分布式令牌桶限流器,确保多实例部署下的并发安全和性能。
令牌桶算法原理与参数设计
令牌桶核心机制:以固定速率向桶中添加令牌,桶有最大容量上限。每个请求消耗一个令牌,桶空时拒绝请求或排队等待。关键参数有两个:rate(令牌生成速率,个/秒)和capacity(桶最大容量)。
rate决定平均处理速率,capacity决定允许的突发流量峰值。例如rate=100、capacity=200表示平均每秒处理100个请求,但短时间内可容忍200个请求的突发。参数设置需根据后端实际处理能力和SLA目标调整。
分布式限流架构设计
单机限流无法应对多实例部署场景。5个服务实例各自限流,总流量是单机限制的5倍,达不到全局限流效果。分布式限流将令牌桶状态存储在Redis中,所有实例共享同一桶状态,通过Lua脚本保证读取-计算-写入的原子性。
架构选型理由:Redis单线程模型天然解决并发竞态;Lua脚本在Redis中原子执行,中间不会被其他命令打断;Go的redis客户端性能优异,单次限流检查延迟在1ms以内。
Lua脚本核心实现
Lua脚本是分布式令牌桶的关键。所有逻辑在Redis服务端原子执行,避免网络往返导致的竞态条件。
-- rate_limiter.lua
-- KEYS[1]: 限流key
-- ARGV[1]: 当前时间戳(毫秒)
-- ARGV[2]: 令牌生成速率(个/秒)
-- ARGV[3]: 桶最大容量
-- ARGV[4]: 每次请求消耗的令牌数
local key = KEYS[1]
local now = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local capacity = tonumber(ARGV[3])
local requested = tonumber(ARGV[4])
local bucket = redis.call('HMGET', key, 'tokens', 'timestamp')
local tokens = tonumber(bucket[1]) or capacity
local last_time = tonumber(bucket[2]) or now
local elapsed = math.max(0, now - last_time)
local new_tokens = (elapsed / 1000) * rate
tokens = math.min(capacity, tokens + new_tokens)
local allowed = 0
if tokens >= requested then
tokens = tokens - requested
allowed = 1
end
local ttl = math.ceil(capacity / rate * 2)
redis.call('HMSET', key, 'tokens', tokens, 'timestamp', now)
redis.call('EXPIRE', key, ttl)
return {allowed, tokens}
Go语言限流器封装
将Lua脚本封装为可复用的Go包,支持多种限流维度(用户ID、API路径、IP等)。
package ratelimiter
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
type TokenBucket struct {
client *redis.Client
script *redis.Script
rate float64
capacity float64
}
func NewTokenBucket(client *redis.Client, rate, capacity float64) *TokenBucket {
return &TokenBucket{
client: client,
script: redis.NewScript(luaScript),
rate: rate,
capacity: capacity,
}
}
func (tb *TokenBucket) Allow(ctx context.Context, key string) (allowed bool, remaining float64, err error) {
return tb.AllowN(ctx, key, 1)
}
func (tb *TokenBucket) AllowN(ctx context.Context, key string, n int) (allowed bool, remaining float64, err error) {
now := float64(time.Now().UnixMilli())
result, err := tb.script.Run(ctx, tb.client,
[]string{fmt.Sprintf("rate_limit:%s", key)},
now, tb.rate, tb.capacity, n,
).Slice()
if err != nil {
return false, 0, fmt.Errorf("限流器执行失败: %w", err)
}
allowed = result[0].(int64) == 1
remaining = result[1].(float64)
return allowed, remaining, nil
}
const luaScript = `
local key = KEYS[1]
local now = tonumber(ARGV[1])
local rate = tonumber(ARGV[2])
local capacity = tonumber(ARGV[3])
local requested = tonumber(ARGV[4])
local bucket = redis.call('HMGET', key, 'tokens', 'timestamp')
local tokens = tonumber(bucket[1]) or capacity
local last_time = tonumber(bucket[2]) or now
local elapsed = math.max(0, now - last_time)
local new_tokens = (elapsed / 1000) * rate
tokens = math.min(capacity, tokens + new_tokens)
local allowed = 0
if tokens >= requested then
tokens = tokens - requested
allowed = 1
end
local ttl = math.ceil(capacity / rate * 2)
redis.call('HMSET', key, 'tokens', tokens, 'timestamp', now)
redis.call('EXPIRE', key, ttl)
return {allowed, tokens}
`
基于限流器的API中间件
将限流器封装为HTTP中间件,对不同维度的请求进行限流。Go中通过中间件函数实现灵活的限流控制。
package middleware
import (
"context"
"net/http"
"strconv"
"yourapp/ratelimiter"
"github.com/gin-gonic/gin"
)
type RateLimitMiddleware struct {
userLimiter *ratelimiter.TokenBucket
ipLimiter *ratelimiter.TokenBucket
globalLimiter *ratelimiter.TokenBucket
}
func NewRateLimitMiddleware(rdb *redis.Client) *RateLimitMiddleware {
return &RateLimitMiddleware{
userLimiter: ratelimiter.NewTokenBucket(rdb, 100, 200),
ipLimiter: ratelimiter.NewTokenBucket(rdb, 50, 100),
globalLimiter: ratelimiter.NewTokenBucket(rdb, 10000, 15000),
}
}
func (m *RateLimitMiddleware) Handle() gin.HandlerFunc {
return func(c *gin.Context) {
ctx := context.Background()
// 全局限流
if ok, _, err := m.globalLimiter.Allow(ctx, "global"); err != nil || !ok {
c.Header("Retry-After", "1")
c.JSON(http.StatusTooManyRequests, gin.H{
"code": 429, "message": "服务繁忙,请稍后再试",
})
c.Abort()
return
}
// IP限流
clientIP := c.ClientIP()
if ok, _, err := m.ipLimiter.Allow(ctx, "ip:"+clientIP); err != nil || !ok {
c.JSON(http.StatusTooManyRequests, gin.H{
"code": 429, "message": "请求过于频繁",
})
c.Abort()
return
}
// 用户限流(已认证用户)
if userID, exists := c.Get("userID"); exists {
key := "user:" + strconv.Itoa(userID.(int))
if ok, remaining, err := m.userLimiter.Allow(ctx, key); err != nil || !ok {
c.Header("X-RateLimit-Remaining", strconv.FormatFloat(remaining, 'f', 0, 64))
c.JSON(http.StatusTooManyRequests, gin.H{
"code": 429, "message": "操作过于频繁,请稍后再试",
})
c.Abort()
return
}
c.Header("X-RateLimit-Remaining", strconv.FormatFloat(remaining, 'f', 0, 64))
}
c.Next()
}
}
消息中间件配合令牌桶处理排队
直接拒绝请求体验不好,可用消息中间件将超限请求排队异步处理。令牌桶作为准入控制,超出速率的请求写入Kafka或RabbitMQ队列,消费者按桶速率处理。
// 异步限流处理:队列消费者
func StartQueueConsumer(rdb *redis.Client, queue <--chan Request) {
limiter := ratelimiter.NewTokenBucket(rdb, 100, 100)
for req := range queue {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
ok, _, err := limiter.Allow(ctx, "queue:"+req.Topic)
cancel()
if ok {
processRequest(req)
} else {
// 令牌不足,重新入队等待
go func(r Request) {
time.Sleep(100 * time.Millisecond)
queue <- r
}(req)
}
}
}
限流器性能测试与调优
分布式限流的性能瓶颈在Redis网络往返。优化方向:使用Redis Pipeline批量处理多个限流检查;使用本地缓存做预过滤,只有本地令牌不足时才访问Redis;连接池调优,避免连接建立开销。
// Redis连接池配置优化
opt, _ := redis.ParseURL("redis://localhost:6379/0")
opt.PoolSize = 50
opt.MinIdleConns = 10
opt.MaxConnAge = 5 * time.Minute
rdb := redis.NewClient(opt)
// 本地预过滤优化
type LocalPrefilter struct {
localTokens chan struct{}
remote *ratelimiter.TokenBucket
}
func (lf *LocalPrefilter) Allow(ctx context.Context, key string) (bool, error) {
select {
case <-lf.localTokens:
return true, nil
default:
ok, _, err := lf.remote.Allow(ctx, key)
return ok, err
}
}
服务治理中限流是核心环节。分布式令牌桶方案的可靠性取决于Redis的可用性,建议Redis配置主从+哨兵或Cluster模式。限流阈值设置需基于压测数据,观察系统在何等QPS下开始出现响应时间劣化,以此为依据设定rate和capacity。上线后持续监控限流触发率,过高说明请求量超出系统承载能力,需要考虑扩容或引流降级。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-shi-xian-fen-bu-shi-ling-pai-tong-xian-liu/