gRPC 协议原理与性能优化

深入 gRPC 基于 HTTP/2 的实现原理,掌握 Protocol Buffers 序列化、四种通信模式与流控调优

gRPC 是 Google 开源的高性能 RPC 框架,基于 HTTP/2 和 Protocol Buffers,已成为微服务间通信的事实标准之一。

一、gRPC 架构概览

1.1 协议栈

┌─────────────────────────────┐
│ gRPC 层                      │
│ - 四种通信模式                │
│ - 拦截器、认证、压缩          │
├─────────────────────────────┤
│ HTTP/2 层                    │
│ - 头部帧、数据帧              │
│ - 流多路复用                  │
├─────────────────────────────┤
│ TLS 层(可选)                │
│ - 证书认证、ALPN              │
├─────────────────────────────┤
│ TCP 层                       │
│ - 可靠传输                    │
└─────────────────────────────┘

1.2 IDL 定义(.proto)

syntax = "proto3";

package order;
option java_package = "com.example.order";
option java_outer_classname = "OrderProto";

// 定义服务
service OrderService {
    // Unary:单次请求-响应
    rpc GetOrder (GetOrderRequest) returns (Order);
    
    // Server Streaming:服务端流
    rpc ListOrders (ListOrdersRequest) returns (stream Order);
    
    // Client Streaming:客户端流
    rpc CreateOrders (stream CreateOrderRequest) returns (OrderBatchResponse);
    
    // Bidirectional Streaming:双向流
    rpc SyncOrders (stream OrderEvent) returns (stream OrderStatus);
}

message GetOrderRequest {
    string order_id = 1;
}

message Order {
    string order_id = 1;
    string user_id = 2;
    double amount = 3;
    OrderStatus status = 4;
    int64 created_at = 5;
}

enum OrderStatus {
    PENDING = 0;
    PAID = 1;
    SHIPPED = 2;
    COMPLETED = 3;
}

message ListOrdersRequest {
    string user_id = 1;
    int32 page = 2;
    int32 page_size = 3;
}

message CreateOrderRequest {
    string user_id = 1;
    repeated OrderItem items = 2;
}

message OrderItem {
    string sku_id = 1;
    int32 quantity = 2;
    double price = 3;
}

message OrderBatchResponse {
    int32 created_count = 1;
    repeated string order_ids = 2;
}

message OrderEvent {
    string order_id = 1;
    OrderStatus new_status = 2;
}

二、Java gRPC 实战

2.1 Maven 配置

<properties>
    <grpc.version>1.65.1</grpc.version>
    <protobuf.version>3.25.3</protobuf.version>
</properties>

<dependencies>
    <dependency>
        <groupId>io.grpc</groupId>
        <artifactId>grpc-netty-shaded</artifactId>
        <version>${grpc.version}</version>
    </dependency>
    <dependency>
        <groupId>io.grpc</groupId>
        <artifactId>grpc-protobuf</artifactId>
        <version>${grpc.version}</version>
    </dependency>
    <dependency>
        <groupId>io.grpc</groupId>
        <artifactId>grpc-stub</artifactId>
        <version>${grpc.version}</version>
    </dependency>
</dependencies>

<build>
    <extensions>
        <extension>
            <groupId>kr.motd.maven</groupId>
            <artifactId>os-maven-plugin</artifactId>
            <version>1.7.1</version>
        </extension>
    </extensions>
    <plugins>
        <plugin>
            <groupId>org.xolstice.maven.plugins</groupId>
            <artifactId>protobuf-maven-plugin</artifactId>
            <version>0.6.1</version>
            <configuration>
                <protocArtifact>com.google.protobuf:protoc:${protobuf.version}:exe:${os.detected.classifier}</protocArtifact>
                <pluginId>grpc-java</pluginId>
                <pluginArtifact>io.grpc:protoc-gen-grpc-java:${grpc.version}:exe:${os.detected.classifier}</pluginArtifact>
            </configuration>
            <executions>
                <execution>
                    <goals>
                        <goal>compile</goal>
                        <goal>compile-custom</goal>
                    </goals>
                </execution>
            </executions>
        </plugin>
    </plugins>
</build>

2.2 服务端实现

@Service
public class OrderGrpcService extends OrderServiceGrpc.OrderServiceImplBase {
    
    @Autowired
    private OrderApplicationService orderService;
    
    // 1. Unary 模式
    @Override
    public void getOrder(GetOrderRequest request, StreamObserver<Order> responseObserver) {
        try {
            OrderDTO dto = orderService.getOrder(request.getOrderId());
            Order order = Order.newBuilder()
                .setOrderId(dto.getOrderId())
                .setUserId(dto.getUserId())
                .setAmount(dto.getAmount())
                .setStatus(mapStatus(dto.getStatus()))
                .setCreatedAt(dto.getCreatedAt().toEpochMilli())
                .build();
            
            responseObserver.onNext(order);
            responseObserver.onCompleted();
        } catch (OrderNotFoundException e) {
            responseObserver.onError(
                Status.NOT_FOUND.withDescription(e.getMessage()).asRuntimeException()
            );
        }
    }
    
    // 2. Server Streaming 模式
    @Override
    public void listOrders(ListOrdersRequest request, StreamObserver<Order> responseObserver) {
        Page<OrderDTO> orders = orderService.listOrders(
            request.getUserId(),
            request.getPage(),
            request.getPageSize()
        );
        
        for (OrderDTO dto : orders.getContent()) {
            responseObserver.onNext(mapToProto(dto));
        }
        responseObserver.onCompleted();
    }
    
    // 3. Client Streaming 模式
    @Override
    public StreamObserver<CreateOrderRequest> createOrders(
        StreamObserver<OrderBatchResponse> responseObserver
    ) {
        return new StreamObserver<>() {
            private final List<String> orderIds = new ArrayList<>();
            
            @Override
            public void onNext(CreateOrderRequest request) {
                String orderId = orderService.createOrder(request);
                orderIds.add(orderId);
            }
            
            @Override
            public void onError(Throwable t) {
                log.error("Client stream error", t);
            }
            
            @Override
            public void onCompleted() {
                responseObserver.onNext(OrderBatchResponse.newBuilder()
                    .setCreatedCount(orderIds.size())
                    .addAllOrderIds(orderIds)
                    .build());
                responseObserver.onCompleted();
            }
        };
    }
    
    // 4. Bidirectional Streaming 模式
    @Override
    public StreamObserver<OrderEvent> syncOrders(StreamObserver<OrderStatus> responseObserver) {
        return new StreamObserver<>() {
            @Override
            public void onNext(OrderEvent event) {
                // 处理客户端发来的事件
                OrderStatus status = orderService.processEvent(event);
                // 实时响应
                responseObserver.onNext(status);
            }
            
            @Override
            public void onError(Throwable t) {
                log.error("Bidirectional stream error", t);
            }
            
            @Override
            public void onCompleted() {
                responseObserver.onCompleted();
            }
        };
    }
}

// 服务端启动
@Configuration
public class GrpcServerConfig {
    
    @Autowired
    private OrderGrpcService orderService;
    
    @Bean(destroyMethod = "shutdown")
    public Server grpcServer() throws IOException {
        Server server = ServerBuilder.forPort(9090)
            .addService(orderService)
            .addService(ServerInterceptors.intercept(
                orderService,
                new AuthInterceptor(),
                new LoggingInterceptor()
            ))
            .maxInboundMessageSize(1024 * 1024 * 10)  // 10MB
            .maxInboundMetadataSize(1024 * 32)         // 32KB
            .build();
        
        server.start();
        return server;
    }
}

2.3 客户端实现

@Service
public class OrderGrpcClient {
    
    private final OrderServiceGrpc.OrderServiceBlockingStub blockingStub;
    private final OrderServiceGrpc.OrderServiceStub asyncStub;
    
    public OrderGrpcClient() {
        ManagedChannel channel = ManagedChannelBuilder
            .forAddress("localhost", 9090)
            .usePlaintext()  // 生产环境使用 TLS
            .maxRetryAttempts(3)
            .build();
        
        this.blockingStub = OrderServiceGrpc.newBlockingStub(channel);
        this.asyncStub = OrderServiceGrpc.newStub(channel);
    }
    
    // Unary 调用
    public Order getOrder(String orderId) {
        GetOrderRequest request = GetOrderRequest.newBuilder()
            .setOrderId(orderId)
            .build();
        
        return blockingStub.getOrder(request);
    }
    
    // Server Streaming
    public List<Order> listOrders(String userId) {
        ListOrdersRequest request = ListOrdersRequest.newBuilder()
            .setUserId(userId)
            .setPage(0)
            .setPageSize(100)
            .build();
        
        Iterator<Order> iterator = blockingStub.listOrders(request);
        List<Order> orders = new ArrayList<>();
        while (iterator.hasNext()) {
            orders.add(iterator.next());
        }
        return orders;
    }
    
    // Client Streaming(异步)
    public CompletableFuture<OrderBatchResponse> createOrders(List<CreateOrderRequest> requests) {
        CompletableFuture<OrderBatchResponse> future = new CompletableFuture<>();
        
        StreamObserver<OrderBatchResponse> responseObserver = new StreamObserver<>() {
            @Override
            public void onNext(OrderBatchResponse response) {
                future.complete(response);
            }
            
            @Override
            public void onError(Throwable t) {
                future.completeExceptionally(t);
            }
            
            @Override
            public void onCompleted() {}
        };
        
        StreamObserver<CreateOrderRequest> requestObserver = asyncStub.createOrders(responseObserver);
        
        for (CreateOrderRequest req : requests) {
            requestObserver.onNext(req);
        }
        requestObserver.onCompleted();
        
        return future;
    }
}

三、gRPC 与 HTTP/2 帧映射

gRPC 消息如何映射到 HTTP/2:

请求:
HEADERS 帧
  :method = POST
  :scheme = http
  :authority = localhost:9090
  :path = /order.OrderService/GetOrder
  content-type = application/grpc
  te = trailers
  grpc-timeout = 10S
  grpc-encoding = gzip

DATA 帧(Length-Prefixed Message)
  +------------------+------------------+
  |  Compressed Flag | Message Length   |  (5 bytes 前缀)
  |  (1 byte)        | (4 bytes)        |
  +------------------+------------------+
  |          Serialized Protobuf         |
  +--------------------------------------+

响应:
HEADERS 帧(初始)
  :status = 200
  content-type = application/grpc

DATA 帧(响应消息)

HEADERS 帧(结束,Trailers)
  grpc-status = 0
  grpc-message = OK

四、性能优化

4.1 连接池与通道复用

@Configuration
public class GrpcChannelConfig {
    
    @Bean
    public Map<String, ManagedChannel> grpcChannels() {
        return Map.of(
            "order-service", createChannel("order-svc", 9090),
            "inventory-service", createChannel("inventory-svc", 9091),
            "payment-service", createChannel("payment-svc", 9092)
        );
    }
    
    private ManagedChannel createChannel(String host, int port) {
        return ManagedChannelBuilder.forAddress(host, port)
            .useTransportSecurity()  // TLS
            .defaultLoadBalancingPolicy("round_robin")
            .enableRetry()
            .maxRetryAttempts(3)
            .maxHedgedAttempts(1)
            .keepAliveTime(60, TimeUnit.SECONDS)
            .keepAliveTimeout(20, TimeUnit.SECONDS)
            .keepAliveWithoutCalls(true)
            .build();
    }
}

4.2 拦截器

// 认证拦截器
public class AuthInterceptor implements ClientInterceptor {
    
    private final String token;
    
    @Override
    public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
        MethodDescriptor<ReqT, RespT> method,
        CallOptions callOptions,
        Channel next
    ) {
        return new ForwardingClientCall.SimpleForwardingClientCall<>(next.newCall(method, callOptions)) {
            @Override
            public void start(Listener<RespT> responseListener, Metadata headers) {
                headers.put(Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER), "Bearer " + token);
                super.start(responseListener, headers);
            }
        };
    }
}

// 超时拦截器
public class TimeoutInterceptor implements ClientInterceptor {
    
    @Override
    public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
        MethodDescriptor<ReqT, RespT> method,
        CallOptions callOptions,
        Channel next
    ) {
        Deadline deadline = Deadline.after(10, TimeUnit.SECONDS);
        return next.newCall(method, callOptions.withDeadline(deadline));
    }
}

五、总结

方面要点
通信模式Unary、Server Stream、Client Stream、Bidirectional
序列化Protocol Buffers,高效紧凑
传输层HTTP/2,流多路复用
优势高性能、强类型、流支持、跨语言
劣势浏览器支持有限(需 gRPC-Web)、调试复杂
优化手段效果
连接复用减少 TCP/TLS 握手开销
消息复用通过流模式减少往返
压缩gzip/deflate 减少传输量
拦截器统一处理认证、日志、超时
负载均衡round_robin / least_request

gRPC 的设计充分利用了 HTTP/2 的多路复用和流特性,配合 Protocol Buffers 的高效序列化,为微服务间通信提供了企业级的解决方案。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「network」更多文章

  1. 网络安全:TLS/SSL、证书与加密通信
  2. 负载均衡算法与高可用架构
  3. DNS 系统与智能解析