Appearance
服务治理:限流、熔断与负载均衡
本篇讲解 go-zero 的服务治理能力,包括限流、熔断、负载均衡和服务发现。这些能力在微服务架构中至关重要:限流保护服务不被流量打垮,熔断防止故障扩散,负载均衡提高资源利用率,服务发现实现动态扩缩容。go-zero 把这些能力内置到框架中,开发者只需配置即可启用。
一、go-zero 内置限流
限流(Rate Limiting)是保护服务的第一道防线,用于在流量超过承受能力时主动拒绝部分请求,避免服务崩溃。
1. 基于 Redis 的分布式限流
go-zero 提供 limit.PeriodLimit(周期限流)和 limit.TokenLimit(令牌桶)两种限流器,都基于 Redis 实现,适合分布式环境。
周期限流(PeriodLimit)
在固定时间窗口内允许 N 次请求:
go
import "github.com/zeromicro/go-zero/core/limit"
// 每秒允许 100 次请求
limiter := limit.NewPeriodLimit(1, 100, rds, "rate_limit:user")
code, err := limiter.Take("user:123")
// code 取值:
// limit.AllowOk = 0 放行
// limit.AllowHit = 1 放行(命中限流但未超)
// limit.AllowOverQuota = 2 拒绝(超限)
// limit.AllowShared = 3 共享窗口(特殊场景)
if code == limit.OverQuota {
return fmt.Errorf("rate limited")
}实现原理:用 Redis 的 INCR 累计计数,首次访问时设置过期(即窗口开始),窗口内累计超过阈值则拒绝。
令牌桶限流(TokenLimit)
更平滑的限流,按固定速率生成令牌,请求消耗令牌:
go
import "github.com/zeromicro/go-zero/core/limit"
// rate: 每秒生成 100 个令牌
// bucket: 桶容量 200(允许突发 200 请求)
limiter := limit.NewTokenLimit(100, 200, rds)
if limiter.Allow() {
// 放行
} else {
// 拒绝
}令牌桶的优势:
- 允许突发流量(桶内有积累的令牌)
- 平均速率稳定(rate 控制)
- 适合大多数业务场景
在中间件中使用
go
// internal/middleware/ratelimitmiddleware.go
package middleware
import (
"net/http"
"strconv"
"github.com/zeromicro/go-zero/core/limit"
"github.com/zeromicro/go-zero/core/stores/redis"
)
type RateLimitMiddleware struct {
limiter *limit.TokenLimiter
}
func NewRateLimitMiddleware(rds *redis.Redis, rate, bucket int) *RateLimitMiddleware {
return &RateLimitMiddleware{
limiter: limit.NewTokenLimit(rate, bucket, rds),
}
}
func (m *RateLimitMiddleware) Handle(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if !m.limiter.Allow() {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusTooManyRequests)
w.Write([]byte(`{"code":429,"msg":"too many requests"}`))
return
}
next(w, r)
}
}2. 单机限流
如果不需要分布式协调,可以用 syncx 包的单机限流:
go
import "github.com/zeromicro/go-zero/core/syncx"
// 每秒 1000 次
limiter := syncx.NewLimit(1000)
if limiter.TryReserve() {
defer limiter.Release()
// 处理请求
} else {
// 拒绝
}单机限流的优势是性能高(无网络开销),适合:
- 单机限流场景
- 实例数固定且不需要全局协调
3. 配置限流规则
在配置文件中定义限流参数:
go
// internal/config/config.go
type Config struct {
rest.RestConf
RateLimit struct {
Rate int `json:",default=100"` // 每秒令牌数
Bucket int `json:",default=200"` // 桶容量
}
Redis redis.RedisConf
}yaml
# etc/user-api.yaml
RateLimit:
Rate: 100
Bucket: 200
Redis:
Host: 127.0.0.1:6379
Type: node在 ServiceContext 中初始化:
go
func NewServiceContext(c config.Config) *ServiceContext {
rds := redis.MustNewRedis(c.Redis)
return &ServiceContext{
Config: c,
RateLimitMiddleware: middleware.NewRateLimitMiddleware(rds, c.RateLimit.Rate, c.RateLimit.Bucket).Handle,
}
}4. 分维度限流
生产环境通常需要多维度限流:
go
// 接口级限流:每个接口独立限流
func (m *RateLimitMiddleware) HandleWithKey(key string) rest.Middleware {
return func(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
// key = "user_api:127.0.0.1"
fullKey := key + ":" + r.RemoteAddr
if !m.limiter.Allow(fullKey) {
w.WriteHeader(http.StatusTooManyRequests)
return
}
next(w, r)
}
}
}
// 用户级限流:限制单用户调用频率
func (l *GetUserLogic) GetUser(req *types.GetUserRequest) (*types.GetUserResponse, error) {
userId := getUserIdFromCtx(l.ctx)
key := fmt.Sprintf("user:%d:query", userId)
code, _ := l.svcCtx.UserLimiter.Take(key)
if code == limit.OverQuota {
return nil, fmt.Errorf("user query too frequent")
}
// ...
}二、go-zero 内置熔断器
熔断器(Circuit Breaker)是防止故障扩散的关键机制。当某个下游服务故障率过高时,熔断器会「断开」,后续请求直接失败或走降级,避免连锁雪崩。
1. googlebreaker 算法
go-zero 的熔断器基于 Google SRE 的自适应熔断算法(googlebreaker),与 Hystrix 的「固定阈值」不同,它是基于「请求成功率」动态决策的。
核心思想:
- 维护一个滑动窗口的请求统计(成功、失败、总量)
- 计算最近一段时间内的成功率
- 用概率公式决定是否放行:
K = 2 / (windowSize * success_rate)
accept_probability = max(0, (total - K * failures) / (total + 1))即:失败越多,放行概率越低;成功率恢复后,自动放行概率上升。
这种算法的优点:
- 自适应:根据实时成功率自动调整,无需手动设阈值
- 平滑:不是「一刀切」断开,而是按概率放行部分请求做探测
- 自恢复:服务恢复后成功率上升,熔断自动解除
2. 自适应熔断
go-zero 的 RPC 客户端和服务端都内置了熔断器,无需手动配置:
go
// RPC 客户端:调用失败时熔断器记录
client, _ := zrpc.NewClient(c.UserRpc)
// 内部已注册 UnaryBreakerInterceptor,自动对每次调用做熔断判断
// RPC 服务端:服务过载时主动拒绝
// 内部已注册 UnaryBreakerInterceptor,自动保护服务熔断触发后:
- 客户端:调用直接返回 error,不会真正发起 RPC
- 服务端:返回
resource_exhausted错误,建议客户端降级
3. 手动使用熔断器
如果需要在业务逻辑中使用熔断(如调用第三方 API),可以用 breaker 包:
go
import "github.com/zeromicro/go-zero/core/breaker"
// 用 name 隔离不同下游的熔断器
err := breaker.Do("third_party:payment", func() error {
// 调用第三方支付 API
resp, err := http.Post(paymentURL, "application/json", body)
if err != nil {
return err
}
if resp.StatusCode >= 500 {
return fmt.Errorf("payment service error")
}
return nil
})
if err != nil {
// 熔断中或调用失败,走降级
return fallbackResponse()
}breaker.Do 的工作机制:
- 检查熔断器状态,如果「断开」则直接返回 err
- 执行传入的函数
- 根据返回结果更新统计(成功/失败)
- 失败率过高时触发熔断
4. 熔断器配置
go-zero 的熔断器是全局自适应的,一般不需要配置。如果需要调整:
go
import "github.com/zeromicro/go-zero/core/breaker"
// 自定义熔断器名称和参数
breaker.NewBreaker("my_service",
breaker.WithName("my_service"),
breaker.WithReason("..."),
)
// 获取熔断器状态
b := breaker.GetBreaker("my_service")
if b != nil {
// b 内部维护统计数据
}5. 熔断与重试的关系
熔断和重试是矛盾的:重试会增加下游压力,熔断是为了减少下游压力。最佳实践:
- 熔断优先于重试:熔断状态下不重试
- 重试次数有限:最多 1-2 次
- 重试要有退避:避免重试风暴
go
func callWithRetry(ctx context.Context, fn func() error) error {
maxRetries := 2
var lastErr error
for i := 0; i <= maxRetries; i++ {
err := breaker.Do("service_a", fn)
if err == nil {
return nil
}
lastErr = err
// 检查是否熔断
if errors.Is(err, breaker.ErrServiceUnavailable) {
return err // 熔断中,不重试
}
// 退避
time.Sleep(time.Duration(i*100) * time.Millisecond)
}
return lastErr
}三、负载均衡
负载均衡(Load Balancing)是把请求分发到多个后端实例的策略。go-zero 的 RPC 客户端内置了多种负载均衡算法。
1. 内置负载均衡策略
go-zero 默认使用 p2c(Power of Two Choices) 算法。p2c 的核心思想:
- 随机选两个后端实例
- 比较两者的「负载分」(基于连接数、延迟、成功率等综合计算)
- 选择负载较低的那个
为什么不用轮询?
- 轮询不考虑实例的实际负载,可能把请求打到繁忙的实例
- p2c 通过随机 + 比较,能近似做到「最闲优先」,但实现比「最闲优先」简单
- 大规模集群下,p2c 表现接近最优
p2c 的负载分计算(go-zero 实现):
load = inflight(正在处理的请求数) * ewma(指数加权移动平均延迟)即:正在处理的请求越多、平均延迟越高的实例,被选中的概率越低。
2. p2c(Power of Two Choices)
p2c 的实现位于 zrpc/internal/balancer/p2c/。每个后端实例维护:
inflight:当前正在处理的请求数latency:EWMA 延迟(最近请求的加权平均)success:成功率lastPick:上次被选中的时间
每次选实例:
go
func (p *p2cPicker) Pick(info PickerInfo) (Conn, func(), error) {
// 随机选两个不同的实例
a := p.randomPick()
b := p.randomPick()
for a == b {
b = p.randomPick()
}
// 比较负载,选较小的
if p.load(a) < p.load(b) {
return a.conn, a.done, nil
}
return b.conn, b.done, nil
}EWMA(Exponentially Weighted Moving Average)的特点:
- 越近的请求权重越大,能快速反映当前状态
- 不需要保存所有历史数据,内存占用小
- 适合实时负载评估
3. 自定义负载均衡
如果需要其他算法(如轮询、一致性哈希),可以实现 grpc.Picker 和 grpc.Resolver:
go
import (
"google.golang.org/grpc/balancer"
"google.golang.org/grpc/balancer/base"
)
const RoundRobinName = "round_robin_custom"
func init() {
balancer.Register(base.NewBalancerBuilder(RoundRobinName,
&rrPickerBuilder{}, base.Config{HealthCheck: true}))
}
type rrPickerBuilder struct{}
func (b *rrPickerBuilder) Build(info base.PickerBuildInfo) balancer.Picker {
var conns []balancer.SubConn
for conn := range info.ReadySCs {
conns = append(conns, conn)
}
return &rrPicker{conns: conns, idx: 0}
}
type rrPicker struct {
conns []balancer.SubConn
idx int
mu sync.Mutex
}
func (p *rrPicker) Pick(info balancer.PickInfo) (balancer.SubConn, func(balancer.PickDoneInfo), error) {
p.mu.Lock()
defer p.mu.Unlock()
if len(p.conns) == 0 {
return nil, nil, fmt.Errorf("no available conn")
}
conn := p.conns[p.idx%len(p.conns)]
p.idx++
return conn, func(info balancer.PickDoneInfo) {}, nil
}使用自定义负载均衡:
go
client, _ := zrpc.NewClient(c.UserRpc,
zrpc.WithDialOption(grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin_custom"}`)),
)实际生产环境,p2c 几乎满足所有场景,不建议自己实现。
4. 负载均衡与权重
go-zero 的 p2c 支持权重(通过 etcd 的 metadata):
yaml
# 实例 1 配置(高性能机器)
Etcd:
Key: user.rpc
Hosts:
- 10.0.0.1:2379
Weight: 100 # 权重(go-zero 扩展字段)
# 实例 2 配置(低性能机器)
Etcd:
Key: user.rpc
Weight: 50权重高的实例会被选中更多次。
四、服务发现
服务发现(Service Discovery)是微服务架构的基础设施,让服务实例能动态感知彼此。go-zero 默认集成 etcd。
1. etcd 服务注册
服务端配置:
yaml
# etc/user.yaml
Name: user.rpc
ListenOn: 0.0.0.0:8080
Etcd:
Hosts:
- 127.0.0.1:2379
Key: user.rpc启动时,go-zero 会:
- 在 etcd 中创建一个 key:
user.rpc/10.0.0.1:8080(前缀是 Key,后缀是实例地址) - 设置 lease(租约),定期续约
- 实例宕机时 lease 过期,etcd 自动删除该 key
客户端通过订阅 user.rpc 前缀,能感知实例上下线。
2. 客户端服务发现
yaml
# etc/user-api.yaml
UserRpc:
Etcd:
Hosts:
- 127.0.0.1:2379
Key: user.rpc
NonBlock: true
Timeout: 2000客户端启动时:
- 从 etcd 拉取
user.rpc前缀下所有 key(即所有实例地址) - 监听 etcd 变更,实例上下线时自动更新本地列表
- 调用时用 p2c 选一个实例
3. 多实例注册
启动多个实例(同一份代码,不同端口):
bash
# 实例 1
go run user.go -f etc/user1.yaml # ListenOn: 0.0.0.0:8080
# 实例 2
go run user.go -f etc/user2.yaml # ListenOn: 0.0.0.0:8081
# 实例 3
go run user.go -f etc/user3.yaml # ListenOn: 0.0.0.0:8082三个实例都用同一个 Etcd.Key,etcd 会维护一个实例列表。客户端自动负载均衡。
4. 直连模式
本地调试时,可以不用 etcd,直接指定目标地址:
yaml
# 直连单个实例
UserRpc:
Endpoints:
- 127.0.0.1:8080
NonBlock: true或者客户端代码:
go
client, _ := zrpc.NewClientWithTarget("127.0.0.1:8080")直连模式不支持多实例负载均衡(除非用 DNS 解析多个 IP)。
5. 服务健康检查
go-zero 的健康检查基于 etcd 的 lease 机制:
- 服务启动时创建 lease(默认 10 秒)
- 服务定期(默认 5 秒)续约 lease
- 服务宕机时 lease 过期,etcd 删除对应 key
- 客户端监听到 key 删除,更新本地实例列表
这种机制是「被动健康检查」:不主动探测实例健康,而是依赖 lease 过期。优点是无需额外网络开销,缺点是故障感知有延迟(最长 lease 时间)。
如果需要主动健康检查,可以配置 gRPC 的 health check:
go
import (
"google.golang.org/grpc/health"
healthpb "google.golang.org/grpc/health/grpc_health_v1"
)
s := zrpc.MustNewServer(c.RpcServerConf, func(grpcServer *grpc.Server) {
pb.RegisterUserServer(grpcServer, srv)
healthpb.RegisterHealthServer(grpcServer, health.NewServer())
})客户端会自动对实例做 health check,剔除不健康的实例。
6. 其他服务发现方案
go-zero 也支持其他注册中心(需自行扩展 Resolver):
- Consul
- Nacos
- Zookeeper
- Kubernetes Service(用 DNS 解析)
K8s 环境下,可以用 Service + ClusterIP,无需 etcd:
yaml
UserRpc:
Endpoints:
- user-rpc-service.default.svc.cluster.local:8080
NonBlock: true五、完整示例:带限流和熔断的微服务
下面演示一个完整的「用户 API + 用户 RPC」服务,包含限流、熔断、负载均衡。
1. 架构
Client
↓
[user-api (HTTP)] ── 限流中间件 ── JWT 中间件 ── Logic
↓
[user-rpc (gRPC)] ── 熔断器 ── Logic ── DB
(3 个实例,p2c 负载均衡)2. user-rpc 多实例配置
etc/user1.yaml:
yaml
Name: user.rpc
ListenOn: 0.0.0.0:8080
Etcd:
Hosts:
- 127.0.0.1:2379
Key: user.rpc
Prometheus:
Host: 0.0.0.0
Port: 9101
Path: /metricsetc/user2.yaml:
yaml
Name: user.rpc
ListenOn: 0.0.0.0:8081
Etcd:
Hosts:
- 127.0.0.1:2379
Key: user.rpc
Prometheus:
Host: 0.0.0.0
Port: 9102
Path: /metrics启动两个实例:
bash
go run user.go -f etc/user1.yaml &
go run user.go -f etc/user2.yaml &3. user-api 限流配置
etc/user-api.yaml:
yaml
Name: user-api
Host: 0.0.0.0
Port: 8888
Redis:
Host: 127.0.0.1:6379
Type: node
RateLimit:
Rate: 200 # 每秒 200 令牌
Bucket: 400 # 桶容量 400
UserRpc:
Etcd:
Hosts:
- 127.0.0.1:2379
Key: user.rpc
NonBlock: true
Timeout: 2000
Prometheus:
Host: 0.0.0.0
Port: 9100
Path: /metrics4. 限流中间件
go
// internal/middleware/ratelimitmiddleware.go
package middleware
import (
"encoding/json"
"net/http"
"github.com/zeromicro/go-zero/core/limit"
"github.com/zeromicro/go-zero/core/stores/redis"
)
type RateLimitMiddleware struct {
limiter *limit.TokenLimiter
}
func NewRateLimitMiddleware(rds *redis.Redis, rate, bucket int) *RateLimitMiddleware {
return &RateLimitMiddleware{
limiter: limit.NewTokenLimit(rate, bucket, rds),
}
}
func (m *RateLimitMiddleware) Handle(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
if !m.limiter.Allow() {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusTooManyRequests)
body, _ := json.Marshal(map[string]any{
"code": 429,
"msg": "too many requests, please retry later",
})
w.Write(body)
return
}
next(w, r)
}
}5. API Logic 调用 RPC(含降级)
go
// internal/logic/user/getuserlogic.go
func (l *GetUserLogic) GetUser(req *types.GetUserRequest) (resp *types.GetUserResponse, err error) {
id, _ := strconv.ParseInt(req.Id, 10, 64)
// 调用 RPC,go-zero 自动应用熔断
rpcResp, err := l.svcCtx.UserRpc.GetUser(l.ctx, &pb.GetUserRequest{Id: id})
if err != nil {
// 熔断或调用失败,走降级
l.Errorf("call user rpc failed: %v, fallback to cache", err)
return l.fallbackFromCache(id)
}
return &types.GetUserResponse{
Id: rpcResp.User.Id,
Name: rpcResp.User.Name,
Email: rpcResp.User.Email,
}, nil
}
// 降级:从 Redis 缓存读取
func (l *GetUserLogic) fallbackFromCache(id int64) (*types.GetUserResponse, error) {
key := fmt.Sprintf("user:cache:%d", id)
val, err := l.svcCtx.Redis.GetCtx(l.ctx, key)
if err != nil || val == "" {
return nil, fmt.Errorf("service unavailable")
}
var resp types.GetUserResponse
if err := json.Unmarshal([]byte(val), &resp); err != nil {
return nil, err
}
return &resp, nil
}6. 验证熔断效果
模拟 RPC 故障(在 user-rpc Logic 里返回 error):
go
// user-rpc/internal/logic/getuserlogic.go
func (l *GetUserLogic) GetUser(in *pb.GetUserRequest) (*pb.GetUserResponse, error) {
return nil, status.Error(codes.Internal, "simulated failure")
}观察现象:
- 前几次调用:API 收到 RPC 错误,走降级
- 失败率上升后:熔断器触发,API 不再真正调用 RPC,直接走降级
- 恢复 RPC(去掉 error):熔断器探测到成功率上升,自动恢复调用
可以在 Prometheus 监控中看到熔断器状态变化(go_zero_zrpc_client_breaker 指标)。
7. 验证负载均衡
在 user-rpc Logic 里打印实例信息:
go
func (l *GetUserLogic) GetUser(in *pb.GetUserRequest) (*pb.GetUserResponse, error) {
l.Infof("instance handling request, id=%d", in.Id)
return &pb.GetUserResponse{...}, nil
}发起多次请求,观察日志,会发现请求被均匀分发到两个实例。
六、服务治理最佳实践
1. 限流策略选择
| 场景 | 推荐策略 |
|---|---|
| 公开 API(登录、注册) | 全局限流 + IP 限流 |
| 用户接口(查询、操作) | 用户级限流 |
| 核心 API(下单、支付) | 接口级限流 + 用户级限流 |
| 后台任务 | 全局限流,避免影响在线业务 |
2. 熔断配置建议
- 默认配置即可(go-zero 自适应)
- 区分下游:每个下游服务用独立熔断器(breaker 用 name 隔离)
- 配合降级:熔断后必须有降级方案,否则用户体验差
3. 负载均衡建议
- 默认 p2c 即可
- 实例性能差异大时配置 weight
- 避免单实例过载:用 K8s HPA 自动扩缩容
4. 服务发现建议
- 生产环境用 etcd(go-zero 默认支持完善)
- K8s 环境可用 Service + DNS,简化部署
- 跨机房部署时,注意 etcd 集群的延迟
5. 监控告警
监控以下指标:
- 限流拒绝率(
http_requests_total中 429 状态码占比) - 熔断器状态(
go_zero_zrpc_client_breaker) - 负载均衡分布(各实例 QPS 是否均衡)
- 服务实例数(etcd 中实例数量)
七、小结
本篇讲解了 go-zero 的服务治理能力。要点回顾:
- 限流:
PeriodLimit(周期限流)和TokenLimit(令牌桶),基于 Redis 实现分布式限流 - 熔断:基于 Google SRE 的自适应算法,无需手动设阈值,自动恢复
- 负载均衡:默认 p2c 算法,基于 EWMA 评估实例负载,比轮询更智能
- 服务发现:默认 etcd,基于 lease 实现被动健康检查
- 完整示例演示了限流 + 熔断 + 降级 + 多实例负载均衡的端到端流程
- 服务治理的核心目标:保护服务、防止故障扩散、提高可用性
下一篇我们将进入配置管理与多环境部署,讲解 go-zero 的配置体系、多环境隔离、配置中心集成和容器化部署。