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、认证、截止时间、拦截器以及流中途取消。

生产检查清单

  1. Proto 字段编号是否保持兼容,删除字段是否 reserved。
  2. 生成器和运行库版本是否固定并同步升级。
  3. 每个调用是否设置截止时间并传播 Context。
  4. 错误是否映射为准确状态码且不泄露内部信息。
  5. 认证、日志、指标和恢复是否由拦截器统一处理。
  6. 客户端连接是否复用,是否使用 TLS 或 mTLS。
  7. 重试是否仅用于幂等调用并限制总预算。
  8. 流式接口是否处理背压、取消和消息大小。
  9. 健康检查、优雅退出和部署摘流顺序是否完整。

gRPC 的强类型协议只是起点。只有把兼容性、截止时间、认证、状态码、重试和生命周期一起设计,服务间通信才能在规模扩大后仍然可控。