0%

go-zero 源码分析 08:zRPC 调用链

我们完整走完了 REST 服务的全部链路——从 goctl 生成代码、配置加载与服务启动,到路由匹配、参数解析、响应写入,再到 11 个中间件的协同配合。现在把视角转到微服务架构的另一端:RPC。

zRPC 增加了什么

gRPC 本身已经是一个功能完备的框架——protobuf 定义了接口契约,protoc-gen-go-grpc 生成了 stub 代码,grpc.Servergrpc.ClientConn 封装了网络传输、序列化和流控。那么问题来了:go-zero 在标准 gRPC 之上到底增加了什么? 如果只是 MustNewServerMustNewClient 两个工厂函数,那跟直接调用 gRPC 原生 API 没有本质区别。

答案藏在一个关键词里:拦截器(Interceptor)。gRPC 的拦截器机制相当于 REST 框架的中间件链——它们都能在请求处理的前后插入横切逻辑。但 gRPC 原生只提供了拦截器的注册入口,没有提供开箱即用的生产级拦截器组合。go-zero 的 zRPC 模块补上的正是这一层:

  • 服务端:自动装配 Trace、Recover、Stat、Prometheus、Breaker、Shedding、Timeout 和 Auth 共 8 个拦截器,它们的顺序设计和职责分工与 REST 中间件一脉相承,但又因为 gRPC 的协议特性而有了不同的侧重点。
  • 客户端:自动装配 Trace、Duration、Prometheus、Breaker 和 Timeout 共 5 个拦截器,保证每次 RPC 调用都有完整的遥测、保护和超时控制。
  • 基础设施:配置驱动的服务注册发现集成、健康检查、代理模式等,让 gRPC 服务从"能跑"升级为"能在生产环境中可靠运行"。

接下来我们先用 goctl 生成并运行一个带业务服务分组的 RPC 工程,让后文中的 server、logic、注册函数和全限定方法名都有实际落点;然后以 RpcServerRpcClient 的创建为入口,分别展开服务端和客户端的拦截器链。在文章最后,我们会用一次成功调用和一次超时调用,串联起客户端和服务端的完整拦截器时序。

先跑起来:goctl 生成带业务分组的 RPC 服务

这个示例只有一个进程和一个监听端口,但业务上分成用户与订单两组,每组各有两个方法。

先说清楚:这里的“服务分组”是什么

本文采用 go-zero Proto DSL 服务分组 中的定义:一个 protobuf service 块就是一个业务分组,块内的 rpc 声明是该分组下的方法。两个业务分组、每组两个方法落实到 proto 语法中,就是两个 service、每个 service 两个 rpc

业务分组(proto service 组内 RPC 方法
UserService CreateUserGetUser
OrderService CreateOrderGetOrder

这里的 业务服务分组 之前讲的 core/service.ServiceGroup 不是同一个概念:

  • proto 服务分组解决的是生成代码如何按业务边界拆目录
  • core/service.ServiceGroup 解决的是一个进程如何编排多个可启动组件的生命周期

本例最终只创建一个 zrpc.RpcServer,两个业务分组都注册到它内部的同一个 grpc.Server,并不会启动两个 RPC server。

在 go-zero 中,RPC 的 业务分组 是通过在 proto 文件中以 service 为维度来进行文件分组,而 Rest 的 业务分组 是通过在 @server 注解中以 group 关键字来进行服务分组

完整的 shop.proto 如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
syntax = "proto3";

package shop;

option go_package = "github.com/zeromicro/go-zero/ai_demo/rpc_service_group/pb";

message CreateUserRequest {
string name = 1;
}

message UserResponse {
int64 id = 1;
string name = 2;
}

message GetUserRequest {
int64 id = 1;
}

service UserService {
rpc CreateUser(CreateUserRequest) returns (UserResponse);
rpc GetUser(GetUserRequest) returns (UserResponse);
}

message CreateOrderRequest {
int64 user_id = 1;
string product = 2;
}

message OrderResponse {
int64 id = 1;
int64 user_id = 2;
string product = 3;
}

message GetOrderRequest {
int64 id = 1;
}

service OrderService {
rpc CreateOrder(CreateOrderRequest) returns (OrderResponse);
rpc GetOrder(GetOrderRequest) returns (OrderResponse);
}

-m 生成分组目录

ai_demo/rpc_service_group 目录执行:

1
2
3
4
5
6
7
8
goctl rpc protoc shop.proto \
--go_out=. \
--go-grpc_out=. \
--zrpc_out=. \
--go_opt=module=github.com/zeromicro/go-zero/ai_demo/rpc_service_group \
--go-grpc_opt=module=github.com/zeromicro/go-zero/ai_demo/rpc_service_group \
--module=github.com/zeromicro/go-zero/ai_demo/rpc_service_group \
-m

关键参数是 -m,也就是 --multiple。它告诉 goctl 当前 proto 包含多个 service,需要启用多服务模式;三个 module 参数让 protobuf 文件落到当前示例的 pb 目录,并生成正确的 Go import path。

生成后的核心目录如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
ai_demo/rpc_service_group/
├── shop.proto
├── shop.go
├── etc/
│ └── shop.yaml
├── pb/
│ ├── shop.pb.go
│ └── shop_grpc.pb.go
├── client/
│ ├── userservice/
│ │ └── userservice.go
│ └── orderservice/
│ └── orderservice.go
└── internal/
├── logic/
│ ├── userservice/
│ │ ├── createuserlogic.go
│ │ └── getuserlogic.go
│ └── orderservice/
│ ├── createorderlogic.go
│ └── getorderlogic.go
├── server/
│ ├── userservice/
│ │ └── userserviceserver.go
│ └── orderservice/
│ └── orderserviceserver.go
└── svc/
└── servicecontext.go

分组体现在三处:clientinternal/logicinternal/server 都以 proto service 名拆出子目录。与此同时,两组仍然共享 shop.go 入口、shop.yaml 配置、protobuf 消息和 ServiceContext这正是 按业务隔离代码、按进程共享基础设施 的效果

goctl 生成的 logic 默认返回空响应。为了让例子可以实际调用,本文在生成代码上补了一个并发安全的内存存储:CreateUser 生成用户 ID,CreateOrder 校验用户并生成订单 ID,两个 Get 方法按 ID 查询。重启进程后数据清空,不引入数据库,以免干扰本文对 zRPC 调用链的观察。

一个端口怎样承载两个业务分组

打开 goctl 生成的 shop.go,可以看到两个分组的 server 被注册到同一个原生 grpc.Server

1
2
3
4
5
6
7
8
9
10
11
12
s := zrpc.MustNewServer(c.RpcServerConf, func(grpcServer *grpc.Server) {
pb.RegisterUserServiceServer(
grpcServer, userserviceServer.NewUserServiceServer(ctx))
pb.RegisterOrderServiceServer(
grpcServer, orderserviceServer.NewOrderServiceServer(ctx))

if c.Mode == service.DevMode || c.Mode == service.TestMode {
reflection.Register(grpcServer)
}
})
defer s.Stop()
s.Start()
  • 这里的注册函数来自 protoc-gen-go-grpc,两个 server 实现来自 goctl
  • zrpc.MustNewServer 只创建一次,因此配置里也只有一个 ListenOn: 0.0.0.0:20091
  • gRPC 根据全限定方法名 /shop.UserService/CreateUser/shop.OrderService/CreateOrder,把请求分发到对应的 service handler

在仓库根目录启动服务:

1
2
go run ./ai_demo/rpc_service_group \
-f ai_demo/rpc_service_group/etc/shop.yaml

开发模式会开启 reflection。下面依次调用两个分组中的四个方法:

1
2
3
4
5
6
7
8
9
10
11
grpcurl -plaintext -d '{"name":"Alice"}' \
127.0.0.1:20091 shop.UserService/CreateUser

grpcurl -plaintext -d '{"id":"1"}' \
127.0.0.1:20091 shop.UserService/GetUser

grpcurl -plaintext -d '{"userId":"1","product":"go-zero book"}' \
127.0.0.1:20091 shop.OrderService/CreateOrder

grpcurl -plaintext -d '{"id":"1"}' \
127.0.0.1:20091 shop.OrderService/GetOrder

两组的 Create/Get 调用分别返回同一个用户和订单:

1
2
3
4
{"id":"1", "name":"Alice"}
{"id":"1", "name":"Alice"}
{"id":"1", "userId":"1", "product":"go-zero book"}
{"id":"1", "userId":"1", "product":"go-zero book"}

以创建订单为例,接下来要分析的真实链路已经出现了:

1
2
3
4
5
6
grpcurl
→ /shop.OrderService/CreateOrder
→ zRPC 服务端拦截器链
→ OrderServiceServer.CreateOrder
→ CreateOrderLogic.CreateOrder
→ ServiceContext

有了这个可运行工程作坐标,下面从配置开始拆解:zrpc.MustNewServer 如何创建原生 gRPC server、注册两个业务分组,并在请求到达 logic 之前装配完整的生产级拦截器链。

服务端与客户端各需要什么

在进入拦截器链之前,先理解 zRPC 的配置结构。服务端和客户端的配置需求不同,但它们共享一个核心模式:配置即契约——配置文件中的每个字段都直接对应一个拦截器或一个基础设施行为的开关。

服务端配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// zrpc/config.go
type RpcServerConf struct {
service.ServiceConf // 内嵌通用服务配置(Name、Log、Mode 等)
ListenOn string // 监听地址,如 "0.0.0.0:8080"
Etcd discov.EtcdConf // etcd 服务发现配置(可选)
Auth bool // 是否开启认证
Redis redis.RedisKeyConf // 认证所需的 Redis 配置
StrictControl bool // 认证的严格模式:Redis 不可用是否拒绝请求
Timeout int64 // 默认超时时间(毫秒),0 表示不限制
CpuThreshold int64 // 降载 CPU 阈值(0-1000,默认 900 即 90%)
Health bool // 是否启用 gRPC 健康检查(默认开启)
Middlewares ServerMiddlewaresConf // 拦截器开关
MethodTimeouts []MethodTimeoutConf // 按方法粒度的超时时间
}

这里有两个嵌套结构值得注意。service.ServiceConf 是第 05 篇中详细讲解过的通用服务配置——它提供了 NameMode(开发/测试/生产)、Log 等基础字段,以及关键的 SetUp() 方法(初始化日志、Prometheus、链路追踪等基础设施)。zRPC 服务端嵌入了它,意味着所有 REST 部分享受到的基础设施能力,zRPC 服务器也天然具备。

ServerMiddlewaresConf 定义了拦截器的开关:

1
2
3
4
5
6
7
8
9
// zrpc/internal/config.go
type ServerMiddlewaresConf struct {
Trace bool `json:",default=true"`
Recover bool `json:",default=true"`
Stat bool `json:",default=true"`
StatConf StatConf
Prometheus bool `json:",default=true"`
Breaker bool `json:",default=true"`
}

五个拦截器全部默认开启。注意服务端没有 Timeout 和 Shedding 的开关——它们由 RpcServerConf.TimeoutRpcServerConf.CpuThreshold 的值是否大于零来隐式控制。Timeout=0 表示不限制超时,CpuThreshold=0 表示不启用降载。这种"值即开关"的设计比显式的布尔字段更直接——配置了阈值就启用,没配置就不启用,不存在"配置了阈值却忘记开开关"的错误。

客户端配置

1
2
3
4
5
6
7
8
9
10
11
12
13
// zrpc/config.go
type RpcClientConf struct {
Etcd discov.EtcdConf // etcd 服务发现配置(三选一)
Endpoints []string // 直连地址列表
Target string // 自定义 target 字符串
App string // 认证:应用标识
Token string // 认证:令牌
NonBlock bool // 是否非阻塞连接(默认 true)
Timeout int64 // 默认调用超时(毫秒,默认 2000)
KeepaliveTime time.Duration // 连接保活间隔
Middlewares ClientMiddlewaresConf // 拦截器开关
BalancerName string // 负载均衡策略(默认 "p2c_ewma")
}

客户端配置中最关键的字段是目标寻址的三选一Endpoints(直连地址列表)、Etcd(服务发现)、Target(自定义 gRPC target 字符串)。BuildTarget() 方法的解析顺序也是这个优先级——先看直连,再看 Target,最后 fallback 到 etcd:

1
2
3
4
5
6
7
8
9
func (cc RpcClientConf) BuildTarget() (string, error) {
if len(cc.Endpoints) > 0 {
return resolver.BuildDirectTarget(cc.Endpoints), nil
} else if len(cc.Target) > 0 {
return cc.Target, nil
}
// ... 用 Etcd 配置构建 discov target
return resolver.BuildDiscovTarget(cc.Etcd.Hosts, cc.Etcd.Key), nil
}

客户端拦截器开关同样有五个:

1
2
3
4
5
6
7
type ClientMiddlewaresConf struct {
Trace bool `json:",default=true"`
Duration bool `json:",default=true"`
Prometheus bool `json:",default=true"`
Breaker bool `json:",default=true"`
Timeout bool `json:",default=true"`
}

注意客户端没有 Recover——客户端不执行服务端的业务逻辑,没有"服务端 handler panic"这样的场景需要 recover。取而代之的是 Duration,它类似于 REST 的 LogHandler 和 MetricHandler 的合体——记录每次调用的耗时、解析慢调用、串联请求内容和错误日志。

配置就绪后,我们分别进入服务端和客户端的创建流程。

服务端:从 NewServer 到 gRPC Server 启动

Server 的类型选择:RpcServer vs RpcPubServer

NewServer 的核心逻辑只有三步:校验配置、创建 server 实例、装配拦截器。但"创建 server 实例"这一步有一个关键分支——是否配置了 etcd:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// zrpc/server.go
func NewServer(c RpcServerConf, register internal.RegisterFn) (*RpcServer, error) {
if err = c.Validate(); err != nil {
return nil, err
}

var server internal.Server
metrics := stat.NewMetrics(c.ListenOn)
serverOptions := []internal.ServerOption{
internal.WithRpcHealth(c.Health),
}

if c.HasEtcd() {
server, err = internal.NewRpcPubServer(c.Etcd, c.ListenOn, serverOptions...)
} else {
server = internal.NewRpcServer(c.ListenOn, serverOptions...)
}
// ...
}

这个 Server 接口是对 gRPC 原生 grpc.Server 的抽象:

1
2
3
4
5
6
7
8
// zrpc/internal/server.go
type Server interface {
AddOptions(options ...grpc.ServerOption)
AddStreamInterceptors(interceptors ...grpc.StreamServerInterceptor)
AddUnaryInterceptors(interceptors ...grpc.UnaryServerInterceptor)
SetName(string)
Start(register RegisterFn) error
}

两种实现——rpcServer 和通过 keepAliveServer 装饰的 rpcServer——共享完全相同的拦截器装配和启动流程,区别仅在于 keepAliveServerStart 之前多了一步 etcd 注册

1
2
3
4
5
6
7
// zrpc/internal/rpcpubserver.go
func (s keepAliveServer) Start(fn RegisterFn) error {
if err := s.registerEtcd(); err != nil {
return err
}
return s.Server.Start(fn)
}

registerEtcd 内部创建了一个 discov.Publisher,它会将自己的地址注册到 etcd 并启动 KeepAlive 协程。这里有一个细节:注册到 etcd 的地址不是配置文件中的 listenOn 原文,而是经过了 figureOutListenOn 的转换——如果绑定地址是 0.0.0.0:8080,框架会先尝试 POD_IP 环境变量(Kubernetes 场景),再尝试获取内网 IP,最终替换为真实可路由的 IP。如果不做这个转换,其他服务从 etcd 拿到的地址是 0.0.0.0:8080,根本无法建立连接。

拦截器的装配顺序

回到 NewServer,server 实例创建完毕后,依次装配三类拦截器:

1
2
3
4
5
6
7
// zrpc/server.go
server.SetName(c.Name)
metrics.SetName(c.Name)

setupStreamInterceptors(server, c) // ① Stream 拦截器
setupUnaryInterceptors(server, c, metrics) // ② Unary 拦截器
setupAuthInterceptors(server, c) // ③ 认证拦截器

为什么是这个顺序?因为 AddStreamInterceptorsAddUnaryInterceptors 都是追加到已有的拦截器列表末尾。先装配基础拦截器,再追加认证拦截器,最终形成的顺序是:

1
2
【Unary 拦截器链(服务端)】
Tracing → Recover → Stat → Prometheus → Breaker → Shedding → Timeout → Auth → Handler
1
2
【Stream 拦截器链(服务端)】
Tracing → Recover → Breaker → Auth → Handler

对比 REST 中间件链(Trace → Log → Prometheus → MaxConns → Breaker → Shedding → Timeout → Recover → Metrics → MaxBytes → Gunzip),zRPC 的拦截器顺序有几个值得注意的差异:

Recover 从倒数第二移到了第二。 REST 中 Recover 在 Timeout 的内侧,因为 Timeout 启动了新 goroutine,Recover 必须在这个新 goroutine 里才能捕获到 handler panic。而 zRPC 的 TimeoutInterceptor 自己已经内置了 panic 保护(通过 panicChan 机制),所以 Recover 可以放到更外层——但这并不意味着 Recover 不重要,因为 Timeout 可能被关闭(Timeout=0),此时 Recover 就是唯一的 panic 防线。

Stat 同时承担了日志和指标的双重角色。 REST 中 Log 和 Metrics 是两个独立的中间件,而 zRPC 用 Stat 一个拦截器完成了两件事——收集延迟指标到 stat.Metrics,同时打日志并标记慢调用。这种合并是合理的:gRPC 的调用比 HTTP 更"结构化",不需要区分"请求日志"和"请求指标"两条管线。

没有 MaxConns。

没有 MaxBytes 和 Gunzip。 gRPC 使用 protobuf 序列化,天然的 schema 约束替代了 body 大小限制的手工校验;而 gRPC 本身不基于 HTTP Content-Encoding 做压缩(gRPC 在传输层有自己的压缩机制),所以不需要 Gunzip。

Start:监听、注册、健康检查、优雅停止

拦截器装配完毕后,RpcServer.Start() 调用 server.Start(register)。这个方法的实现简洁但有层次:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
// zrpc/internal/rpcserver.go
func (s *rpcServer) Start(register RegisterFn) error {
lis, err := net.Listen("tcp", s.address)
// ...

unaryInterceptorOption := grpc.ChainUnaryInterceptor(s.unaryInterceptors...)
streamInterceptorOption := grpc.ChainStreamInterceptor(s.streamInterceptors...)
options := append(s.options, unaryInterceptorOption, streamInterceptorOption)
server := grpc.NewServer(options...)
register(server)

// health check
if s.health != nil {
grpc_health_v1.RegisterHealthServer(server, s.health)
s.health.Resume()
}
s.healthManager.MarkReady()
health.AddProbe(s.healthManager)

// graceful stop
waitForCalled := proc.AddShutdownListener(func() {
if s.health != nil {
s.health.Shutdown()
}
server.GracefulStop()
})
defer waitForCalled()

return server.Serve(lis)
}

逐层来看这个方法做了几件事:

第一,grpc.ChainUnaryInterceptorgrpc.ChainStreamInterceptor 将拦截器列表融合为两个 grpc.ServerOption 注意这不是简单的逐个链式调用——gRPC 的 ChainUnaryInterceptor 返回一个"将多个拦截器串联为一个"的选项,框架免去了手工拼接的麻烦。

第二,register(server) 是 goctl 生成的代码中提供的注册函数。 以 greet 服务为例,生成的 register 函数长这样:

1
2
3
func RegisterGreetServer(s *grpc.Server, svr GreetServer) {
s.RegisterService(&Greet_ServiceDesc, svr)
}

这遵循了 gRPC 标准模式——protobuf 生成的 RegisterXxxServer 函数将服务实现注册到 grpc.Server 中。go-zero 没有在这一层做魔改,完全保持与 gRPC 生态的兼容。

第三,健康检查分两套。 gRPC 标准的 grpc_health_v1.HealthServer 暴露的是 gRPC 协议的健康检查服务(供 gRPC 客户端调用 Check 方法)。而 healthManager 是 go-zero 内部的探针接口(第 05 篇中介绍的 Probe 接口),供 health.AddProbe 注册为 Kubernetes 的 liveness/readiness 探针端点。两者互补——一个对外,一个对内。

第四,优雅停止通过 proc.AddShutdownListener 注册。 先关闭健康检查(标记服务不可用,从负载均衡中摘除),再调用 server.GracefulStop()(等待正在处理的请求完成后再关闭连接)。这个流程与 REST 的优雅停止设计完全一致。

嵌套启动:gRPC 日志接管

还有一个难以察觉但影响日常运维的细节——gRPC 的日志接管。在 internal/rpclogger.goinit() 函数中:

1
2
3
func init() {
grpclog.SetLoggerV2(new(Logger))
}

go-zero 用自定义的 Logger 替换了 gRPC 内置的日志输出器。gRPC 原生日志会打印大量 transport 层的 debug 和 info 日志,在生产环境中是噪音。go-zero 的做法是:Error 和 Fatal 级别的日志转发到 logx.Error(不丢失关键错误信息),Info 和 Warning 级别的日志直接丢弃。这意味着你在生产日志中看到的每一条 gRPC 相关记录,都是框架有意为之的行为日志,而不是传输层的噪音。

服务端拦截器逐个剖析

TracingInterceptor:在 Span 中串联跨服务调用链

服务端的 Tracing 拦截器与 REST 的 TraceHandler 功能一致,但实现方式因 gRPC 的 metadata 机制而有所不同:

1
2
3
4
5
6
7
8
9
10
11
// zrpc/internal/serverinterceptors/tracinginterceptor.go
func UnaryTracingInterceptor(ctx context.Context, req any, info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (any, error) {
ctx, span := startSpan(ctx, info.FullMethod)
defer span.End()

ztrace.MessageReceived.Event(ctx, 1, req)
resp, err := handler(ctx, req)
// ... 设置 span 状态和属性 ...
return resp, nil
}

startSpan 的实现揭示了一个关键差异——trace context 的传播方式。REST 中使用 HTTP Header 做 propagation,而 gRPC 使用 gRPC metadata

1
2
3
4
5
6
func startSpan(ctx context.Context, method string) (context.Context, trace.Span) {
md, ok := metadata.FromIncomingContext(ctx)
// ...提取 trace context 从 metadata 中...
tr := otel.Tracer(ztrace.TraceName)
return tr.Start(trace.ContextWithRemoteSpanContext(ctx, spanCtx), name, ...)
}

gRPC metadata 等价于 HTTP Header 的 key-value 结构,但它是通过 gRPC 协议原生传递的,不依赖 HTTP/2 帧头。OpenTelemetry 的 propagator 能将一个 metadata.MD 当作 TextMapCarrier 来处理,所以 REST 和 gRPC 之间的 trace context 可以无缝传播——这正是第 11 篇会深入讲解的主题。

Stream 版本的 Tracing 拦截器需要更多工作。因为流式 RPC 可能跨越多条消息,仅仅在拦截器入口处创建 span 是不够的——还需要追踪每条消息的收发:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
func (w *serverStream) RecvMsg(m any) error {
err := w.ServerStream.RecvMsg(m)
if err == nil {
w.receivedMessageID++
ztrace.MessageReceived.Event(w.Context(), w.receivedMessageID, m)
}
return err
}

func (w *serverStream) SendMsg(m any) error {
err := w.ServerStream.SendMsg(m)
w.sentMessageID++
ztrace.MessageSent.Event(w.Context(), w.sentMessageID, m)
return err
}

通过包装 grpc.ServerStream,在每次 RecvMsgSendMsg 时记录 trace event,最终在 span 中可以看到流式调用的每条消息的时间和内容。

RecoverInterceptor:比 REST 版本更简洁

1
2
3
4
5
6
7
8
// zrpc/internal/serverinterceptors/recoverinterceptor.go
func UnaryRecoverInterceptor(ctx context.Context, req any, _ *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (resp any, err error) {
defer handleCrash(func(r any) {
err = toPanicError(ctx, r)
})
return handler(ctx, req)
}

两个差异值得注意。第一,不像 REST 的 RecoverHandler 那样写 http.StatusInternalServerError 和堆栈到 ResponseWriter,这里的 toPanicError 返回的是 status.Errorf(codes.Internal, "panic: %v", r)——将一个 recover 到的 panic 转换为标准的 gRPC 错误。gRPC 框架会自动将这个 status.Error 序列化为 grpc-statusgrpc-message trailer,客户端通过 status.FromError(err) 就能还原出状态码和消息。

第二,这里的 Recover 在服务端拦截器链的第二位,在 Timeout 的外面。还记得为什么吗?我们稍后在 Timeout 拦截器中会看到,它内部已经通过 panicChan 做了一次 panic 捕获。但 Timeout 是可开关的——当 Timeout=0 时超时被禁用,此时 Recover 就是唯一的 panic 防线。所以放在 Stat 等拦截器外面是正确的——Cover 的范围最大化。

StatInterceptor:日志与指标的双重身份

Stat 拦截器相当于 REST 中 LogHandler + MetricHandler 的合体:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// zrpc/internal/serverinterceptors/statinterceptor.go
func UnaryStatInterceptor(metrics *stat.Metrics, conf StatConf) grpc.UnaryServerInterceptor {
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (resp any, err error) {
startTime := timex.Now()
defer func() {
duration := timex.Since(startTime)
metrics.Add(stat.Task{Duration: duration}) // 指标收集
logDuration(ctx, info.FullMethod, req, duration,
staticNotLoggingContentMethods, conf.SlowThreshold) // 日志记录
}()
return handler(ctx, req)
}
}

metrics.Add 将请求耗时写入 stat.Metrics,供后台每分钟聚合一次并打印统计报告(QPS、延迟分位数等)。logDuration 则提供了请求级别的日志:

  • 正常请求(耗时 < 慢阈值):记录请求内容和耗时。
  • 慢请求(耗时 >= 慢阈值,默认 500ms):以 Slowf 级别记录,日志中带上 slowcall 标记。生产中可以通过 grep slowcall 快速发现性能瓶颈。
  • 忽略内容的请求:通过 IgnoreContentMethods 配置指定的方法(通常是大请求体或二进制数据的方法),日志中只记录方法名和耗时,不记录请求体。

慢请求日志还附带客户端地址——通过 peer.FromContext(ctx) 从 gRPC 上下文中提取对端 IP,方便定位是哪个客户端触发的慢调用。

PrometheusInterceptor:面向 Prometheus 的标准化指标

1
2
3
4
5
6
7
8
9
// zrpc/internal/serverinterceptors/prometheusinterceptor.go
func UnaryPrometheusInterceptor(ctx context.Context, req any,
info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (any, error) {
startTime := timex.Now()
resp, err := handler(ctx, req)
metricServerReqDur.Observe(timex.Since(startTime).Milliseconds(), info.FullMethod)
metricServerReqCodeTotal.Inc(info.FullMethod, strconv.Itoa(int(status.Code(err))))
return resp, err
}

暴露 rpc_server_requests_duration_msrpc_server_requests_code_total 两个指标。与 REST 版本的关键区别在于标签维度——这里只按 method(gRPC 完整方法名,如 /greet.GreetService/SayHello)和 code(gRPC 状态码)切分,没有 HTTP path 和 HTTP method 的维度。gRPC 的全限定方法名已经包含了足够的信息,不需要像 REST 那样需要 path+method 双重标签。

为什么在 Stat 之后? Prometheus 和 Stat 都会统计延迟,但它们的作用不同。Prometheus 是"拉模式"的标准化指标——通过 /metrics 端点被 Prometheus Server 定时抓取,适合 Grafana 图表和 PromQL 告警。Stat 是"推模式"的本地统计——每分钟打印一行聚合报告,适合开发阶段快速扫一眼终端。两者放在相近的位置(都在 Trace/Recover 内侧),保证了它们测量的范围一致——都只统计真正进入 handler 的请求。

BreakerInterceptor:基于方法名的独立熔断

服务端断路器与 REST 版本的差异体现在两个地方:

1
2
3
4
5
6
7
8
9
10
11
12
// zrpc/internal/serverinterceptors/breakerinterceptor.go
func UnaryBreakerInterceptor(ctx context.Context, req any, info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (resp any, err error) {
breakerName := info.FullMethod
err = breaker.DoWithAcceptableCtx(ctx, breakerName, func() error {
var err error
resp, err = handler(ctx, req)
return err
}, serverSideAcceptable)

return resp, convertError(err)
}

第一,断路器的命名。 REST 断路器以 method + path 命名,zRPC 直接使用 info.FullMethod(如 /greet.GreetService/SayHello)。这是 gRPC 的优点——全限定方法名就是天然的、唯一的断路器标识。

第二,错误可接受性判断。 serverSideAcceptable 是一个函数,它决定了哪些错误是"可接受的"(不计入断路器失败计数):

1
2
3
4
5
6
func serverSideAcceptable(err error) bool {
if errorx.In(err, context.DeadlineExceeded, breaker.ErrServiceUnavailable) {
return false
}
return codes.Acceptable(err)
}

这个分两步判断:

  • context.DeadlineExceeded:超时错误不可接受。函数式调用超时通常意味着下游处理能力不足——这不是偶然错误,是真实的容量信号。
  • breaker.ErrServiceUnavailable:断路器本身的"熔断中"错误不可接受。如果不排除这类错误,熔断器自身的拒绝行为会被计为新的失败,导致熔断器永远无法恢复。
  • codes.Acceptable(err):对其他 gRPC 错误做分类判定。具体规则是——DeadlineExceededInternalUnavailableDataLossUnimplementedResourceExhausted 这些错误码是不可接受的;其他错误码(如 InvalidArgumentNotFoundPermissionDeniedOK)是可接受的

这个分类背后的设计思想是:断路器关心的是"服务是否健康",而不是"请求是否正确"。 客户端传入非法参数导致的 InvalidArgument 不应该触发熔断——那只是客户端的错误,不代表服务端出了问题。而 Internal(服务端内部错误)和 Unavailable(服务不可用)是真正的健康信号。

特别要注意 Unimplemented 也在不可接受列表中。如果客户端调了一个服务端没有实现的方法,这可能意味着版本不匹配——在滚动更新期间,旧服务端可能还没有部署新方法。将这种错误视为不可接受,可以在版本不一致时触发熔断,保护客户端不被持续的协议错误拖垮。

第三,错误转换。 断路器拒绝请求时返回的是 breaker.ErrServiceUnavailable,这个错误对 gRPC 客户端是需要被正确理解的:

1
2
3
4
5
6
func convertError(err error) error {
if errors.Is(err, breaker.ErrServiceUnavailable) {
return status.Error(gcodes.Unavailable, err.Error())
}
return err
}

将其转换为 gRPC 标准的 codes.Unavailable 状态码,客户端才能通过 status.Code(err) 正确识别并触发自己的 retry 或 fallback 逻辑。

SheddingInterceptor:CPU 过载时的自适应降载

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// zrpc/internal/serverinterceptors/sheddinginterceptor.go
func UnarySheddingInterceptor(shedder load.Shedder, metrics *stat.Metrics) grpc.UnaryServerInterceptor {
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (val any, err error) {
sheddingStat.IncrementTotal()
promise, err = shedder.Allow()
if err != nil {
metrics.AddDrop()
sheddingStat.IncrementDrop()
err = status.Error(codes.ResourceExhausted, err.Error())
return
}
// ...
return handler(ctx, req)
}
}

降载器被拒绝时返回 gRPC 标准的 codes.ResourceExhausted,语义上非常准确——“服务器资源耗尽,无法处理请求”。与 REST 中降载失败返回 HTTP 503 一样,这里返回的 ResourceExhausted 也进入了 codes.Acceptable 的判定——它是不可接受的,会被断路器计为失败。这是正确的——如果 CPU 持续过载导致大部分请求被降载,断路器也应该最终跳闸,形成一个整体性的保护。

降载器的成功/失败判定与 REST 一致:context.DeadlineExceeded 算失败(可能是 CPU 过高导致处理超时),其他情况算通过。

TimeoutInterceptor:比 REST 更优雅的超时处理

zRPC 服务端的超时拦截器有完整的 panic 安全机制和方法级超时定制:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
// zrpc/internal/serverinterceptors/timeoutinterceptor.go
func UnaryTimeoutInterceptor(timeout time.Duration,
methodTimeouts ...MethodTimeoutConf) grpc.UnaryServerInterceptor {
timeouts := buildMethodTimeouts(methodTimeouts)
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (any, error) {
t := getTimeoutByUnaryServerInfo(info.FullMethod, timeouts, timeout)
ctx, cancel := context.WithTimeout(ctx, t)
defer cancel()

var resp any
var err error
var lock sync.Mutex
done := make(chan struct{})
panicChan := make(chan any, 1)
go func() {
defer func() {
if p := recover(); p != nil {
panicChan <- fmt.Sprintf("%+v\n\n%s", p,
strings.TrimSpace(string(debug.Stack())))
}
}()
lock.Lock()
defer lock.Unlock()
resp, err = handler(ctx, req)
close(done)
}()

select {
case p := <-panicChan:
panic(p)
case <-done:
lock.Lock()
defer lock.Unlock()
return resp, err
case <-ctx.Done():
// ... 返回 Canceled 或 DeadlineExceeded 状态码
}
}
}

这个实现有三个值得深入的点:

第一,方法级超时。 通过 MethodTimeouts 配置,可以为不同方法设置不同的超时时间。例如:

1
2
3
4
Timeout: 2000  # 全局默认 2 秒
MethodTimeouts:
- FullMethod: /greet.GreetService/LongRunningTask
Timeout: 10s

getTimeoutByUnaryServerInfo 先查找方法级配置,找不到再用全局默认值。这是一个比 REST 更灵活的机制——REST 的超时只能按路由组设置,zRPC 可以精确到每个方法。

第二,内置 panic 保护。 注意 defer recover() 在 goroutine 内部——与 REST 的 TimeoutHandler 不同,zRPC 的 TimeoutInterceptor 自己在 handler goroutine 中捕获 panic,通过 panicChan 传递给主 goroutine 后重新抛出。但这并不会让 Recover 变成多余——当 Timeout=0(超时被禁用)时,TimeouterInterceptor 根本不会被注册,Recover 仍是唯一的防线。

第三,context.Canceledcontext.DeadlineExceeded 的区分。 与 REST 中 499 vs 503 的区分一样,zRPC 将 Canceled 映射为 gRPC 的 codes.Canceled,将 DeadlineExceeded 映射为 codes.DeadlineExceeded。这两个错误码在 codes.Acceptable 中有不同的待遇——DeadlineExceeded 不可接受(触发断路器),但 Canceled 不在列表中(实际上它不在 serverSideAcceptable 的检查范围内,会被 codes.Acceptable 中的 default 分支标记为可接受)。等一下——让我再确认一下。

实际上 Canceled 不在 codes.Acceptable 的显式匹配列表中(codes.Canceled 没有被列在 DeadlineExceeded, Internal, Unavailable, DataLoss, Unimplemented, ResourceExhausted 中),所以它走到 default 分支返回 true——可接受。这与 REST 的设计一致——客户端主动取消连接(用户关闭了浏览器/客户端超时了)不应该计为服务端的失败。

AuthInterceptor:App/Token 认证

认证拦截器的实现本身很简单——从 gRPC metadata 中提取 app 和 token,验证通过后放行:

1
2
3
4
5
6
7
8
9
10
// zrpc/internal/serverinterceptors/authinterceptor.go
func UnaryAuthorizeInterceptor(authenticator *auth.Authenticator) grpc.UnaryServerInterceptor {
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler) (any, error) {
if err := authenticator.Authenticate(ctx); err != nil {
return nil, err
}
return handler(ctx, req)
}
}

但认证器的内部实现更值得关注:

1
2
3
4
5
6
7
// zrpc/internal/auth/auth.go
func (a *Authenticator) validate(app, token string) error {
expect, err := a.cache.Take(app, func() (any, error) {
return a.store.Hget(a.key, app)
})
// ...
}

这里使用 collection.Cache(基于 TimingWheel 的本地缓存)缓存 app→token 的映射关系,默认过期时间 5 分钟。每次认证请求不直接查 Redis,而是先从本地缓存拿——Take 方法集成了 SingleFlight,保证同一个 app 的并发请求只会查一次 Redis。缓存未命中或过期时才回源 Redis。

严格模式(StrictControl: true)下,如果 Redis 不可达,认证直接失败——这是一种"宁可拒绝所有请求也不放过未认证请求"的安全策略。非严格模式下,Redis 不可达时认证通过(放行)——优先保障可用性,适用于内部服务间调用的场景。

客户端:从 Target 解析到连接建立

看完了服务端,把视角切换到调用方。NewClient 的创建流程主要做四件事:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
// zrpc/client.go
func NewClient(c RpcClientConf, options ...ClientOption) (Client, error) {
var opts []ClientOption
// 1. 认证凭证
if c.HasCredential() {
opts = append(opts, WithDialOption(grpc.WithPerRPCCredentials(&auth.Credential{...})))
}
// 2. 阻塞/非阻塞模式
if c.NonBlock { opts = append(opts, WithNonBlock())
} else { opts = append(opts, WithBlock()) }
// 3. 超时
if c.Timeout > 0 { opts = append(opts, WithTimeout(...)) }
// 4. Keepalive
if c.KeepaliveTime > 0 { opts = append(opts, WithDialOption(...)) }

// 5. 负载均衡策略
svcCfg := makeLBServiceConfig(c.BalancerName)
opts = append(opts, WithDialOption(grpc.WithDefaultServiceConfig(svcCfg)))

// 6. 构建 target
target, err := c.BuildTarget()

// 7. 建立连接
client, err := internal.NewClient(target, c.Middlewares, opts...)
}

Target 的构建与 Resolver 注册

BuildTarget 我们已经在配置部分讨论过——三选一:直连 target、自定义 target、etcd target。但这里的 target 字符串并不是最终传给 grpc.DialContext 的直连地址,而是 gRPC 的 URI 格式

1
2
3
4
5
// 直连模式
direct:///127.0.0.1:8080,127.0.0.1:8081

// etcd 模式
discov://etcd-host:2379/greet.rpc

gRPC 根据 target 的 scheme(direct://discov:// 等)选择对应的 resolver。go-zero 在 init() 中注册了所有这些 resolver:

1
2
3
4
// zrpc/internal/client.go
func init() {
resolver.Register()
}

resolver.Register() 内部调用了所有 builder 的注册函数,将 directdiscovetcdk8s 等 scheme 绑定到各自的 resolver builder。这部分细节我们留到下一篇(第 09 篇服务发现与负载均衡)中深入展开。

连接建立与 dial 流程

internal.NewClient 调用 dial 建立连接:

1
2
3
4
5
6
7
8
// zrpc/internal/client.go
func (c *client) dial(server string, opts ...ClientOption) error {
options := c.buildDialOptions(opts...)
timeCtx, cancel := context.WithTimeout(context.Background(), dialTimeout)
defer cancel()
conn, err := grpc.DialContext(timeCtx, server, options...)
// ... 错误处理 ...
}

这里的 dialTimeout 固定为 3 秒——不是配置中的 TimeoutTimeout 是每次 RPC 调用的超时,而 dialTimeout 是连接建立阶段的超时。如果 3 秒内连不上目标地址,直接报错。这避免了在服务发现失败或网络不可达时无限期阻塞。

错误信息中有一个巧妙的处理——当超时发生时,把 target 字符串的最后一段(/ 之后的部分,即服务名)提取出来,生成更友好的错误提示:

1
2
3
4
5
6
7
8
if errors.Is(err, context.DeadlineExceeded) {
pos := strings.LastIndexByte(server, separator)
if 0 < pos && pos < len(server)-1 {
service = server[pos+1:]
}
}
return fmt.Errorf("rpc dial: %s, error: %s, make sure rpc service %q is already started",
server, err.Error(), service)

这让排查连接问题更容易——日志中能看到具体的服务名而不是冗长的 etcd 地址。

客户端拦截器的装配

buildUnaryInterceptors 按固定顺序装配五个客户端拦截器:

1
2
3
4
5
6
7
8
9
func (c *client) buildUnaryInterceptors(timeout time.Duration) []grpc.UnaryClientInterceptor {
var interceptors []grpc.UnaryClientInterceptor
if c.middlewares.Trace { interceptors = append(interceptors, clientinterceptors.UnaryTracingInterceptor) }
if c.middlewares.Duration { interceptors = append(interceptors, clientinterceptors.DurationInterceptor) }
if c.middlewares.Prometheus { interceptors = append(interceptors, clientinterceptors.PrometheusInterceptor) }
if c.middlewares.Breaker { interceptors = append(interceptors, clientinterceptors.BreakerInterceptor) }
if c.middlewares.Timeout { interceptors = append(interceptors, clientinterceptors.TimeoutInterceptor(timeout)) }
return interceptors
}

最终形成的客户端拦截器链:

1
2
【Unary 拦截器链(客户端)】
Tracing → Duration → Prometheus → Breaker → Timeout → Invoker(实际发起 RPC 调用)

注意这个顺序与服务端不完全对称。服务端是 Trace → Recover → Stat → Prometheus → Breaker → Shedding → Timeout → Auth,而客户端是 Trace → Duration → Prometheus → Breaker → Timeout。我们逐个对比。

客户端拦截器逐个剖析

TracingInterceptor(客户端) 创建 client span,将 trace context inject 到 gRPC outgoing metadata 中——这是跨服务链路追踪的另一半。服务端的 TracingInterceptor 做 extract,客户端的 TracingInterceptor 做 inject,两者配合完成 trace context 的跨进程传播:

1
2
3
4
5
6
7
8
9
10
// zrpc/internal/clientinterceptors/tracinginterceptor.go
func startSpan(ctx context.Context, method, target string) (context.Context, trace.Span) {
md, ok := metadata.FromOutgoingContext(ctx)
tr := otel.Tracer(ztrace.TraceName)
name, attr := ztrace.SpanInfo(method, target)
ctx, span := tr.Start(ctx, name, trace.WithSpanKind(trace.SpanKindClient), ...)
ztrace.Inject(ctx, otel.GetTextMapPropagator(), &md)
ctx = metadata.NewOutgoingContext(ctx, md)
return ctx, span
}

DurationInterceptor(客户端) 相当于服务端 Stat 中日志部分的"镜像"。它记录每次调用的耗时、解析慢调用、在失败时打印完整错误日志:

1
2
3
4
5
6
7
8
9
10
11
12
13
func DurationInterceptor(ctx context.Context, method string, req, reply any,
cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
start := timex.Now()
err := invoker(ctx, method, req, reply, cc, opts...)
if err != nil {
// 失败:记录请求内容 + 错误
logger.Errorf("fail - %s - %v - %s", serverName, req, err.Error())
} else if elapsed > slowThreshold.Load() {
// 慢调用:记录请求和响应
logger.Slowf("[RPC] ok - slowcall - %s - %v - %v", serverName, req, reply)
}
return err
}

在调用失败时,Duration 记录了完整的请求内容——这对于排查"下游收到了什么导致失败"非常有用。而对于忽略日志内容的方法名(通过 DontLogClientContentForMethod 设置),失败时只记录方法名和错误信息。

PrometheusInterceptor(客户端) 暴露 rpc_client_requests_duration_msrpc_client_requests_code_total,与服务端的 rpc_server_requests_* 形成对称的指标对。两者结合可以在 Prometheus 中构建"客户端视角 vs 服务端视角"的延迟和错误率对比——如果客户端 p99 远高于服务端 p99,说明网络延迟是瓶颈,而非服务端处理能力。

BreakerInterceptor(客户端) 使用的是 gRPC 拦截器的签名,但它复用了同一个 breaker 实现和同一套 codes.Acceptable 判定:

1
2
3
4
5
6
7
func BreakerInterceptor(ctx context.Context, method string, req, reply any,
cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
breakerName := path.Join(cc.Target(), method)
return breaker.DoWithAcceptableCtx(ctx, breakerName, func() error {
return invoker(ctx, method, req, reply, cc, opts...)
}, codes.Acceptable)
}

注意断路器名字是用 cc.Target() 和 method 拼接的——etcd-host:2379/greet.rpc/greet.GreetService/SayHello。这与服务端的 info.FullMethod 不同——客户端熔断是按目标服务+方法粒度的,同一个方法调用不同下游实例不会被熔断连累。

这里直接使用 codes.Acceptable(不经过 serverSideAcceptablecontext.DeadlineExceeded 额外排除)。这意味着客户端侧,DeadlineExceeded 仍然是不可接受的——客户端超时通常意味着下游响应太慢,这正是断路器应该保护的场景。

TimeoutInterceptor(客户端) 实现比服务端简单,但支持一个独特的"单次调用覆盖"能力:

1
2
3
4
5
6
7
8
9
func TimeoutInterceptor(timeout time.Duration) grpc.UnaryClientInterceptor {
return func(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn,
invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
t := getTimeoutFromCallOptions(opts, timeout)
ctx, cancel := context.WithTimeout(ctx, t)
defer cancel()
return invoker(ctx, method, req, reply, cc, opts...)
}
}

getTimeoutFromCallOptions 遍历 grpc.CallOption 列表,检查是否有 TimeoutCallOption(通过 WithCallTimeout 传入)。如果有,就用单次调用的超时时间覆盖默认值。这意味着同一客户端可以在每次调用时动态决定超时时间——不需要为长查询和短查询分别创建客户端实例:

1
2
3
4
5
// 短查询使用默认超时(配置中的 500ms)
client.SayHello(ctx, &req)

// 长查询使用 10 秒超时
client.SayHello(ctx, &req, zrpc.WithCallTimeout(10*time.Second))

错误码判定:断路器决策的核心依据

在服务端和客户端的断路器拦截器中,我们多次提到了 codes.Acceptable——它决定了哪些 gRPC 错误不计入断路器的失败计数。这里集中深入分析一下。

1
2
3
4
5
6
7
8
9
10
// zrpc/internal/codes/accept.go
func Acceptable(err error) bool {
switch status.Code(err) {
case codes.DeadlineExceeded, codes.Internal, codes.Unavailable,
codes.DataLoss, codes.Unimplemented, codes.ResourceExhausted:
return false
default:
return true
}
}

这个函数区分了两类错误:

不可接受的错误(6 个),计入断路器失败计数:

错误码 含义 为什么不可接受
DeadlineExceeded 处理超时 通常是下游处理能力不足的信号
Internal 服务端内部错误 服务本身出了问题
Unavailable 服务不可达 服务健康状态的直接信号
DataLoss 数据丢失 严重错误,可能磁盘或存储故障
Unimplemented 方法未实现 版本不匹配信号,可能整个服务需要更新
ResourceExhausted 资源耗尽(被降载) CPU 或内存过载,需要容量保护

可接受的错误(其余所有),不计入断路器失败计数:

错误码示例 含义 为什么可接受
InvalidArgument 参数错误 客户端问题,不是服务端故障
NotFound 资源不存在 正常的业务结果
PermissionDenied 权限不足 认证/授权问题,非故障
OK 成功 当然可接受
Canceled 客户端取消 客户端行为,非故障

这个分类的精妙之处在于:断路器只关心"服务是否健康",不关心"请求是否正确"。 如果客户端每天传 100 次非法参数,每次都返回 InvalidArgument,把它们计为失败会导致断路器跳闸,正常请求也被拒绝——这就产生了"错误参数的级联故障"。区分两类错误,保证了断路器的准确性。

RpcProxy:为不同的 App/Token 缓存连接

除了标准的 RpcClient,zRPC 还提供了一个 RpcProxy——它为每次带有不同认证凭证的请求自动管理连接池:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
// zrpc/proxy.go
type RpcProxy struct {
backend string
clients map[string]Client
options []internal.ClientOption
singleFlight syncx.SingleFlight
lock sync.Mutex
}

func (p *RpcProxy) TakeConn(ctx context.Context) (*grpc.ClientConn, error) {
cred := auth.ParseCredential(ctx)
key := cred.App + "/" + cred.Token
val, err := p.singleFlight.Do(key, func() (any, error) {
p.lock.Lock()
client, ok := p.clients[key]
p.lock.Unlock()
if ok { return client, nil }

client, err := NewClientWithTarget(p.backend, opts...)
// ...
p.clients[key] = client
return client, nil
})
return val.(Client).Conn(), nil
}

设计思路很清晰:每个 App/Token 组合对应一个连接,通过 SingleFlight 保证并发安全且不重复创建。ParseCredential 从 gRPC incoming metadata 中提取 app 和 token——这个函数定义在 auth/credential.go 中,与认证拦截器共享同一套 metadata key。

当 Gateway 或 HTTP 代理层需要将多个租户的请求转发到同一个后端 RPC 服务时,RpcProxy 比每次创建新的 RpcClient 更高效——连接复用,不会因为请求量大而耗尽文件描述符。

两种调用走一遍:成功与超时

成功的 Unary RPC 调用

现在把客户端和服务端串联起来,追踪一次完整的 RPC 调用:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
【客户端】
业务代码: client.SayHello(ctx, &req)
→ Tracing: 创建 client span, inject trace context 到 outgoing metadata
→ Duration: 启动计时
→ Prometheus: 启动计时
→ Breaker: Allow() 成功, 放行
→ Timeout: ctx.WithTimeout(2 秒)
→ invoker (gRPC 底层): 序列化 protobuf → 发送 HTTP/2 DATA 帧
...

【网络传输】 gRPC metadata + protobuf body 通过 HTTP/2 传输

【服务端】
gRPC Server 接收请求
→ Tracing: extract trace context 从 incoming metadata, 创建 server span
→ Recover: defer recover
→ Stat: 启动计时
→ Prometheus: 启动计时
→ Breaker: Allow() 成功
→ Shedding: Allow() 成功
→ Timeout: ctx.WithTimeout(2 秒), 启动 handler goroutine
→ Auth: app/token 验证通过
→ Handler: 业务逻辑执行 (~15ms)
← 返回 resp, nil
← 返回
← promise.Pass() (无错误)
← promise.Accept()
← Prometheus: observe duration + code="0" (OK)
← Stat: metrics.Add(duration); logDuration: [RPC] 127.0.0.1:52341 - /greet.GreetService/SayHello - 15ms
← (recover 未触发)
← span: SetStatus(OK), End()

【网络传输】 gRPC response + trailer (grpc-status: 0)

【客户端】
← invoker 返回
← Timeout: cancel()
← Breaker: promise.Accept() (err == nil)
← Prometheus: observe duration + code="0"
← Duration: elapsed=15ms < slow threshold, 不打印日志
← Tracing: MessageReceived event, span.SetStatus(OK), End()
业务代码: resp, err := ... ← err == nil

这个流程展示了 zRPC 拦截器链的核心价值——一次业务调用在框架层面上自动获得了 trace 传播、Prometheus 指标、熔断保护、降载保护、超时控制和结构化日志。业务代码只需要写 client.SayHello(ctx, &req)return &resp, nil

超时的 Unary RPC 调用

如果下游服务处理太慢,2 秒后超时触发:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
【服务端 Timeout goroutine 内】
Handler 执行中... → 2 秒超时触发
→ ctx.Done() ← context.DeadlineExceeded
→ 返回 status.Error(codes.DeadlineExceeded, ...)
→ handler goroutine 中的 lock 最终释放(handler 返回)
→ Shedding: errors.Is(err, context.DeadlineExceeded) → promise.Fail()
→ Breaker: serverSideAcceptable → context.DeadlineExceeded → false (不可接受) → promise.Reject()
→ Prometheus: code="4" (DeadlineExceeded)
→ Stat: 记录 slowcall 日志 (elapsed >= 500ms)

【客户端】
invoker 返回 DeadlineExceeded 错误
→ Breaker: codes.Acceptable → DeadlineExceeded → false (不可接受) → promise.Reject()
→ Prometheus: code="4"
→ Duration: logger.Errorf("fail - ... - context deadline exceeded")
→ Tracing: span.SetStatus(Error, "context deadline exceeded")

这里的关键观察是:一个超时触发了双方断路器的失败计数。 服务端断路器收到 DeadlineExceeded → Reject,客户端断路器也收到 DeadlineExceeded → Reject。如果下游持续超时,双方的断路器会先后跳闸——服务端先跳(因为超时发生在服务端处理过程中),客户端接着跳(因为后续调用被服务端断路器直接拒绝返回 Unavailable)。这是一套自洽的保护链路。

总结

本文以 RpcServerRpcClient 的创建为入口,完整走完了 zRPC 的服务端和客户端调用链:

服务端配置驱动拦截器链。 8 个拦截器按 Tracing → Recover → Stat → Prometheus → Breaker → Shedding → Timeout → Auth 的顺序叠加。这个顺序与 REST 中间件链共享同一套分层逻辑——观测层在最外(Tracing、Stat、Prometheus),保护层在中间(Breaker、Shedding、Timeout),恢复层内嵌在保护层中(Recover 在 Timeout 可能被关闭时提供兜底),认证层在最内(Auth 决定"是否能访问",在一切保护之后)。

客户端拦截器链是服务端的简化镜像。 5 个拦截器(Tracing → Duration → Prometheus → Breaker → Timeout)保证了每次调用的可观测性和弹性保护。客户端不需要 Recover(不执行业务逻辑)、不需要 Shedding(降载是服务端的事)、不需要 Auth(认证凭证在连接建立时通过 WithPerRPCCredentials 注入)。

错误码判定是断路器准确性的核心。 codes.Acceptable 将 6 个 gRPC 错误码标记为"不可接受",其余为"可接受"。这个分类体现了"断路器关心服务健康,不关心请求是否正确"的设计原则,避免了参数错误等客户端问题触发级联熔断。

zRPC 与标准 gRPC 的关系。 zRPC 没有魔改 gRPC 的 stub 生成、序列化或传输层——它完全在拦截器这个 gRPC 原生的扩展点之上构建增值能力。这保证了与 gRPC 生态的完全兼容——你可以继续使用 protoc 生成的代码、标准的 grpc.ClientConn 和任何 gRPC 工具链。

从下一篇开始,我们将深入服务发现与负载均衡——在客户端拦截器链的 Timeout 和 invoker 之间,gRPC 的 resolver 和 balancer 是如何根据 etcd watch 结果更新地址列表、P2C 算法如何选择最优节点的。