gRPC基于HTTP/2和Protocol Buffers构建,提供高性能、强类型的微服务通信方案。相比REST的JSON文本序列化,gRPC使用二进制Protobuf编码,消息体积更小、解析速度更快。HTTP/2的多路复用和头部压缩机制使gRPC在内部服务间通信场景下性能优势明显,适合高吞吐、低延迟的微服务架构。
Protocol Buffers服务定义与代码生成
使用.proto文件定义服务接口和消息格式。一个典型的用户服务定义:
syntax = "proto3";
package user.v1;
option go_package = "github.com/example/proto/user/v1;userv1";
service UserService {
rpc GetUser(GetUserRequest) returns (GetUserResponse);
rpc ListUsers(ListUsersRequest) returns (stream User);
rpc CreateUserBatch(stream CreateUserRequest) returns (CreateUserBatchResponse);
rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}
message GetUserRequest {
int64 id = 1;
}
message GetUserResponse {
User user = 1;
}
message User {
int64 id = 1;
string name = 2;
string email = 3;
UserStatus status = 4;
int64 created_at = 5;
}
enum UserStatus {
USER_STATUS_UNSPECIFIED = 0;
USER_STATUS_ACTIVE = 1;
USER_STATUS_INACTIVE = 2;
USER_STATUS_BANNED = 3;
}
message ListUsersRequest {
int32 page_size = 1;
string page_token = 2;
UserStatus status_filter = 3;
}
message CreateUserRequest {
string name = 1;
string email = 2;
}
message CreateUserBatchResponse {
repeated int64 ids = 1;
}
message ChatMessage {
string content = 1;
int64 timestamp = 2;
}
gRPC支持四种调用模式:一元调用(Unary)、服务端流(Server Streaming)、客户端流(Client Streaming)和双向流(Bidirectional Streaming)。使用protoc工具生成各语言客户端和服务端代码:
protoc --go_out=. --go-grpc_out=. --proto_path=proto proto/user/v1/user.proto
Go服务端实现与注册
生成代码后实现服务接口,以Go为例:
package server
import (
"context"
"database/sql"
userv1 "github.com/example/proto/user/v1"
)
type UserServer struct {
userv1.UnimplementedUserServiceServer
db *sql.DB
}
func (s *UserServer) GetUser(ctx context.Context, req *userv1.GetUserRequest) (*userv1.GetUserResponse, error) {
var user userv1.User
err := s.db.QueryRowContext(ctx,
"SELECT id, name, email, status, created_at FROM users WHERE id = ?",
req.Id,
).Scan(&user.Id, &user.Name, &user.Email, &user.Status, &user.CreatedAt)
if err != nil {
return nil, status.Errorf(codes.NotFound, "user %d not found", req.Id)
}
return &userv1.GetUserResponse{User: &user}, nil
}
func (s *UserServer) ListUsers(req *userv1.ListUsersRequest, stream userv1.UserService_ListUsersServer) error {
rows, err := s.db.Query(
"SELECT id, name, email, status FROM users WHERE status = ? LIMIT ?",
req.StatusFilter, req.PageSize,
)
if err != nil {
return status.Errorf(codes.Internal, "query failed: %v", err)
}
defer rows.Close()
for rows.Next() {
var user userv1.User
if err := rows.Scan(&user.Id, &user.Name, &user.Email, &user.Status); err != nil {
return status.Errorf(codes.Internal, "scan failed: %v", err)
}
if err := stream.Send(&user); err != nil {
return err
}
}
return nil
}
启动gRPC服务器并注册服务:
func main() {
lis, err := net.Listen("tcp", ":50051")
if err != nil {
log.Fatalf("listen failed: %v", err)
}
grpcServer := grpc.NewServer(
grpc.UnaryInterceptor(unaryInterceptor),
grpc.StreamInterceptor(streamInterceptor),
)
userv1.RegisterUserServiceServer(grpcServer, &UserServer{db: db})
log.Println("gRPC server listening on :50051")
grpcServer.Serve(lis)
}
拦截器中间件与请求链设计
gRPC拦截器类似HTTP中间件,在请求处理前后插入横切逻辑。一元拦截器用于日志、认证、限流等:
func unaryInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
start := time.Now()
// 认证检查
if err := checkAuth(ctx); err != nil {
return nil, err
}
// 请求日志
log.Printf("gRPC call: %s, request: %+v", info.FullMethod, req)
// 执行handler
resp, err := handler(ctx, req)
// 响应日志与耗时
duration := time.Since(start)
if err != nil {
log.Printf("gRPC error: %s, duration: %v, error: %v", info.FullMethod, duration, err)
metrics.GRPCErrors.WithLabelValues(info.FullMethod, status.Code(err).String()).Inc()
} else {
log.Printf("gRPC done: %s, duration: %v", info.FullMethod, duration)
metrics.GRPCLatency.WithLabelValues(info.FullMethod).Observe(duration.Seconds())
}
return resp, err
}
流式拦截器处理流请求:
func streamInterceptor(srv interface{}, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
ctx := ss.Context()
if err := checkAuth(ctx); err != nil {
return err
}
wrappedStream := &serverStreamWrapper{
ServerStream: ss,
method: info.FullMethod,
}
return handler(srv, wrappedStream)
}
type serverStreamWrapper struct {
grpc.ServerStream
method string
}
func (w *serverStreamWrapper) SendMsg(m interface{}) error {
log.Printf("stream send: %s, msg: %+v", w.method, m)
return w.ServerStream.SendMsg(m)
}
func (w *serverStreamWrapper) RecvMsg(m interface{}) error {
err := w.ServerStream.RecvMsg(m)
if err == nil {
log.Printf("stream recv: %s, msg: %+v", w.method, m)
}
return err
}
多拦截器链式组合
实际项目需要多个拦截器协同工作,如认证-日志-限流-恢复。Go gRPC支持传入多个拦截器但需要手动链式包装,或使用go-grpc-middleware库:
import grpcMiddleware "github.com/grpc-ecosystem/go-grpc-middleware"
grpcServer := grpc.NewServer(
grpcMiddleware.WithUnaryServerChain(
recovery.UnaryServerInterceptor(), // panic恢复
logging.UnaryServerInterceptor(), // 请求日志
auth.UnaryServerInterceptor(), // JWT认证
ratelimit.UnaryServerInterceptor(), // 限流
),
grpcMiddleware.WithStreamServerChain(
recovery.StreamServerInterceptor(),
logging.StreamServerInterceptor(),
auth.StreamServerInterceptor(),
),
)
拦截器按声明顺序执行,每个拦截器调用handler前是请求预处理阶段,调用后是响应后处理阶段。执行顺序类似洋葱模型:外层拦截器先执行预处理,内层先执行后处理。
客户端调用与连接池管理
gRPC客户端使用grpc.Dial建立连接,底层复用HTTP/2连接,支持多路复用:
conn, err := grpc.Dial(
"localhost:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(`{
"loadBalancingPolicy": "round_robin",
"methodConfig": [{
"name": [{"service": "user.v1.UserService"}],
"retryPolicy": {
"maxAttempts": 3,
"initialBackoff": "0.1s",
"maxBackoff": "1s",
"retryableStatusCodes": ["UNAVAILABLE", "DEADLINE_EXCEEDED"]
}
}]
}`),
grpc.WithUnaryInterceptor(clientLogInterceptor),
)
if err != nil {
log.Fatalf("dial failed: %v", err)
}
defer conn.Close()
client := userv1.NewUserServiceClient(conn)
resp, err := client.GetUser(ctx, &userv1.GetUserRequest{Id: 1})
Service Config配置了负载均衡策略和重试策略。DNS解析多个地址时round_robin策略自动轮询。重试策略仅对幂等操作安全,需谨慎配置retryableStatusCodes。
连接池在gRPC中由底层HTTP/2连接自动管理。单个TCP连接支持100个并发流(HTTP/2默认),可通过grpc.WithMaxConcurrentStreams调整。高并发场景下建议保持连接复用,避免频繁建连。
错误处理与状态码规范
gRPC定义了标准状态码,服务端通过status.Errorf返回带状态码的错误:
switch {
case err == sql.ErrNoRows:
return nil, status.Error(codes.NotFound, "user not found")
case strings.Contains(err.Error(), "duplicate"):
return nil, status.Error(codes.AlreadyExists, "email already registered")
case strings.Contains(err.Error(), "connection refused"):
return nil, status.Error(codes.Unavailable, "database unavailable")
default:
return nil, status.Error(codes.Internal, err.Error())
}
客户端根据状态码执行重试或向上传递错误:status.Code(err)获取状态码,status.Convert(err).Message()获取错误消息。建议在API网关层将gRPC状态码映射为HTTP状态码,如NotFound对应404、AlreadyExists对应409、PermissionDenied对应403。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-wei-fu-wu-tong-xin-shi-zhan-protocolbuffers-ding-yi-yu/