gRPC是Google开源的高性能RPC框架,基于HTTP/2协议传输,使用Protocol Buffers作为接口定义语言和序列化格式。相比REST+JSON,gRPC在吞吐量、延迟和类型安全方面具有显著优势,二进制序列化体积约为JSON的1/3到1/10,HTTP/2多路复用避免队头阻塞,适合微服务间高并发内部通信。gRPC支持四种调用模式:Unary单向调用、Server Streaming服务端流、Client Streaming客户端流、Bidirectional Streaming双向流,覆盖同步请求、日志推送、文件上传、实时通信等场景。本文从Protobuf定义、服务实现、拦截器、负载均衡四个层面展开实战。
Protobuf接口定义与代码生成
Protobuf(Protocol Buffers)是gRPC的接口定义语言(IDL),通过.proto文件定义服务接口和消息结构。protoc编译器将.proto文件编译为各语言的源代码,保证跨语言类型安全。
// proto/user_service.proto
syntax = "proto3";
package user.v1;
option go_package = "github.com/example/proto/user/v1;userv1";
option java_package = "com.example.proto.user.v1";
option java_multiple_files = true;
import "google/protobuf/timestamp.proto";
import "google/protobuf/empty.proto";
// 用户服务定义
service UserService {
// Unary: 创建用户
rpc CreateUser(CreateUserRequest) returns (CreateUserResponse);
// Unary: 查询用户
rpc GetUser(GetUserRequest) returns (GetUserResponse);
// Unary: 更新用户
rpc UpdateUser(UpdateUserRequest) returns (UpdateUserResponse);
// Server Streaming: 批量查询用户(服务端流式返回)
rpc ListUsers(ListUsersRequest) returns (stream User);
// Client Streaming: 批量创建用户(客户端流式发送)
rpc BatchCreateUsers(stream CreateUserRequest) returns (BatchCreateResponse);
// Bidirectional Streaming: 实时用户状态变更通知
rpc WatchUserStatus(stream WatchRequest) returns (stream UserStatusEvent);
}
// 消息定义
message User {
int64 id = 1;
string username = 2;
string email = 3;
string phone = 4;
UserStatus status = 5;
google.protobuf.Timestamp created_at = 6;
google.protobuf.Timestamp updated_at = 7;
repeated string roles = 8; // repeated = 数组
map<string, string> metadata = 9; // map类型
}
enum UserStatus {
USER_STATUS_UNSPECIFIED = 0;
USER_STATUS_ACTIVE = 1;
USER_STATUS_INACTIVE = 2;
USER_STATUS_SUSPENDED = 3;
}
message CreateUserRequest {
string username = 1;
string email = 2;
string phone = 3;
string password = 4;
}
message CreateUserResponse {
User user = 1;
}
message GetUserRequest {
int64 id = 1;
}
message GetUserResponse {
User user = 1;
}
message UpdateUserRequest {
int64 id = 1;
optional string email = 2; // optional字段
optional string phone = 3;
optional UserStatus status = 4;
}
message UpdateUserResponse {
User user = 1;
}
message ListUsersRequest {
int32 page_size = 1;
string page_token = 2;
UserStatus status_filter = 3;
}
message BatchCreateResponse {
int32 success_count = 1;
int32 failure_count = 2;
repeated string errors = 3;
}
message WatchRequest {
int64 user_id = 1;
}
message UserStatusEvent {
int64 user_id = 1;
UserStatus old_status = 2;
UserStatus new_status = 3;
google.protobuf.Timestamp event_time = 4;
}
# 生成Go代码
protoc --go_out=. --go_opt=paths=source_relative --go-grpc_out=. --go-grpc_opt=paths=source_relative proto/user_service.proto
# 生成Java代码
protoc --java_out=src/main/java --grpc-java_out=src/main/java proto/user_service.proto
# 使用buf工具管理proto(推荐)
# buf.yaml
version: v1
breaking:
use:
- FILE
lint:
use:
- DEFAULT
# buf.gen.yaml
version: v1
plugins:
- plugin: go
out: gen/go
opt: paths=source_relative
- plugin: go-grpc
out: gen/go
opt: paths=source_relative
# 执行代码生成
buf generate
gRPC服务端实现与流式调用处理
Go语言实现gRPC服务端,需注册服务并实现.proto中定义的所有RPC方法。流式RPC通过流对象(stream)进行消息发送和接收。
// server/main.go - gRPC服务端实现
package main
import (
"context"
"log"
"net"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/types/known/timestamppb"
pb "github.com/example/proto/user/v1"
)
type userServiceServer struct {
pb.UnimplementedUserServiceServer
users map[int64]*pb.User
nextID int64
}
func NewUserServiceServer() *userServiceServer {
return &userServiceServer{
users: make(map[int64]*pb.User),
nextID: 1,
}
}
// Unary RPC: 创建用户
func (s *userServiceServer) CreateUser(ctx context.Context, req *pb.CreateUserRequest) (*pb.CreateUserResponse, error) {
// 参数校验
if req.GetUsername() == "" {
return nil, status.Error(codes.InvalidArgument, "username is required")
}
if req.GetEmail() == "" {
return nil, status.Error(codes.InvalidArgument, "email is required")
}
user := &pb.User{
Id: s.nextID,
Username: req.GetUsername(),
Email: req.GetEmail(),
Phone: req.GetPhone(),
Status: pb.UserStatus_USER_STATUS_ACTIVE,
CreatedAt: timestamppb.Now(),
UpdatedAt: timestamppb.Now(),
Roles: []string{"user"},
}
s.users[s.nextID] = user
s.nextID++
return &pb.CreateUserResponse{User: user}, nil
}
// Server Streaming: 批量查询用户
func (s *userServiceServer) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
count := 0
for _, user := range s.users {
// 根据状态过滤
if req.GetStatusFilter() != pb.UserStatus_USER_STATUS_UNSPECIFIED &&
user.GetStatus() != req.GetStatusFilter() {
continue
}
// 逐条发送,模拟分页
if err := stream.Send(user); err != nil {
return status.Errorf(codes.Internal, "send error: %v", err)
}
count++
// 限制返回数量
if req.GetPageSize() > 0 && count >= int(req.GetPageSize()) {
break
}
// 模拟延迟
time.Sleep(10 * time.Millisecond)
}
return nil
}
// Bidirectional Streaming: 实时用户状态监控
func (s *userServiceServer) WatchUserStatus(stream pb.UserService_WatchUserStatusServer) error {
// 接收客户端订阅请求,推送状态变更事件
for {
req, err := stream.Recv()
if err != nil {
return err
}
// 模拟状态变更事件推送
event := &pb.UserStatusEvent{
UserId: req.GetUserId(),
NewStatus: pb.UserStatus_USER_STATUS_ACTIVE,
EventTime: timestamppb.Now(),
}
if err := stream.Send(event); err != nil {
return err
}
}
}
func main() {
lis, err := net.Listen("tcp", ":50051")
if err != nil {
log.Fatalf("failed to listen: %v", err)
}
grpcServer := grpc.NewServer()
pb.RegisterUserServiceServer(grpcServer, NewUserServiceServer())
log.Println("gRPC server listening on :50051")
if err := grpcServer.Serve(lis); err != nil {
log.Fatalf("failed to serve: %v", err)
}
}
gRPC拦截器与中间件链式调用
拦截器(Interceptor)是gRPC的AOP机制,类似Web框架中的中间件。通过拦截器统一处理日志记录、认证鉴权、链路追踪、指标采集、限流熔断等横切关注点。
// interceptor/interceptors.go - gRPC拦截器
package interceptor
import (
"context"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/metadata"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)
// Unary服务端拦截器:日志记录
func LoggingUnaryInterceptor(
ctx context.Context, req interface{},
info *grpc.UnaryServerInfo, handler grpc.UnaryHandler,
) (resp interface{}, err error) {
start := time.Now()
resp, err = handler(ctx, req)
duration := time.Since(start)
code := status.Code(err)
log.Printf("method=%s duration=%v code=%s req_size=%d",
info.FullMethod, duration, code, protoSize(req))
return resp, err
}
// Unary服务端拦截器:认证鉴权
func AuthUnaryInterceptor(
ctx context.Context, req interface{},
info *grpc.UnaryServerInfo, handler grpc.UnaryHandler,
) (resp interface{}, err error) {
// 从metadata中提取token
md, ok := metadata.FromIncomingContext(ctx)
if !ok {
return nil, status.Error(codes.Unauthenticated, "missing metadata")
}
tokens := md.Get("authorization")
if len(tokens) == 0 {
return nil, status.Error(codes.Unauthenticated, "missing auth token")
}
// 验证token
userID, err := validateToken(tokens[0])
if err != nil {
return nil, status.Error(codes.Unauthenticated, "invalid token")
}
// 将用户信息注入context
ctx = context.WithValue(ctx, "userID", userID)
return handler(ctx, req)
}
// 限流拦截器(令牌桶)
func RateLimitUnaryInterceptor(
ctx context.Context, req interface{},
info *grpc.UnaryServerInfo, handler grpc.UnaryHandler,
) (resp interface{}, err error) {
if !rateLimiter.Allow() {
return nil, status.Error(codes.ResourceExhausted, "rate limit exceeded")
}
return handler(ctx, req)
}
// 注册拦截器链
func main() {
grpcServer := grpc.NewServer(
grpc.ChainUnaryInterceptor(
LoggingUnaryInterceptor,
AuthUnaryInterceptor,
RateLimitUnaryInterceptor,
),
grpc.ChainStreamInterceptor(
LoggingStreamInterceptor,
AuthStreamInterceptor,
),
)
// ...
}
// Stream拦截器
func LoggingStreamInterceptor(
srv interface{}, ss grpc.ServerStream,
info *grpc.StreamServerInfo, handler grpc.StreamHandler,
) error {
start := time.Now()
err := handler(srv, ss)
log.Printf("stream method=%s duration=%v err=%v",
info.FullMethod, time.Since(start), err)
return err
}
gRPC客户端负载均衡与服务发现集成
gRPC客户端内置负载均衡支持,支持round-robin、pick-first等策略。配合服务发现组件(Consul、etcd、Nacos)实现动态服务地址更新,无需依赖外部负载均衡器。
// client/main.go - gRPC客户端负载均衡
package main
import (
"context"
"log"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/balancer/roundrobin"
_ "google.golang.org/grpc/balancer/roundrobin"
pb "github.com/example/proto/user/v1"
)
func main() {
// 方式1:静态地址 + round-robin负载均衡
conn, err := grpc.Dial(
"dns:///user-service:50051", // DNS解析多个A记录
grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
if err != nil {
log.Fatalf("did not connect: %v", err)
}
defer conn.Close()
client := pb.NewUserServiceClient(conn)
// Unary调用(带超时)
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
resp, err := client.CreateUser(ctx, &pb.CreateUserRequest{
Username: "testuser",
Email: "test@example.com",
})
if err != nil {
log.Fatalf("create user failed: %v", err)
}
log.Printf("created user: %v", resp.GetUser())
// 方式2:自定义Resolver集成Consul服务发现
// grpc.WithResolvers(consul.NewResolver()) // 注册Consul resolver
// conn, _ = grpc.Dial(
// "consul://127.0.0.1:8500/user-service",
// grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`),
// )
// 方式3:客户端流式调用
stream, err := client.ListUsers(ctx, &pb.ListUsersRequest{PageSize: 10})
if err != nil {
log.Fatalf("list users failed: %v", err)
}
for {
user, err := stream.Recv()
if err != nil {
break
}
log.Printf("user: id=%d username=%s", user.GetId(), user.GetUsername())
}
}
生产环境中gRPC服务部署在Kubernetes时,可直接利用Service的DNS解析实现客户端负载均衡,每个Pod实例对应一个DNS A记录。对于跨集群通信,可引入Envoy或Linkerd作为L7代理,在代理层实现更灵活的流量管理策略。gRPC的健康检查通过gRPC Health Checking Protocol实现,Kubernetes中可使用grpc-health-probe替代TCP探针,准确反映服务就绪状态。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-wei-fu-wu-tong-xin-kuang-jia-yu-protobuf-jie-kou-ding/