Go gRPC 实战:Protobuf、拦截器与服务治理
gRPC 使用 Protocol Buffers 定义强类型接口,基于 HTTP/2 提供一元调用和流式通信。它适合服务间通信,但生产系统还需要截止时间、状态码、认证、TLS、健康检查和可观测性。
准备工具
安装 protoc 后,再安装 Go 代码生成插件:
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest
初始化模块并添加依赖:
go mod init example.com/user-service
go get google.golang.org/grpc
go get google.golang.org/protobuf
团队项目应在工具文件、构建脚本或 CI 镜像中固定生成器版本,避免不同开发环境生成不一致代码。
定义 Protobuf 接口
创建 api/user/v1/user.proto:
syntax = "proto3";
package user.v1;
option go_package = "example.com/user-service/gen/user/v1;userv1";
import "google/protobuf/timestamp.proto";
service UserService {
rpc GetUser(GetUserRequest) returns (GetUserResponse);
rpc ListUsers(ListUsersRequest) returns (stream User);
}
message GetUserRequest {
int64 id = 1;
}
message GetUserResponse {
User user = 1;
}
message ListUsersRequest {
int32 page_size = 1;
string page_token = 2;
}
message User {
int64 id = 1;
string name = 2;
string email = 3;
google.protobuf.Timestamp created_at = 4;
}
生成 Go 代码:
protoc \
--go_out=. --go_opt=paths=source_relative \
--go-grpc_out=. --go-grpc_opt=paths=source_relative \
api/user/v1/user.proto
字段编号一旦发布就不能复用。删除字段时应使用 reserved 保留编号和名称:
message User {
reserved 5;
reserved "legacy_name";
}
新增字段通常向后兼容,修改字段类型、语义或编号可能破坏客户端。协议变更应经过兼容性检查。
实现服务端
type Server struct {
userv1.UnimplementedUserServiceServer
users UserStore
}
func (s *Server) GetUser(ctx context.Context, request *userv1.GetUserRequest) (*userv1.GetUserResponse, error) {
if request.GetId() < 1 {
return nil, status.Error(codes.InvalidArgument, "id must be positive")
}
user, err := s.users.Find(ctx, request.GetId())
if errors.Is(err, ErrUserNotFound) {
return nil, status.Error(codes.NotFound, "user not found")
}
if err != nil {
return nil, status.Error(codes.Internal, "failed to load user")
}
return &userv1.GetUserResponse{
User: &userv1.User{
Id: user.ID,
Name: user.Name,
Email: user.Email,
CreatedAt: timestamppb.New(user.CreatedAt),
},
}, nil
}
服务端应把数据库错误转换为稳定 gRPC 状态码,不向客户端暴露 SQL 和内部路径。
常用状态码:
InvalidArgument:参数格式或范围错误。Unauthenticated:缺少或无法验证身份。PermissionDenied:身份有效但没有权限。NotFound:资源不存在。AlreadyExists:唯一资源冲突。FailedPrecondition:当前资源状态不允许操作。ResourceExhausted:配额、并发或容量不足。Unavailable:暂时不可用,客户端可能重试。DeadlineExceeded:未在截止时间内完成。Internal:未公开的内部错误。
启动 gRPC Server
listener, err := net.Listen("tcp", ":9090")
if err != nil {
return fmt.Errorf("listen: %w", err)
}
server := grpc.NewServer(
grpc.ChainUnaryInterceptor(
requestIDInterceptor,
authInterceptor,
loggingInterceptor,
recoveryInterceptor,
),
)
userv1.RegisterUserServiceServer(server, &Server{users: store})
if err := server.Serve(listener); err != nil {
return fmt.Errorf("serve gRPC: %w", err)
}
拦截器顺序会影响认证、日志和恢复行为,应保持明确并通过测试验证。
一元拦截器
func loggingInterceptor(
ctx context.Context,
req any,
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (any, error) {
start := time.Now()
response, err := handler(ctx, req)
code := status.Code(err)
slog.InfoContext(ctx, "grpc request",
"method", info.FullMethod,
"code", code.String(),
"duration", time.Since(start),
)
return response, err
}
不要默认记录完整请求和响应,其中可能包含密码、令牌或个人信息。日志和指标使用 FullMethod 作为低基数维度。
Panic 恢复:
func recoveryInterceptor(
ctx context.Context,
req any,
info *grpc.UnaryServerInfo,
handler grpc.UnaryHandler,
) (response any, err error) {
defer func() {
if recovered := recover(); recovered != nil {
slog.ErrorContext(ctx, "grpc panic",
"method", info.FullMethod,
"panic", recovered,
"stack", string(debug.Stack()),
)
err = status.Error(codes.Internal, "internal server error")
}
}()
return handler(ctx, req)
}
元数据与认证
客户端通过 Metadata 传递令牌:
ctx := metadata.AppendToOutgoingContext(ctx, "authorization", "Bearer "+token)
response, err := client.GetUser(ctx, request)
服务端读取并验证:
func authenticate(ctx context.Context) (Principal, error) {
values, ok := metadata.FromIncomingContext(ctx)
if !ok {
return Principal{}, status.Error(codes.Unauthenticated, "missing metadata")
}
authorization := values.Get("authorization")
if len(authorization) != 1 {
return Principal{}, status.Error(codes.Unauthenticated, "missing token")
}
return verifyToken(authorization[0])
}
Metadata 键会规范为小写。认证完成后可以把主体放入派生 Context,业务层不应重复解析令牌。
客户端连接与截止时间
credentials, err := credentials.NewClientTLSFromFile("ca.crt", "users.internal")
if err != nil {
return err
}
connection, err := grpc.NewClient(
"dns:///users.internal:9090",
grpc.WithTransportCredentials(credentials),
)
if err != nil {
return fmt.Errorf("create gRPC client: %w", err)
}
defer connection.Close()
client := userv1.NewUserServiceClient(connection)
ctx, cancel := context.WithTimeout(context.Background(), 800*time.Millisecond)
defer cancel()
response, err := client.GetUser(ctx, &userv1.GetUserRequest{Id: 42})
gRPC 默认不会替调用自动决定业务截止时间。每个调用都应继承上游 Context,并在总体预算内设置超时。
不要在每次请求中创建连接。ClientConn 可并发使用,应在进程生命周期内复用。
重试与幂等性
只有可重试状态和幂等操作才能自动重试。查询通常幂等;创建订单、扣款等写操作需要幂等键和服务端去重。
服务配置示例:
{
"methodConfig": [{
"name": [{"service": "user.v1.UserService", "method": "GetUser"}],
"timeout": "0.8s",
"retryPolicy": {
"MaxAttempts": 3,
"InitialBackoff": "0.05s",
"MaxBackoff": "0.2s",
"BackoffMultiplier": 2,
"RetryableStatusCodes": ["UNAVAILABLE"]
}
}]
}
重试会增加下游流量,应限制总尝试次数,并让退避包含抖动。服务端过载时盲目重试会形成重试风暴。
服务端流式响应
func (s *Server) ListUsers(request *userv1.ListUsersRequest, stream grpc.ServerStreamingServer[userv1.User]) error {
users, err := s.users.List(stream.Context(), int(request.GetPageSize()))
if err != nil {
return status.Error(codes.Internal, "failed to list users")
}
for _, user := range users {
if err := stream.Send(toProto(user)); err != nil {
return err
}
}
return nil
}
流式处理应持续检查 Context,限制单条消息大小和总流量,并考虑慢客户端带来的背压。不要一次把无限数据加载到内存后再发送。
生成代码的泛型流接口取决于 grpc-go 版本;升级生成器和运行库时应一起更新并重新生成代码。
TLS 与双向 TLS
服务端启用 TLS:
serverCredentials, err := credentials.NewServerTLSFromFile("server.crt", "server.key")
if err != nil {
return err
}
server := grpc.NewServer(grpc.Creds(serverCredentials))
内部零信任环境可以使用 mTLS,同时验证客户端证书。证书需要自动续期、轮换和过期监控,不能只在首次部署时手工生成。
健康检查与反射
gRPC 标准健康检查可供负载均衡器和编排平台使用:
healthServer := health.NewServer()
grpc_health_v1.RegisterHealthServer(server, healthServer)
healthServer.SetServingStatus("", grpc_health_v1.HealthCheckResponse_SERVING)
服务尚未完成依赖初始化或准备退出时,应切换为 NOT_SERVING。
反射便于 grpcurl 调试:
reflection.Register(server)
生产环境是否开放反射取决于网络边界和安全策略。
优雅退出
done := make(chan struct{})
go func() {
server.GracefulStop()
close(done)
}()
select {
case <-done:
case <-time.After(10 * time.Second):
server.Stop()
}
先从服务发现或负载均衡中摘除实例,再停止接收新请求。GracefulStop 等待现有 RPC 结束,超过预算后使用 Stop 强制关闭。
使用 bufconn 测试
bufconn 可以在内存连接上测试完整 gRPC 编解码和拦截器:
listener := bufconn.Listen(1024 * 1024)
server := grpc.NewServer()
userv1.RegisterUserServiceServer(server, testServer)
go server.Serve(listener)
t.Cleanup(server.Stop)
dialer := func(context.Context, string) (net.Conn, error) {
return listener.Dial()
}
connection, err := grpc.NewClient(
"passthrough:///bufnet",
grpc.WithContextDialer(dialer),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
测试参数校验、状态码、Metadata、认证、截止时间、拦截器以及流中途取消。
生产检查清单
- Proto 字段编号是否保持兼容,删除字段是否 reserved。
- 生成器和运行库版本是否固定并同步升级。
- 每个调用是否设置截止时间并传播 Context。
- 错误是否映射为准确状态码且不泄露内部信息。
- 认证、日志、指标和恢复是否由拦截器统一处理。
- 客户端连接是否复用,是否使用 TLS 或 mTLS。
- 重试是否仅用于幂等调用并限制总预算。
- 流式接口是否处理背压、取消和消息大小。
- 健康检查、优雅退出和部署摘流顺序是否完整。
gRPC 的强类型协议只是起点。只有把兼容性、截止时间、认证、状态码、重试和生命周期一起设计,服务间通信才能在规模扩大后仍然可控。