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 的高效序列化,为微服务间通信提供了企业级的解决方案。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。