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 调用可以抽象为:

RPC=Method+Metadata+Message Stream+Status\text{RPC} = \text{Method}+ \text{Metadata}+ \text{Message Stream}+ \text{Status}

其中:

  • Method:服务名和方法名,例如 /wr.order.OrderService/GetOrder
  • Metadata:请求头和响应头,通常承载认证令牌、链路追踪标识、租户信息等。
  • Message Stream:一个或多个 Protobuf 消息。
  • Status:最终状态,例如 OKNOT_FOUNDDEADLINE_EXCEEDED

“Stream”在这里有两层含义:

  1. HTTP/2 本身通过一个 stream 承载一次 RPC。
  2. 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)承担两项职责:

  1. 描述消息结构;
  2. 描述 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 的作用是防止后续开发者重新使用已经删除的字段编号或名称。

兼容性可以用一条约束表示:

Reuse(field number)=false\text{Reuse(field number)} = \text{false}

字段编号的语义必须保持稳定。改变字段类型也不能只看 Java 类型是否“能转换”,还要看 Protobuf wire type 是否兼容。例如某些整数类型之间在 wire level 上可兼容,但可能发生溢出或语义变化;stringbytes 也不能因为都表现为字节序列就任意替换。

2.3 proto3 默认值与 presence

在普通 proto3 标量字段中,未设置的 string 通常读取为空字符串,int32 读取 0bool 读取 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 和方法描述;
  • GetOrderRequestOrderReply 等:由 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
  • onCompletedonError 不能同时调用;
  • 终止后不能继续发送消息。

对 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);

然后在调用时通过 withDeadlineAfterwithInterceptors 等方法派生新的 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)

框架会以当前时间为基准计算截止时刻:

D=T0+BD = T_0 + B

其中:

  • T0T_0:调用开始时刻;
  • BB:调用预算,例如 800 ms;
  • DD:绝对截止时间。

在调用过程中,剩余预算是:

R(t)=DtR(t) = D - t

R(t)0R(t) \leq 0 时,调用应该被取消,并返回 DEADLINE_EXCEEDED

这与只给数据库设置 800 ms 超时不同。RPC 总耗时还包括排队、连接建立、名称解析、服务端业务、响应传输和客户端处理:

Trpc=Tqueue+Tconnect+Tserver+Tnetwork+TclientT_{\text{rpc}} = T_{\text{queue}}+ T_{\text{connect}}+ T_{\text{server}}+ T_{\text{network}}+ T_{\text{client}}

如果每一层都重新设置 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 到期通常包含以下路径:

  1. 客户端取消本次 RPC;
  2. 服务端的 Context.current().isCancelled() 变为 true
  3. 服务端继续执行的阻塞任务不会被框架自动杀死;
  4. 客户端收到 DEADLINE_EXCEEDED
  5. 服务端如果之后尝试发送消息,发送可能失败。

因此,业务代码需要配合取消:

Context context = Context.current();

while (hasMoreWork()) {
    if (context.isCancelled()) {
        cleanup();
        return;
    }

    processOneItem();
}

对 JDBC、远程 HTTP 客户端、文件读写等阻塞操作,必须进一步配置它们自己的超时或取消机制。gRPC Deadline 不能自动中断任意第三方阻塞调用。

6.4 Deadline 与重试的乘法关系

如果一次调用允许重试,单次尝试预算不能脱离总 Deadline。设总预算为 BB,已有耗时为 EE,那么下一次尝试最多可用:

BnextBEB_{\text{next}} \leq B-E

如果重试策略错误地为每次尝试设置固定 1 秒,而总 Deadline 只有 800 ms,那么第二次尝试即使开始,也没有真实的剩余时间。

只有满足以下条件时,自动重试才可能安全:

  1. 操作具有幂等性,或者请求携带幂等键;
  2. 状态码属于可重试集合;
  3. 重试不会突破总 Deadline;
  4. 有最大尝试次数和退避;
  5. 服务端不会因重复请求产生重复副作用。

UNAVAILABLE 通常表示暂时不可用,可能适合重试;INVALID_ARGUMENTPERMISSION_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));

但这只是串行化发送,不会无限扩大接收能力。若生产速度大于消费速度:

λproducer>λconsumer\lambda_{\text{producer}} > \lambda_{\text{consumer}}

积压会增长:

Q(t)Q(0)+(λproducerλconsumer)tQ(t) \approx Q(0) + (\lambda_{\text{producer}}-\lambda_{\text{consumer}})t

所以流式服务必须设置有界队列、丢弃策略、批量策略或应用层确认机制。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-idtenant-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 后至少验证:

  1. 旧客户端调用新服务;
  2. 新客户端调用旧服务;
  3. 增加字段时旧客户端仍能读取核心字段;
  4. 删除字段后编号被 reserved
  5. 枚举增加值时旧客户端能处理未知值;
  6. 必须区分“缺失”和“默认值”的字段使用 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、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。