This repository has been archived on 2026-05-24. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
all-in-kingsoft/hhs/gRPC/5. 中间件与拦截器/14-Unary 与 Stream 拦截器.md
T
2026-05-11 19:02:38 +08:00

284 lines
10 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
tags: [gRPC, Interceptor, Middleware, Go]
create time: 2026-05-11 16:00
---
# Unary 与 Stream 拦截器
## 概述
Interceptor 是 gRPC 的「插件系统」——在每个 RPC 调用执行前后注入逻辑。它和 HTTP middleware 概念类似,但接口更底层、更灵活。本文档完整覆盖服务端和客户端的 Unary / Stream 拦截器签名、链式调用原理、Recovery、错误码映射等核心模式——理解这些是你实现鉴权、日志、重试的前提。
> [!tip] Interceptor 是单例
> Interceptor 在 server/client 初始化时注册一次,之后对每个请求生效。不要在 interceptor 里持有 per-request 状态。
## 正文
### Unary Interceptor 签名
Unary(普通 RPC)拦截器的核心类型如下:
```go
type UnaryServerInterceptor func(
ctx context.Context,
req interface{},
info *UnaryServerInfo,
handler UnaryHandler,
) (interface{}, error)
type UnaryServerInfo struct {
Server string
FullMethod string // e.g. "user.v1.UserService/CreateUser"
}
```
`handler` 就是真正的业务方法实现。你可以选择在调用 handler 之前做任何事(比如鉴权),也可以在调用之后做后处理(比如记录日志)。关键技巧:**你可以在调用 handler 之前或之后插入逻辑**。
```go
func MyInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
// BEFORE: 前置逻辑 — 鉴权、校验、埋点
result, err := handler(ctx, req) // 调用真正 handler
// AFTER: 后置逻辑 — 日志、指标、错误处理
return result, err
}
```
### Stream Interceptor 签名
Streaming RPC 的拦截器有所不同,因为数据是通过流传递的:
```go
type StreamServerInterceptor func(
srv interface{},
ss ServerStream,
info *StreamServerInfo,
handler StreamHandler,
) error
```
注意几点差异:
- 第一个参数是 `srv`(服务实例),而非 `ctx`——stream 的 context 通过 `ss.Context()` 获取
- 返回的是整条 stream 的错误,不是单个 message 的错误
- 你无法直接修改发送/接收的消息内容
### 客户端拦截器
服务端拦截器处理入站请求,而客户端拦截器包裹出站调用。它们的签名略有不同:
```go
// 客户端 Unary
type UnaryClientInterceptor func(
ctx context.Context,
method string,
req any,
reply any,
cc *grpc.ClientConn,
invoker grpc.UnaryInvoker,
opts ...grpc.CallOption,
) error
// 客户端 Stream
type StreamClientInterceptor func(
ctx context.Context,
desc *StreamDesc,
cc *grpc.ClientConn,
method string,
streamer grpc.Streamer,
opts ...grpc.CallOption,
) (grpc.ClientStream, error)
```
关键差异:
- `method` 是路径名如 `/user.v1.UserService/CreateUser`,非完整 method string
- `req` 和 `reply` 都是 `any`——你可以反序列化后检查响应内容
- `opts ...grpc.CallOption` 允许链式追加 CallOption(比如超时、metadata)
- Client Stream Interceptor 返回 `grpc.ClientStream`,而非 `error`——真正的错误在后续收发消息时抛出
```go
func TimeoutInterceptor(ctx context.Context, method string, req any, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
ctx, cancel := context.WithTimeout(ctx, time.Second*5)
defer cancel()
return invoker(ctx, method, req, reply, cc, opts...)
}
```
这段代码自动为每个 RPC 调用添加 5 秒超时,无需手动在每个 call 中设置——这是客户端拦截器最常见的用途之一。
### 手动构建 Interceptor Chain
gRPC Go 原生支持链式调用,我们先手动实现一个 chain 来理解其原理:
```go
func chainUnaryInterceptors(interceptors ...grpc.UnaryServerInterceptor) grpc.UnaryServerInterceptor {
n := len(interceptors)
return func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
ch := handler
for i := n - 1; i >= 0; i-- {
finalHandler := ch
ch = func(c context.Context, r interface{}) (interface{}, error) {
return interceptors[i](c, r, info, finalHandler)
}
}
return ch(ctx, req)
}
}
```
这段代码的关键在于从右到左包裹——最后一个 interceptor 最先被传入,离 handler 最近。请求进来时执行顺序是 **A → B → C → handler**,返回时反向通过每一层。**外层 interceptor 能捕获内层的一切异常**(包括 panic 和 error),这正是链式拦截器的核心设计。
> [!question] 为什么循环要从 n-1 到 0?
> 因为最后一个 interceptor 应该最先执行(最靠近 handler),这样才能保证第一个 interceptor 在最外层捕获所有下游异常。
### gRPC 官方推荐方式
实际使用中直接使用 gRPC 内置的 chain 函数:
```go
server := grpc.NewServer(
grpc.ChainUnaryInterceptor(
logInterceptor, // 第 1 层(最外层)
authInterceptor, // 第 2 层
recoveryInterceptor, // 第 3 层(最内层)
),
grpc.ChainStreamInterceptor(
logStreamInterceptor,
authStreamInterceptor,
),
)
```
执行顺序与手动 chain 一致:请求到达时从左到右依次进入每一层(A → B → C → handler),返回时反向退出。**最外层最先看到请求、最后看到响应**。recovery interceptor 放在最内侧以捕获所有 panic。如果 recovery 放最外侧,它会先捕获其他 interceptor 抛出的异常而非让业务处理——这些不是 bug,而是不满足条件的正常错误流。
### 完整示例:Request Logger
下面是一个实用的请求日志 interceptor:
```go
func RequestLogger(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
start := time.Now()
resp, err := handler(ctx, req)
dur := time.Since(start)
log.Printf("rpc: %s %s %.2fs err=%v",
info.FullMethod,
reflect.TypeOf(req),
dur.Seconds(),
err,
)
return resp, err
}
```
这个 interceptor 做了三件事:记录开始时间、调用 handler、打印耗时和错误信息。它可以作为所有 interceptor 链的基础层。
### 实战示例:Panic Recovery
服务崩掉一个 handler 不应影响整个进程,用 `recover()` 兜底:
```go
func RecoveryInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
defer func() {
if r := recover(); r != nil {
log.Error("panic recovered", "error", r, "method", info.FullMethod)
}
}()
return handler(ctx, req)
}
```
> [!tip] Panic Recovery 必须放最内层
> 如果把 recovery 放在最外侧,它会吞掉其他 interceptor(比如 auth)主动返回的错误——这些错误不是 bug,不该被 recover。所以 recover 应该离 handler 最近,确保只捕获真正的 panic。
### 错误处理规范
Interceptor 中返回错误时,**不要直接返回 `errors.New`**——要用 `status.Errorf` 映射为 gRPC status code:
```go
import "google.golang.org/grpc/status"
// ✗ 错误做法
return nil, errors.New("user not found")
// ✓ 正确做法
return nil, status.Error(codes.NotFound, "user not found")
// ✓ 带细节的正确做法
return nil, status.Errorf(codes.InvalidArgument, "invalid email: %v", err)
```
| 场景 | 推荐 Code | 含义 |
|------|-----------|------|
| 参数校验失败 | `codes.InvalidArgument` | 客户端传参有问题 |
| 资源不存在 | `codes.NotFound` | ID 对应的记录不存在 |
| 未认证 | `codes.Unauthenticated` | Token 缺失或无效 |
| 无权限 | `codes.PermissionDenied` | 认证通过但无权访问 |
| 超时 | `codes.DeadlineExceeded` | 处理时间超出限制 |
| 内部错误 | `codes.Internal` | 服务端意外 panic 或 DB 故障 |
| 限流 | `codes.ResourceExhausted` | 超出速率上限 |
> [!tip] Client 侧重试判定
> 客户端 interceptor(如重试)依据 status code 决定是否重试:只有 `Unavailable`、`DeadlineExceeded`、`ResourceExhausted` 等可恢复 code 才触发重试。错误的 code 映射会导致不该重试的请求被反复发送。
### Interceptor vs StatsHandler
| 维度 | Interceptor | StatsHandler |
|------|-------------|--------------|
| 能力 | 修改 req/res、控制流程 | 纯观测(metrics/tracing) |
| 可写 | 可以改返回值 | 只读 |
| 性能 | 较高开销 | 更低(异步) |
| 适用 | Auth, Recovery, RateLimit | Metrics, Tracing, Profiling |
如果你在追求高性能的可观测性,优先选 StatsHandler;如果需要修改请求/响应或控制执行流程,Interceptor 是唯一选择。
### Context 传递规则
Interceptor 中可以向 context 注入信息,下游 handler 可以读取。服务端和客户端都有各自的传递方向:
```go
type ctxKey struct{}
func AuthInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
token := extractToken(ctx)
claims, _ := jwt.Parse(token)
ctx = context.WithValue(ctx, ctxKey{}, claims)
return handler(ctx, req)
}
```
**服务端**:从 incoming metadata 中提取身份信息 → 写入 context → 传给 handler
**客户端**:从 local context 读取 token → 写入 outgoing metadata → 发送给上游
注意事项:
- value key 必须定义为专用不可比较的类型(如上 `ctxKey` struct),避免包间冲突
- 不要在 interceptor 里阻塞或做耗时操作,否则会影响所有下游请求
- context value 不应传递大对象或敏感明文——token 解析后的 claims 可以传,原始密码不行
- 如果上游已经注入了相同 key 的 value,下游会覆盖它——确保 chain 中每个步骤使用唯一 key
### 常见 Interceptor 模式总结
```mermaid
flowchart LR
subgraph ClientChain["客户端链"]
A1["Timeout\n自动超时"] --> A2["Retry\n错误重试"]
end
subgraph ServerChain["服务端链"]
B1["Recovery\nPanic 兜底"] --> B2["Auth\n鉴权校验"] --> B3["Logger\n记录耗时"]
end
ClientChain -->|"gRPC call"| ServerChain
style A1 fill:#00B6BC,color:#fff
style A2 fill:#FFD43B
style B1 fill:#EE5A24,color:#fff
style B2 fill:#FFD43B
style B3 fill:#00B6BC,color:#fff
```
一个典型生产环境的 interceptor chain 结构如上:**客户端侧**做超时控制、重试容错;**服务端侧**做 panic 恢复、鉴权和日志。每一层职责单一,便于测试和维护。
## 关联笔记
- [[hhs/gRPC/5. 中间件与拦截器/15-元数据与鉴权]]
- [[hhs/gRPC/5. 中间件与拦截器/16-日志与链路追踪]]