Java 基础体系 · 第 96/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java gRPC:Protobuf、Unary、Stream、拦截器、Deadline 和治理
gRPC 是一套基于 HTTP/2 的 RPC(Remote Procedure Call,远程过程调用)框架。客户端调用远程服务时,代码形态可以接近本地方法调用:
OrderReply reply = blockingStub.getOrder(request);
但这行代码背后包含了序列化、HTTP/2 帧传输、连接复用、流控、超时、取消、状态码、拦截器和服务治理。理解 gRPC,不能只记住“用 .proto 生成 Java 类”,还必须理解一次 RPC 的数据流和失败边界。
本文以 Java 25 LTS 为编译与运行环境,使用 grpc-java 的 Java API。Java 25 是应用运行时版本;gRPC 的具体 API 仍由 grpc-java 版本决定,因此工程中必须锁定并验证 grpc-java、protobuf-java、protoc 和 protobuf Maven 插件的兼容版本。下面示例使用一组可复现的 Maven 版本,升级时应在 CI 中重新执行生成、编译和集成测试。
一、先建立模型:一次 gRPC 调用到底包含什么
一次 gRPC 调用可以抽象为:
其中:
- Method:服务名和方法名,例如
/wr.order.OrderService/GetOrder。 - Metadata:请求头和响应头,通常承载认证令牌、链路追踪标识、租户信息等。
- Message Stream:一个或多个 Protobuf 消息。
- Status:最终状态,例如
OK、NOT_FOUND、DEADLINE_EXCEEDED。
“Stream”在这里有两层含义:
- HTTP/2 本身通过一个 stream 承载一次 RPC。
- gRPC 方法还可以定义连续发送多个 Protobuf 消息的流式语义。
因此,Unary RPC 也使用 HTTP/2 stream,只是应用层消息数量通常是一请求一响应;Server Streaming 则是在同一个 HTTP/2 stream 中返回多个响应消息。
一次调用的典型顺序如下:
sequenceDiagram
participant C as Java 客户端
participant I1 as ClientInterceptor
participant H as HTTP/2 Channel
participant SI as ServerInterceptor
participant S as ServiceImpl
C->>I1: 调用 Stub
I1->>H: Metadata + Method + Deadline
H->>SI: 创建服务端调用
SI->>S: onMessage(request)
S-->>SI: onNext(response)
SI-->>H: 响应消息 + trailers
H-->>C: response 或 Status
关键点是:业务方法并不直接操作 TCP 连接。客户端 Stub、Channel、拦截器、HTTP/2 传输层和服务实现之间由框架连接起来。
二、Protobuf:不仅是“序列化格式”
2.1 .proto 是服务契约
Protobuf(Protocol Buffers)承担两项职责:
- 描述消息结构;
- 描述 RPC 服务和方法。
下面定义一个订单服务:
syntax = "proto3";
package wr.order.v1;
option java_multiple_files = true;
option java_package = "com.example.order.proto";
option java_outer_classname = "OrderProto";
service OrderService {
rpc GetOrder(GetOrderRequest) returns (OrderReply);
rpc WatchOrder(WatchOrderRequest) returns (stream OrderEvent);
rpc UploadOrderEvents(stream OrderEvent) returns (UploadSummary);
rpc Chat(stream OrderEvent) returns (stream OrderEvent);
}
message GetOrderRequest {
string order_id = 1;
}
message WatchOrderRequest {
string order_id = 1;
}
message OrderReply {
string order_id = 1;
string status = 2;
int64 version = 3;
}
message OrderEvent {
string order_id = 1;
int64 version = 2;
string type = 3;
}
message UploadSummary {
int32 accepted = 1;
}
字段后的数字不是注释,而是 field number。Protobuf 二进制编码使用字段编号识别字段,而不是使用字段名。因此:
string order_id = 1;
序列化时会编码字段编号 1 及其值。字段名可以修改,但字段编号一旦发布就不能随意复用。
2.2 字段演进的形式化约束
假设旧版本消息定义为:
message User {
string name = 1;
int32 age = 2;
}
新版本希望增加邮箱:
message User {
string name = 1;
int32 age = 2;
string email = 3;
}
旧客户端读取新消息时,不认识字段 3,会将其视为未知字段并忽略;新客户端读取旧消息时,字段 3 不存在,得到默认值。这构成常见的前向和后向兼容基础。
但以下修改破坏兼容性:
// 错误:把原来的字段 1 改成了其他语义
int64 name = 1;
// 错误:删除字段后重新复用 2
string email = 2;
删除字段时应保留编号:
message User {
reserved 2;
reserved "age";
string name = 1;
}
reserved 的作用是防止后续开发者重新使用已经删除的字段编号或名称。
兼容性可以用一条约束表示:
字段编号的语义必须保持稳定。改变字段类型也不能只看 Java 类型是否“能转换”,还要看 Protobuf wire type 是否兼容。例如某些整数类型之间在 wire level 上可兼容,但可能发生溢出或语义变化;string 与 bytes 也不能因为都表现为字节序列就任意替换。
2.3 proto3 默认值与 presence
在普通 proto3 标量字段中,未设置的 string 通常读取为空字符串,int32 读取 0,bool 读取 false。这意味着“没有传值”和“显式传了默认值”可能无法区分。
需要区分 presence 时,可以使用 optional:
message UpdateOrderRequest {
string order_id = 1;
optional string remark = 2;
}
生成的 Java API 会提供 presence 判断,例如:
UpdateOrderRequest request = UpdateOrderRequest.newBuilder()
.setOrderId("o-100")
.build();
if (request.hasRemark()) {
System.out.println(request.getRemark());
}
如果业务语义中“清空字段”和“不修改字段”不同,不能仅依赖普通标量默认值。
2.4 Maven 生成 Java 类和 gRPC Stub
一个最小 Maven 配置如下:
<project>
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>grpc-order-demo</artifactId>
<version>1.0.0</version>
<properties>
<maven.compiler.release>25</maven.compiler.release>
<grpc.version>1.66.0</grpc.version>
<protobuf.version>3.25.3</protobuf.version>
<protobuf.plugin.version>0.6.1</protobuf.plugin.version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-bom</artifactId>
<version>${grpc.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-netty-shaded</artifactId>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-protobuf</artifactId>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-stub</artifactId>
</dependency>
<dependency>
<groupId>org.apache.tomcat</groupId>
<artifactId>annotations-api</artifactId>
<version>6.0.53</version>
<scope>provided</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.xolstice.maven.plugins</groupId>
<artifactId>protobuf-maven-plugin</artifactId>
<version>${protobuf.plugin.version}</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>
<plugin>
<groupId>kr.motd.maven</groupId>
<artifactId>os-maven-plugin</artifactId>
<version>1.7.1</version>
<extensions>true</extensions>
</plugin>
<plugin>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.13.0</version>
<configuration>
<release>25</release>
</configuration>
</plugin>
</plugins>
</build>
</project>
将 .proto 文件放在:
src/main/proto/order.proto
然后执行:
mvn clean compile
预期结果是生成两类代码:
OrderServiceGrpc:包含服务基类、Stub 和方法描述;GetOrderRequest、OrderReply等:由 Protobuf 生成的不可变消息类。
如果只执行 Protobuf 编译而没有执行 compile-custom,消息类可能生成了,但 gRPC Service 和 Stub 不会生成,最终会出现 cannot find symbol: OrderServiceGrpc。
三、四种 RPC 形态:消息数量与方向决定 API
在 .proto 中,stream 修饰请求或响应:
| 类型 | 请求 | 响应 | 常见 Java Stub |
|---|---|---|---|
| Unary | 1 | 1 | blockingStub、异步 Stub |
| Server Streaming | 1 | 多 | blockingStub 返回 Iterator |
| Client Streaming | 多 | 1 | 异步 Stub 接收 StreamObserver |
| Bidirectional Streaming | 多 | 多 | 异步 Stub 接收 StreamObserver |
这四种形态不是性能开关,而是协议语义。把一个本应连续传输的业务过程硬塞进多个 Unary 调用,会丢失流的顺序、背压和取消语义;把简单查询设计成双向流,则会增加状态管理复杂度。
3.1 Unary:一个请求对应一个响应
生成的服务端基类大致包含:
public abstract static class OrderServiceImplBase
implements BindableService {
public void getOrder(
GetOrderRequest request,
StreamObserver<OrderReply> responseObserver) {
...
}
}
实现服务:
import com.example.order.proto.GetOrderRequest;
import com.example.order.proto.OrderReply;
import com.example.order.proto.OrderServiceGrpc;
import io.grpc.Status;
import io.grpc.stub.StreamObserver;
public final class OrderService
extends OrderServiceGrpc.OrderServiceImplBase {
@Override
public void getOrder(
GetOrderRequest request,
StreamObserver<OrderReply> responseObserver) {
if (request.getOrderId().isBlank()) {
responseObserver.onError(
Status.INVALID_ARGUMENT
.withDescription("order_id must not be blank")
.asRuntimeException());
return;
}
if (request.getOrderId().equals("missing")) {
responseObserver.onError(
Status.NOT_FOUND
.withDescription("order not found")
.asRuntimeException());
return;
}
OrderReply reply = OrderReply.newBuilder()
.setOrderId(request.getOrderId())
.setStatus("PAID")
.setVersion(7)
.build();
responseObserver.onNext(reply);
responseObserver.onCompleted();
}
}
服务端必须遵守响应观察者的终止规则:
- 成功路径:
onNext零次或多次,然后onCompleted; - 失败路径:调用一次
onError; onCompleted和onError不能同时调用;- 终止后不能继续发送消息。
对 Unary 来说,通常只发送一次 onNext。如果发送两次,客户端可能得到协议错误或未定义的业务结果,因此不要把 StreamObserver 当作普通集合输出接口。
客户端可以使用阻塞 Stub:
OrderReply reply = blockingStub
.withDeadlineAfter(800, java.util.concurrent.TimeUnit.MILLISECONDS)
.getOrder(
GetOrderRequest.newBuilder()
.setOrderId("o-100")
.build());
异常处理:
try {
OrderReply reply = blockingStub
.withDeadlineAfter(800, TimeUnit.MILLISECONDS)
.getOrder(request);
System.out.println(reply.getStatus());
} catch (io.grpc.StatusRuntimeException e) {
io.grpc.Status.Code code = e.getStatus().getCode();
if (code == io.grpc.Status.Code.NOT_FOUND) {
// 业务上的“不存在”
} else if (code == io.grpc.Status.Code.DEADLINE_EXCEEDED) {
// 本次调用没有在预算内完成
} else {
// 记录 status、description 和 cause,交给统一错误处理
}
}
StatusRuntimeException 表示 RPC 失败,不应直接把它转换成 HTTP 200 或普通字符串后丢失状态码语义。
3.2 Server Streaming:一次请求,服务端连续返回多个消息
服务端实现:
@Override
public void watchOrder(
WatchOrderRequest request,
StreamObserver<OrderEvent> responseObserver) {
try {
for (int version = 1; version <= 3; version++) {
if (io.grpc.Context.current().isCancelled()) {
return;
}
responseObserver.onNext(OrderEvent.newBuilder()
.setOrderId(request.getOrderId())
.setVersion(version)
.setType("UPDATED")
.build());
}
responseObserver.onCompleted();
} catch (RuntimeException e) {
responseObserver.onError(
Status.INTERNAL
.withDescription("watch failed")
.withCause(e)
.asRuntimeException());
}
}
阻塞 Stub 的调用结果是 Iterator:
Iterator<OrderEvent> events = blockingStub
.withDeadlineAfter(2, TimeUnit.SECONDS)
.watchOrder(
WatchOrderRequest.newBuilder()
.setOrderId("o-100")
.build());
while (events.hasNext()) {
OrderEvent event = events.next();
System.out.printf("%s: version=%d%n",
event.getType(), event.getVersion());
}
hasNext() 可能阻塞,且在服务端返回错误时可能抛出 StatusRuntimeException。所以不能只在创建 Iterator 的地方捕获异常:
try {
Iterator<OrderEvent> events = blockingStub.watchOrder(request);
while (events.hasNext()) {
consume(events.next());
}
} catch (StatusRuntimeException e) {
handleRpcFailure(e);
}
Server Streaming 适合事件订阅、分页结果、批量导出等场景,但“服务端持续发送”不等于“永远不设置 Deadline”。长连接应该明确生命周期、取消方式、重连方式和最大持续时间,否则连接泄漏会表现为服务端线程、内存或订阅对象持续增长。
3.3 Client Streaming:多个请求消息汇聚为一个响应
服务端基类方法形态是:
public StreamObserver<OrderEvent> uploadOrderEvents(
StreamObserver<UploadSummary> responseObserver)
实现示例:
@Override
public StreamObserver<OrderEvent> uploadOrderEvents(
StreamObserver<UploadSummary> responseObserver) {
return new StreamObserver<>() {
private int accepted;
@Override
public void onNext(OrderEvent event) {
if (!event.getOrderId().isBlank()) {
accepted++;
}
}
@Override
public void onError(Throwable throwable) {
// 客户端取消、网络断开或客户端发送失败
System.err.println("upload cancelled: " + throwable);
}
@Override
public void onCompleted() {
responseObserver.onNext(
UploadSummary.newBuilder()
.setAccepted(accepted)
.build());
responseObserver.onCompleted();
}
};
}
客户端使用异步 Stub:
CountDownLatch done = new CountDownLatch(1);
StreamObserver<UploadSummary> responseObserver =
new StreamObserver<>() {
@Override
public void onNext(UploadSummary summary) {
System.out.println("accepted = " + summary.getAccepted());
}
@Override
public void onError(Throwable t) {
handleRpcFailure(t);
done.countDown();
}
@Override
public void onCompleted() {
done.countDown();
}
};
StreamObserver<OrderEvent> requestObserver =
asyncStub
.withDeadlineAfter(3, TimeUnit.SECONDS)
.uploadOrderEvents(responseObserver);
try {
requestObserver.onNext(event("o-100", 1));
requestObserver.onNext(event("o-100", 2));
requestObserver.onCompleted();
if (!done.await(5, TimeUnit.SECONDS)) {
throw new IllegalStateException("client wait timeout");
}
} catch (RuntimeException e) {
requestObserver.onError(e);
}
这里的 onCompleted() 表示“请求方向发送结束”,并不等于整个 RPC 已成功;服务端仍可能返回错误。只有响应观察者收到 onCompleted(),客户端才知道响应方向正常结束。
3.4 Bidirectional Streaming:两个方向独立推进
双向流中,客户端和服务端都可以在对方完成之前发送消息:
@Override
public StreamObserver<OrderEvent> chat(
StreamObserver<OrderEvent> responseObserver) {
return new StreamObserver<>() {
@Override
public void onNext(OrderEvent event) {
responseObserver.onNext(
OrderEvent.newBuilder()
.setOrderId(event.getOrderId())
.setVersion(event.getVersion())
.setType("ACK")
.build());
}
@Override
public void onError(Throwable t) {
// 对端取消或连接失败
}
@Override
public void onCompleted() {
responseObserver.onCompleted();
}
};
}
双向流的难点不是 API,而是并发和生命周期:
onNext回调通常由 gRPC 的执行线程调用;- 不能在回调中执行长时间阻塞的数据库或 HTTP 调用;
- 发送方必须考虑接收方速度,不能无限制地把消息堆在内存中;
- 关闭一侧不一定意味着另一侧立即发送完;
- 任一侧取消后,另一侧继续发送最终都会失败。
如果业务需要严格的应用层背压、窗口确认或消息重放,必须在 Protobuf 消息中设计序号、确认号和幂等键,不能假设 TCP 或 HTTP/2 会替业务自动完成这些事情。
四、Channel、Stub 和服务端生命周期
4.1 Channel 是长期资源,Stub 是轻量视图
客户端通常创建一个长期复用的 ManagedChannel:
ManagedChannel channel = ManagedChannelBuilder
.forAddress("localhost", 50051)
.usePlaintext()
.build();
OrderServiceGrpc.OrderServiceBlockingStub blockingStub =
OrderServiceGrpc.newBlockingStub(channel);
然后在调用时通过 withDeadlineAfter、withInterceptors 等方法派生新的 Stub。Stub 本身通常是不可变的轻量对象,可以按请求派生;Channel 则持有连接、线程和名称解析等资源,不能每次请求都创建。
关闭过程:
channel.shutdown();
if (!channel.awaitTermination(5, TimeUnit.SECONDS)) {
channel.shutdownNow();
}
usePlaintext() 只适合本地开发或受控测试。生产环境应配置 TLS,并根据部署环境配置服务器证书校验、客户端证书或其他认证机制。
4.2 一个可运行的服务端和客户端
服务端:
import io.grpc.Server;
import io.grpc.ServerBuilder;
public final class Main {
public static void main(String[] args) throws Exception {
Server server = ServerBuilder
.forPort(50051)
.addService(new OrderService())
.build()
.start();
System.out.println("gRPC server started at :50051");
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
server.shutdown();
try {
if (!server.awaitTermination(5, TimeUnit.SECONDS)) {
server.shutdownNow();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
server.shutdownNow();
}
}));
server.awaitTermination();
}
}
客户端:
public final class ClientMain {
public static void main(String[] args) {
ManagedChannel channel = ManagedChannelBuilder
.forAddress("localhost", 50051)
.usePlaintext()
.build();
try {
var stub = OrderServiceGrpc.newBlockingStub(channel);
var reply = stub
.withDeadlineAfter(800, TimeUnit.MILLISECONDS)
.getOrder(GetOrderRequest.newBuilder()
.setOrderId("o-100")
.build());
System.out.printf(
"order=%s status=%s version=%d%n",
reply.getOrderId(),
reply.getStatus(),
reply.getVersion());
} finally {
channel.shutdown();
}
}
}
先启动服务端,再运行客户端,预期输出:
order=o-100 status=PAID version=7
如果客户端请求 order_id = "missing",服务端返回 NOT_FOUND,客户端不会得到 OrderReply,而是收到 StatusRuntimeException。
五、拦截器:统一处理调用,但不替代业务逻辑
拦截器(Interceptor)是在 RPC 生命周期的统一切入点。它适合处理横切逻辑,例如:
- 请求日志和耗时;
- Trace ID;
- 认证令牌;
- 指标;
- 统一异常映射;
- 审计字段;
- 限流或租户校验。
拦截器不应承载订单扣款、库存扣减等业务流程,因为它无法自然表达业务状态,也容易让调用链变得隐蔽。
5.1 Metadata:传输附加信息
Metadata 是键值对,但二进制键和 ASCII 键的命名有规则:
- 普通键使用 ASCII;
- 二进制键以
-bin结尾; - 不应把大对象放进 Metadata;
- Metadata 不是业务消息替代品。
定义一个客户端拦截器:
import io.grpc.CallOptions;
import io.grpc.Channel;
import io.grpc.ClientCall;
import io.grpc.ClientInterceptor;
import io.grpc.ClientCall.Listener;
import io.grpc.Metadata;
import io.grpc.MethodDescriptor;
import java.util.UUID;
public final class TraceClientInterceptor
implements ClientInterceptor {
private static final Metadata.Key<String> TRACE_ID =
Metadata.Key.of("x-trace-id", Metadata.ASCII_STRING_MARSHALLER);
@Override
public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
MethodDescriptor<ReqT, RespT> method,
CallOptions callOptions,
Channel next) {
ClientCall<ReqT, RespT> call =
next.newCall(method, callOptions);
return new ForwardingClientCall.SimpleForwardingClientCall<>(
call) {
@Override
public void start(
Listener<RespT> responseListener,
Metadata headers) {
headers.put(TRACE_ID, UUID.randomUUID().toString());
super.start(responseListener, headers);
}
};
}
}
应用到 Channel:
ManagedChannel channel = ManagedChannelBuilder
.forAddress("localhost", 50051)
.usePlaintext()
.intercept(new TraceClientInterceptor())
.build();
服务端可以通过 ServerInterceptor 读取它:
import io.grpc.Context;
import io.grpc.Contexts;
import io.grpc.Metadata;
import io.grpc.ServerCall;
import io.grpc.ServerCallHandler;
import io.grpc.ServerInterceptor;
import io.grpc.ServerCall.Listener;
import io.grpc.Status;
public final class TraceServerInterceptor
implements ServerInterceptor {
public static final Context.Key<String> TRACE_ID =
Context.key("trace-id");
private static final Metadata.Key<String> TRACE_HEADER =
Metadata.Key.of("x-trace-id", Metadata.ASCII_STRING_MARSHALLER);
@Override
public <ReqT, RespT> Listener<ReqT> interceptCall(
ServerCall<ReqT, RespT> call,
Metadata headers,
ServerCallHandler<ReqT, RespT> next) {
String traceId = headers.get(TRACE_HEADER);
if (traceId == null || traceId.isBlank()) {
call.close(
Status.UNAUTHENTICATED
.withDescription("missing trace id"),
new Metadata());
return new ServerCall.Listener<>() {};
}
Context context = Context.current()
.withValue(TRACE_ID, traceId);
return Contexts.interceptCall(context, call, headers, next);
}
}
服务实现中读取:
String traceId = TraceServerInterceptor.TRACE_ID.get();
Context 适合传播请求范围的数据和取消信号,但不要把它当作通用可变容器。不要在其中放置连接池、全局缓存或需要显式关闭的资源。
5.2 服务端异常边界
业务代码中的异常不应直接穿透到客户端:
@Override
public void getOrder(
GetOrderRequest request,
StreamObserver<OrderReply> observer) {
try {
OrderReply result = repository.find(request.getOrderId());
observer.onNext(result);
observer.onCompleted();
} catch (OrderNotFoundException e) {
observer.onError(Status.NOT_FOUND
.withDescription("order does not exist")
.asRuntimeException());
} catch (IllegalArgumentException e) {
observer.onError(Status.INVALID_ARGUMENT
.withDescription(e.getMessage())
.asRuntimeException());
} catch (Exception e) {
observer.onError(Status.INTERNAL
.withDescription("unexpected server failure")
.withCause(e)
.asRuntimeException());
}
}
withCause(e) 主要服务于服务端日志和调试,不能假设底层异常堆栈会安全地原样传到客户端。生产环境还应避免把 SQL、文件路径或内部主机名放进 description。
六、Deadline:时间预算,不是简单的 socket timeout
6.1 Deadline 的定义
Deadline 是一个绝对时间点,表示 RPC 必须在该时间点前完成。客户端写:
stub.withDeadlineAfter(800, TimeUnit.MILLISECONDS)
框架会以当前时间为基准计算截止时刻:
其中:
- :调用开始时刻;
- :调用预算,例如 800 ms;
- :绝对截止时间。
在调用过程中,剩余预算是:
当 时,调用应该被取消,并返回 DEADLINE_EXCEEDED。
这与只给数据库设置 800 ms 超时不同。RPC 总耗时还包括排队、连接建立、名称解析、服务端业务、响应传输和客户端处理:
如果每一层都重新设置 800 ms,整个调用链可能远超 800 ms。
6.2 Deadline 传播和下游调用
服务端处理请求时可以读取当前 Context 的 Deadline:
Deadline deadline = Context.current().getDeadline();
if (deadline != null) {
long remainingNanos =
deadline.timeRemaining(TimeUnit.NANOSECONDS);
System.out.println("remaining nanos = " + remainingNanos);
}
当服务端调用下游 gRPC 服务时,应显式使用剩余预算或设置更短的本地预算:
long remainingMillis = Math.max(
1,
Context.current()
.getDeadline()
.timeRemaining(TimeUnit.MILLISECONDS));
DownstreamReply reply = downstreamStub
.withDeadlineAfter(remainingMillis, TimeUnit.MILLISECONDS)
.call(request);
实际代码需要处理 Context.current().getDeadline() 为 null 的情况。对于具有明确 SLA 的服务,更推荐在入口设置总预算,在每个下游调用上使用不超过剩余预算的 Deadline。
一个三层调用例子:
客户端总预算:800 ms
├── 网关处理:最多 100 ms
├── 订单服务:剩余约 700 ms
│ ├── 数据库:最多 300 ms
│ └── 库存服务:最多 350 ms
└── 响应序列化和网络传输:使用剩余时间
如果库存服务仍配置固定 800 ms,就可能发生:
客户端 800 ms 已超时
库存服务仍在执行,直到自己的 800 ms 超时
这会造成客户端已经失败、服务端资源却继续占用的“超时失控”。
6.3 Deadline 到期时发生什么
Deadline 到期通常包含以下路径:
- 客户端取消本次 RPC;
- 服务端的
Context.current().isCancelled()变为true; - 服务端继续执行的阻塞任务不会被框架自动杀死;
- 客户端收到
DEADLINE_EXCEEDED; - 服务端如果之后尝试发送消息,发送可能失败。
因此,业务代码需要配合取消:
Context context = Context.current();
while (hasMoreWork()) {
if (context.isCancelled()) {
cleanup();
return;
}
processOneItem();
}
对 JDBC、远程 HTTP 客户端、文件读写等阻塞操作,必须进一步配置它们自己的超时或取消机制。gRPC Deadline 不能自动中断任意第三方阻塞调用。
6.4 Deadline 与重试的乘法关系
如果一次调用允许重试,单次尝试预算不能脱离总 Deadline。设总预算为 ,已有耗时为 ,那么下一次尝试最多可用:
如果重试策略错误地为每次尝试设置固定 1 秒,而总 Deadline 只有 800 ms,那么第二次尝试即使开始,也没有真实的剩余时间。
只有满足以下条件时,自动重试才可能安全:
- 操作具有幂等性,或者请求携带幂等键;
- 状态码属于可重试集合;
- 重试不会突破总 Deadline;
- 有最大尝试次数和退避;
- 服务端不会因重复请求产生重复副作用。
UNAVAILABLE 通常表示暂时不可用,可能适合重试;INVALID_ARGUMENT 和 PERMISSION_DENIED 通常不适合重试;DEADLINE_EXCEEDED 是否重试则取决于剩余预算和操作语义,不能统一处理。
七、取消、并发和流控
7.1 取消是正常控制流
取消可能来自:
- 客户端主动取消;
- Deadline 到期;
- HTTP/2 连接断开;
- 服务端主动拒绝;
- 进程关闭。
服务端可以注册取消回调:
Context.current().addListener(
context -> {
if (context.isCancelled()) {
System.out.println("rpc cancelled: "
+ context.cancellationCause());
}
},
command -> command.run());
实际项目中应确保回调只做轻量操作,例如设置停止标志、取消下游任务、释放订阅资源,而不是在回调中再次执行长时间阻塞操作。
7.2 StreamObserver 不是线程安全的业务队列
gRPC Java 的回调可能运行在框架管理的执行器中。业务代码若从多个线程同时调用同一个 responseObserver.onNext,必须自己保证并发安全和消息顺序。
一种简单的做法是用单线程执行器串行发送:
ExecutorService sender = Executors.newSingleThreadExecutor();
sender.execute(() -> responseObserver.onNext(event));
但这只是串行化发送,不会无限扩大接收能力。若生产速度大于消费速度:
积压会增长:
所以流式服务必须设置有界队列、丢弃策略、批量策略或应用层确认机制。HTTP/2 流控可以限制网络层发送窗口,但它不是业务层的无限内存保护,也不替代队列容量设计。
八、状态码:把失败分类,而不是只记录异常字符串
gRPC 使用 Status.Code 描述调用结果。常见语义如下:
| Code | 常见含义 |
|---|---|
OK |
调用成功 |
INVALID_ARGUMENT |
请求参数不合法 |
UNAUTHENTICATED |
缺少或无法验证身份 |
PERMISSION_DENIED |
身份有效但没有权限 |
NOT_FOUND |
资源不存在 |
ALREADY_EXISTS |
创建资源时已存在 |
FAILED_PRECONDITION |
当前状态不满足操作前提 |
ABORTED |
并发冲突或事务中止 |
RESOURCE_EXHAUSTED |
配额、限流或资源耗尽 |
UNAVAILABLE |
暂时不可用,可能是连接或实例问题 |
DEADLINE_EXCEEDED |
超过时间预算 |
INTERNAL |
服务内部错误 |
UNIMPLEMENTED |
方法未实现 |
错误分类的因果关系很重要:
请求字段为空
-> INVALID_ARGUMENT
资源不存在
-> NOT_FOUND
令牌无效
-> UNAUTHENTICATED
服务实例暂时不可达
-> UNAVAILABLE
服务处理超过总预算
-> DEADLINE_EXCEEDED
程序未捕获的内部异常
-> INTERNAL
不要把所有异常都映射为 INTERNAL。这样客户端无法判断是否应修正请求、刷新认证、等待后重试或报警。
九、治理:从“能调用”到“可运行”
gRPC 治理不是单独的某个 API,而是对名称解析、负载均衡、连接、认证、超时、重试、限流和可观测性的组合管理。
9.1 名称解析和负载均衡
客户端通常连接一个目标地址:
ManagedChannel channel = ManagedChannelBuilder
.forTarget("dns:///order-service.example.com")
.useTransportSecurity()
.build();
dns:/// 表示通过 DNS 解析目标。实际生产环境也可能使用注册中心、xDS 或自定义 NameResolver。解析器负责把逻辑服务名转换成后端地址;负载均衡器再决定请求发送到哪个地址。
需要区分:
- 连接级负载均衡:选择一个后端建立连接;
- RPC 级负载均衡:在调用层选择具体后端;
- 流式 RPC:一个长流通常绑定一个后端,不能像独立 Unary 请求一样逐消息迁移。
因此,长时间 Server Streaming 连接会导致连接分布与短请求不同。治理系统必须单独观察长流数量、持续时间和后端倾斜。
9.2 TLS 与认证
开发环境可以使用:
.usePlaintext()
生产环境不能因为“内网”就默认安全。至少需要考虑:
- 服务端证书校验;
- 客户端如何验证目标身份;
- 是否使用双向 TLS;
- Metadata 中的令牌是否会被安全传输;
- 证书轮换期间已有 Channel 如何更新。
认证拦截器负责读取凭证并附加 Metadata,但真正的身份校验应在服务端完成。客户端传来的 user-id、tenant-id 只能作为声明,不能直接当作可信身份。
9.3 限流与资源隔离
限流可以放在 ServerInterceptor 中做入口拒绝,但仅依赖拦截器计数还不够。不同 RPC 的资源成本可能不同:
GetOrder:一次数据库读取
WatchOrder:长期订阅和推送资源
UploadOrderEvents:持续占用解析和存储资源
如果三者共享同一个无界线程池或无限制连接数,低成本 Unary 请求可能被长流拖垮。更可靠的做法是按方法、租户或优先级设置:
- 最大并发 Unary 请求;
- 最大长流数;
- 单流消息速率;
- 最大消息大小;
- 单租户配额;
- 独立执行器或资源池。
消息大小限制也很重要。gRPC Java 可以在 Channel 或 Server 上配置最大入站消息大小,例如:
ManagedChannel channel = ManagedChannelBuilder
.forAddress("localhost", 50051)
.usePlaintext()
.maxInboundMessageSize(4 * 1024 * 1024)
.build();
这不是完整的业务防护。一个合法大小的消息仍可能触发高昂的数据库操作,因此还需要业务字段限制和请求级配额。
9.4 可观测性
一次 RPC 至少应记录或导出以下维度:
- 完整方法名;
- Status Code;
- 总耗时;
- Deadline 剩余时间或是否超时;
- 请求和响应消息数量;
- 消息字节数;
- 重试次数;
- 取消原因;
- 后端地址或实例标识;
- trace ID;
- 认证主体和租户,但必须遵循脱敏规则。
不要把完整 Protobuf 请求直接写入日志。请求可能含有令牌、手机号、地址或大字段;即使没有隐私问题,也可能产生巨量日志和序列化开销。
一个合理的诊断顺序是:
客户端看到 DEADLINE_EXCEEDED
-> 查看客户端总耗时和 Deadline
-> 查看服务端是否收到请求
-> 查看服务端 Context 是否已取消
-> 查看数据库/下游调用剩余预算
-> 查看连接、DNS、TLS 和排队耗时
-> 判断是慢、不可达、拒绝还是客户端提前取消
只看“服务端接口平均耗时”是不够的,因为超时可能发生在客户端排队、名称解析或连接建立阶段。
十、Spring Boot 集成时应保持边界清晰
Spring Boot 可以管理配置、生命周期、依赖注入和监控,但 gRPC 服务注册、端口和拦截器仍取决于所使用的 gRPC Spring 集成方案。Spring Framework 和 Spring Boot 官方参考文档并不等同于 grpc-java API 文档,因此不能把某个第三方 starter 的配置项误认为 gRPC 标准。
在 Spring Boot 应用中,建议保持以下边界:
Spring Bean
-> 管理 Repository、Service、配置和生命周期
gRPC ServiceImpl
-> 将 Protobuf 请求转换为应用命令
应用层
-> 执行业务规则和事务
gRPC 层
-> 将领域结果和异常映射为 Protobuf 与 Status
例如:
@Component
public final class GrpcOrderService
extends OrderServiceGrpc.OrderServiceImplBase {
private final OrderApplicationService applicationService;
public GrpcOrderService(
OrderApplicationService applicationService) {
this.applicationService = applicationService;
}
@Override
public void getOrder(
GetOrderRequest request,
StreamObserver<OrderReply> observer) {
try {
OrderResult result =
applicationService.getOrder(request.getOrderId());
observer.onNext(OrderReply.newBuilder()
.setOrderId(result.id())
.setStatus(result.status())
.setVersion(result.version())
.build());
observer.onCompleted();
} catch (OrderNotFoundException e) {
observer.onError(Status.NOT_FOUND
.withDescription("order not found")
.asRuntimeException());
}
}
}
这里的 @Component 只负责让 Spring 管理对象;“如何把这个 Bean 加入 gRPC Server”是具体集成库的职责,不能假设所有 starter 的自动配置方式相同。生产环境应检查:
- Bean 是否只注册一次;
- Spring 容器关闭时 gRPC Server 是否优雅停止;
- 拦截器是否实际作用于目标服务;
- 端口、TLS、健康检查和反射服务是否按环境配置;
grpc.version与 starter 支持范围是否匹配。
十一、常见错误及其失败表现
错误一:每个请求都创建 Channel
错误方式:
ManagedChannel channel = ManagedChannelBuilder
.forAddress(host, port)
.usePlaintext()
.build();
return OrderServiceGrpc.newBlockingStub(channel).getOrder(request);
问题是连接、线程和名称解析资源会快速膨胀,且请求结束后还可能忘记关闭 Channel。正确方向是复用长期 Channel,用 Stub 派生调用选项。
错误二:用 Thread.sleep 模拟服务耗时,却不检查取消
Thread.sleep(10_000);
observer.onNext(reply);
客户端 500 ms 超时后,服务端线程仍可能继续睡眠 10 秒,随后再尝试发送。测试看起来只是“客户端超时”,生产中却会累积大量无效工作。应使用可取消的任务和下游超时,并检查 Context.current().isCancelled()。
错误三:把 Deadline 当成服务端执行终止器
Deadline 能传播取消信号,但不会强制终止任意 Java 线程,也不会自动中断 JDBC、文件系统或第三方 HTTP 调用。若下游库没有超时,gRPC 已返回失败后,下游操作仍可能继续。
错误四:在每次失败后无条件重试
以下逻辑具有明显风险:
for (int i = 0; i < 3; i++) {
try {
return blockingStub.createOrder(request);
} catch (StatusRuntimeException e) {
// 无条件重试
}
}
如果第一次请求已经在服务端成功写库,只是响应丢失,第二次请求可能重复创建订单。重试前必须确认幂等性、状态码和剩余 Deadline。
错误五:把业务错误作为 INTERNAL
catch (Exception e) {
observer.onError(Status.INTERNAL.asRuntimeException());
}
如果 e 实际是参数错误或资源不存在,客户端会错误地报警、重试或返回 500。应在边界处建立明确的异常到 Status 映射,并对未知异常统一记录日志后返回 INTERNAL。
错误六:无限制地积压流消息
生产者快于消费者时,内存队列持续增长。服务端最终可能因为堆外缓冲、堆内对象或连接数耗尽而雪崩。流式接口应定义最大消息数、最大持续时间、最大并发流数和取消后的清理行为。
十二、测试和验证
12.1 验证 Protobuf 兼容性
每次修改 .proto 后至少验证:
- 旧客户端调用新服务;
- 新客户端调用旧服务;
- 增加字段时旧客户端仍能读取核心字段;
- 删除字段后编号被
reserved; - 枚举增加值时旧客户端能处理未知值;
- 必须区分“缺失”和“默认值”的字段使用
optional。
12.2 验证 Deadline
可以在服务端临时加入延迟:
try {
Thread.sleep(1_000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
客户端设置 100 ms:
stub.withDeadlineAfter(100, TimeUnit.MILLISECONDS)
.getOrder(request);
预期客户端抛出:
StatusRuntimeException:
DEADLINE_EXCEEDED
但这个测试不能证明业务线程被中断,只能证明 RPC 层在预算到期后向客户端报告了超时。还应检查服务端日志是否观察到 Context 取消,以及延迟任务是否真正停止。
12.3 验证流结束和取消
Server Streaming 测试应覆盖:
- 正常收到多个消息后
onCompleted; - 中途服务端返回错误;
- 客户端 Deadline 到期;
- 客户端主动取消;
- 服务端停止发送并释放订阅资源;
- 重连后是否重复收到事件;
- 事件是否需要序号和幂等处理。
流式测试不能只断言“收到了三条消息”,还要断言终止状态和资源清理。
十三、最终的设计判断
选择 gRPC 的核心理由不是“Java 方法调用看起来方便”,而是它把以下内容纳入了一个明确的契约:
Protobuf 消息结构
+ RPC 方向和消息数量
+ HTTP/2 传输
+ Metadata
+ Deadline 与取消
+ Status 错误语义
+ 拦截器
+ 名称解析与负载均衡
+ TLS、限流和可观测性
可以用以下原则判断设计是否合理:
- 请求查询通常使用 Unary;
- 服务端持续产生结果时使用 Server Streaming;
- 客户端上传大量消息并最终汇总时使用 Client Streaming;
- 双方需要同时持续通信时才使用 Bidirectional Streaming;
- 字段编号是长期协议资产,不能复用;
- Deadline 是整条调用链的时间预算,不是单层 socket 超时;
Context取消必须传递到业务和下游资源;- 拦截器处理横切逻辑,业务状态留在服务实现和应用层;
- 重试必须建立在幂等性、状态码和剩余预算之上;
- 长流必须单独治理连接数、生命周期、背压和重连;
- Spring Boot 可以管理应用生命周期,但不能掩盖 gRPC 传输和协议边界。
当这些关系被正确建立后,gRPC 才不仅是“生成 Stub 后调用方法”,而是一套可演进、可取消、可诊断并能够接受生产故障的服务通信协议。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Spring Integration:Channel、Adapter、Gateway、流程和错误处理
- 下一篇:Java Kafka:Producer、Consumer、分区、Offset、事务和再均衡
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论