gRPC服务开发实战:Protocol Buffers定义与跨语言调用实现

gRPC协议与微服务架构选型

gRPC是Google开源的高性能RPC框架,基于HTTP/2协议和Protocol Buffers序列化格式,在微服务架构中广泛应用于服务间通信。与REST API相比,gRPC的binary编码在传输体积上比JSON小3-10倍,HTTP/2多路复用避免了HTTP/1.1的连接数限制,端到端延迟通常降低30-50%。

在后端开发的API接口规范场景中,gRPC的四种调用模式覆盖了不同的通信需求:Unary RPC(一元调用,类似传统请求-响应)、Server Streaming(服务端流)、Client Streaming(客户端流)和Bidirectional Streaming(双向流)。服务治理场景下,gRPC原生支持拦截器、超时、重试和负载均衡。

Protocol Buffers消息定义与代码生成

Protocol Buffers(protobuf)是gRPC的接口定义语言(IDL)和序列化协议。通过.proto文件定义服务接口和消息结构,protoc编译器生成各语言的类型安全代码。

// user.proto - 用户服务定义
syntax = "proto3";

package user.v1;

option go_package = "github.com/yunthe/user-service/proto/user";
option java_package = "com.yunthe.user.proto";
option java_multiple_files = true;

// 用户服务定义
service UserService {
  // Unary: 创建用户
  rpc CreateUser(CreateUserRequest) returns (CreateUserResponse);
  
  // Server Streaming: 批量查询用户
  rpc ListUsers(ListUsersRequest) returns (stream User);
  
  // Client Streaming: 批量上传用户
  rpc UploadUsers(stream User) returns (UploadSummary);
  
  // Bidirectional Streaming: 实时聊天
  rpc ChatStream(stream ChatMessage) returns (stream ChatMessage);
}

// 消息定义
message User {
  int64 id = 1;
  string name = 2;
  string email = 3;
  UserStatus status = 4;
  google.protobuf.Timestamp created_at = 5;
}

enum UserStatus {
  USER_STATUS_UNSPECIFIED = 0;
  USER_STATUS_ACTIVE = 1;
  USER_STATUS_INACTIVE = 2;
  USER_STATUS_BANNED = 3;
}

message CreateUserRequest {
  string name = 1;
  string email = 2;
}

message CreateUserResponse {
  User user = 1;
}

message ListUsersRequest {
  int32 page_size = 1;
  string page_token = 2;
}

message UploadSummary {
  int32 total = 1;
  int32 success = 2;
  int32 failed = 3;
}

message ChatMessage {
  int64 user_id = 1;
  string content = 2;
  google.protobuf.Timestamp sent_at = 3;
}

使用protoc生成Go和Java代码:

# 安装protoc和插件
protoc --version  # 需3.15+

# Go插件
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest

# Java插件(Maven集成方式更常用)
# pom.xml中添加 protobuf-maven-plugin

# 生成Go代码
protoc --proto_path=.   --go_out=. --go_opt=paths=source_relative   --go-grpc_out=. --go-grpc_opt=paths=source_relative   user.proto

# 生成Java代码
protoc --proto_path=. --java_out=src/main/java user.proto

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/grpc/credentials"
    userv1 "github.com/yunthe/user-service/proto/user"
)

type UserServiceServer struct {
    userv1.UnimplementedUserServiceServer
    users map[int64]*userv1.User
}

// Unary RPC: 创建用户
func (s *UserServiceServer) CreateUser(ctx context.Context, req *userv1.CreateUserRequest) (*userv1.CreateUserResponse, error) {
    id := time.Now().Unix()
    user := &userv1.User{
        Id:    id,
        Name:  req.GetName(),
        Email: req.GetEmail(),
        Status: userv1.UserStatus_USER_STATUS_ACTIVE,
        CreatedAt: time.Now().UTC(),
    }
    s.users[id] = user
    return &userv1.CreateUserResponse{User: user}, nil
}

// Server Streaming: 批量查询
func (s *UserServiceServer) ListUsers(req *userv1.ListUsersRequest, stream userv1.UserService_ListUsersServer) error {
    count := 0
    for _, user := range s.users {
        if err := stream.Send(user); err != nil {
            return status.Errorf(codes.Internal, "发送失败: %v", err)
        }
        count++
        if count >= int(req.GetPageSize()) {
            break
        }
    }
    return nil
}

// Bidirectional Streaming: 实时聊天
func (s *UserServiceServer) ChatStream(stream userv1.UserService_ChatStreamServer) error {
    for {
        msg, err := stream.Recv()
        if err != nil {
            return err
        }
        // 处理消息并回复
        reply := &userv1.ChatMessage{
            UserId:   msg.GetUserId(),
            Content:  "收到: " + msg.GetContent(),
            SentAt:   time.Now().UTC(),
        }
        if err := stream.Send(reply); err != nil {
            return err
        }
    }
}

func main() {
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatalf("监听失败: %v", err)
    }
    
    // TLS加密
    creds, err := credentials.NewServerTLSFromFile("server.crt", "server.key")
    
    s := grpc.NewServer(
        grpc.Creds(creds),
        grpc.UnaryInterceptor(loggingInterceptor),
    )
    
    userv1.RegisterUserServiceServer(s, &UserServiceServer{
        users: make(map[int64]*userv1.User),
    })
    
    log.Println("gRPC服务启动在 :50051")
    if err := s.Serve(lis); err != nil {
        log.Fatalf("启动失败: %v", err)
    }
}

// 拦截器:日志记录
func loggingInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    start := time.Now()
    resp, err := handler(ctx, req)
    log.Printf("%s 耗时: %v 错误: %v", info.FullMethod, time.Since(start), err)
    return resp, err
}

Java客户端跨语言调用实现

package com.yunthe.user;

import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.TimeUnit;

public class UserClient {
    private final ManagedChannel channel;
    private final UserServiceGrpc.UserServiceBlockingStub blockingStub;
    private final UserServiceGrpc.UserServiceStub asyncStub;

    public UserClient(String host, int port) {
        this.channel = ManagedChannelBuilder.forAddress(host, port)
            .useTransportSecurity()  // TLS
            .build();
        this.blockingStub = UserServiceGrpc.newBlockingStub(channel);
        this.asyncStub = UserServiceGrpc.newStub(channel);
    }

    // Unary调用
    public void createUser(String name, String email) {
        CreateUserRequest req = CreateUserRequest.newBuilder()
            .setName(name)
            .setEmail(email)
            .build();
        CreateUserResponse resp = blockingStub.createUser(req);
        System.out.println("创建成功: " + resp.getUser().getName());
    }

    // Server Streaming调用
    public void listUsers(int pageSize) {
        ListUsersRequest req = ListUsersRequest.newBuilder()
            .setPageSize(pageSize)
            .build();
        blockingStub.listUsers(req).forEachRemaining(user -> {
            System.out.println("用户: " + user.getName() + " 状态: " + user.getStatus());
        });
    }

    // Bidirectional Streaming调用
    public void chat() {
        StreamObserver<ChatMessage> requestObserver = asyncStub.chatStream(
            new StreamObserver<ChatMessage>() {
                @Override
                public void onNext(ChatMessage value) {
                    System.out.println("收到: " + value.getContent());
                }
                @Override
                public void onError(Throwable t) {
                    t.printStackTrace();
                }
                @Override
                public void onCompleted() {
                    System.out.println("聊天结束");
                }
            }
        );
        // 发送消息
        requestObserver.onNext(ChatMessage.newBuilder()
            .setUserId(1)
            .setContent("Hello")
            .build());
        requestObserver.onCompleted();
    }

    public void shutdown() throws InterruptedException {
        channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);
    }

    public static void main(String[] args) throws InterruptedException {
        UserClient client = new UserClient("localhost", 50051);
        client.createUser("张三", "zhangsan@yunthe.com");
        client.listUsers(10);
        client.chat();
        client.shutdown();
    }
}

gRPC拦截器与服务治理

gRPC拦截器(Interceptor)类似中间件机制,可在请求处理前后插入横切逻辑,如认证、日志、限流、链路追踪。

// 客户端拦截器:自动重试
func retryInterceptor(ctx context.Context, method string, req, reply interface{}, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
    var lastErr error
    for i := 0; i < 3; i++ {
        ctx, cancel := context.WithTimeout(ctx, 5*time.Second)
        defer cancel()
        
        lastErr = invoker(ctx, method, req, reply, cc, opts...)
        if lastErr == nil {
            return nil
        }
        if status.Code(lastErr) == codes.Unavailable {
            time.Sleep(time.Duration(i+1) * time.Second)
            continue
        }
        return lastErr
    }
    return lastErr
}

// 服务端拦截器:JWT认证
func authInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    md, ok := metadata.FromIncomingContext(ctx)
    if !ok {
        return nil, status.Errorf(codes.Unauthenticated, "缺少metadata")
    }
    token := md.Get("authorization")
    if len(token) == 0 {
        return nil, status.Errorf(codes.Unauthenticated, "缺少token")
    }
    // 验证token
    if !validateToken(token[0]) {
        return nil, status.Errorf(codes.Unauthenticated, "token无效")
    }
    return handler(ctx, req)
}

// 客户端连接时传入拦截器
conn, _ := grpc.Dial("localhost:50051",
    grpc.WithTransportCredentials(creds),
    grpc.WithUnaryInterceptor(retryInterceptor),
)
defer conn.Close()

在高并发设计场景中,gRPC客户端使用连接池和keepalive机制维持长连接,减少握手开销。配合Envoy或Nginx作为gRPC代理,可实现负载均衡和流量分发。消息中间件如Kafka可与gRPC结合,实现事件驱动的异步通信模式。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/grpc-fu-wu-kai-fa-shi-zhan-protocolbuffers-ding-yi-yu-kua/

(0)
小编小编
上一篇 3小时前
下一篇 3小时前

相关推荐