Go语言实现分布式令牌桶限流:Redis Lua原子操作保证高并发安全

限流是高并发系统中保护服务稳定性的第一道防线。令牌桶算法相比固定窗口和滑动窗口,支持突发流量且实现简洁,是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/

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

相关推荐