微服务限流的必要性
高并发微服务架构中,限流是保障系统稳定性的核心手段。当下游服务响应变慢或遭遇突发流量时,上游服务的请求会堆积,引发级联故障。限流通过主动丢弃超额请求来保护系统,避免被拖垮。Go语言因其轻量级协程模型,在微服务限流实现上有天然优势。
令牌桶算法的Go实现
令牌桶是应用最广泛的限流算法,核心思路是以固定速率向桶中放入令牌,请求到达时取走令牌,桶空则拒绝。它允许瞬时突发流量(桶中有积累的令牌),同时保证长期平均速率不超过设定值。
package ratelimit
import (
"sync"
"time"
)
type TokenBucket struct {
mu sync.Mutex
rate float64 // 每秒放入令牌数
capacity float64 // 桶容量
tokens float64 // 当前令牌数
lastTime time.Time // 上次填充时间
}
func NewTokenBucket(rate, capacity float64) *TokenBucket {
return &TokenBucket{
rate: rate,
capacity: capacity,
tokens: capacity,
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 > tb.capacity {
tb.tokens = tb.capacity
}
if tb.tokens >= 1 {
tb.tokens -= 1
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(50 * time.Millisecond):
continue
}
}
}
分布式限流:Redis加Lua的原子实现
单机限流无法应对多实例部署。Redis加Lua脚本实现分布式令牌桶,保证原子性操作:
-- rate_limiter.lua
local key = KEYS[1]
local rate = tonumber(ARGV[1])
local capacity = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local data = redis.call("HMGET", key, "tokens", "last_time")
local tokens = tonumber(data[1]) or capacity
local last_time = tonumber(data[2]) or now
local elapsed = now - last_time
tokens = math.min(capacity, tokens + elapsed * rate)
local allowed = 0
if tokens >= 1 then
tokens = tokens - 1
allowed = 1
end
redis.call("HMSET", key, "tokens", tokens, "last_time", now)
redis.call("EXPIRE", key, math.ceil(capacity / rate) + 10)
return allowed
Go端调用:
func (r *RedisRateLimiter) Allow(ctx context.Context, key string) (bool, error) {
now := time.Now().UnixMilli() / 1000.0
result, err := r.luaScript.Run(
ctx, r.client,
[]string{fmt.Sprintf("ratelimit:%s", key)},
r.rate, r.capacity, now,
).Int()
if err != nil {
return false, err
}
return result == 1, nil
}
自适应限流:根据系统负载动态调整
固定阈值限流的问题是:阈值设高则保护不了系统,设低则浪费容量。自适应限流根据实时系统指标动态调整限流阈值。Netflix的梯度限流和阿里Sentinel的自适应策略是两种成熟方案。
以下实现基于CPU利用率和平均响应时间的自适应限流:
package ratelimit
import (
"sync/atomic"
"time"
)
type AdaptiveLimiter struct {
maxQPS int64
currentLimit atomic.Int64
// 采样窗口
totalRequests atomic.Int64
totalLatency atomic.Int64 // 微秒
lastSample time.Time
}
func (al *AdaptiveLimiter) Start() {
al.lastSample = time.Now()
go al.adjustLoop()
}
func (al *AdaptiveLimiter) adjustLoop() {
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
for range ticker.C {
al.adjust()
}
}
func (al *AdaptiveLimiter) adjust() {
now := time.Now()
elapsed := now.Sub(al.lastSample).Seconds()
al.lastSample = now
reqCount := al.totalRequests.Swap(0)
latSum := al.totalLatency.Swap(0)
if reqCount == 0 {
return
}
// 计算QPS和平均延迟
qps := float64(reqCount) / elapsed
avgLatency := float64(latSum) / float64(reqCount) / 1000.0
// 获取CPU使用率
cpuUsage := getCPUUsage()
currentLimit := al.currentLimit.Load()
// 调整策略
switch {
case cpuUsage > 0.8 || avgLatency > 500:
// 负载过高,降低限流阈值
newLimit := int64(float64(currentLimit) * 0.7)
if newLimit < 10 {
newLimit = 10
}
al.currentLimit.Store(newLimit)
case cpuUsage < 0.5 && avgLatency < 200:
// 负载健康,逐步提升限流阈值
newLimit := int64(float64(currentLimit) * 1.1)
if newLimit > al.maxQPS {
newLimit = al.maxQPS
}
al.currentLimit.Store(newLimit)
}
}
限流在微服务网关层的集成
限流的最佳部署位置是API网关层,而非每个微服务内部。在网关统一限流可以减少服务内部的重复实现:
// Gin中间件集成
func RateLimitMiddleware(limiter RateLimiter) gin.HandlerFunc {
return func(c *gin.Context) {
key := c.ClientIP()
if c.GetHeader("X-API-Key") != "" {
key = c.GetHeader("X-API-Key") // 按API Key限流
}
if !limiter.Allow() {
c.JSON(http.StatusTooManyRequests, gin.H{
"error": "rate limit exceeded",
"retry_after": 1,
})
c.Abort()
return
}
c.Next()
}
}
网关层限流支持按IP、按用户、按API Key多维度配置。后端微服务内部可以再做一层细粒度限流(如按租户、按资源类型),形成网关粗粒度加服务细粒度的双层限流架构。
限流指标的可观测性
限流效果需要量化监控。Prometheus指标暴露是标准做法:
var (
limitTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "rate_limit_total",
Help: "Rate limit check total",
},
[]string{"result"}, // allowed, rejected
)
limitCurrent = promauto.NewGaugeVec(
prometheus.GaugeOpts{
Name: "rate_limit_current_qps",
Help: "Current QPS",
},
[]string{"service"},
)
)
通过rate(limit_total{result="rejected"}[5m])计算拒绝速率,结合服务错误率和延迟P99,评估限流阈值是否合理。限流不是一劳永逸的配置,需要根据业务增长持续调优。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-wei-fu-wu-xian-liu-shi-xian-ling-pai-tong-yu-zi/