Java 语言工程实践指南¶
在企业级后端与微服务生态(如 Spring Cloud、Dubbo、分布式大数据计算)中,Java 占据着举足轻重的地位。gRPC-Java 是官方维护的纯 Java 原生实现,基于高性能网络框架 Netty 构建。
1. 构建工具与依赖配置¶
在 Java 项目中,通常使用 Maven 或 Gradle 配合编译器插件自动生成 Protobuf 与 gRPC 代码。
① Maven 配置 (pom.xml)¶
[!TIP] 最佳实践:优先使用
grpc-netty-shaded微服务项目(特别是 Spring Boot)中常引入多种依赖(如 Spring WebFlux、Elasticsearch、Gateway),它们内部可能自带不同版本的 Netty。使用grpc-netty-shaded会将 Netty 的类路径进行重命名重包隔离(Package Relocation),彻底杜绝令人头疼的 Netty 依赖地狱(NoSuchMethodError/ClassNotFoundException)。
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example.grpc</groupId>
<artifactId>grpc-java-demo</artifactId>
<version>1.0.0</version>
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<grpc.version>1.62.2</grpc.version>
<protobuf.version>3.25.3</protobuf.version>
<protoc.version>3.25.3</protoc.version>
</properties>
<dependencies>
<!-- gRPC 核心依赖 (推荐使用 shaded 隔离 Netty 冲突) -->
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-netty-shaded</artifactId>
<version>${grpc.version}</version>
<scope>runtime</scope>
</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>
<!-- Java 9+ 必须补充 javax.annotation 依赖 -->
<dependency>
<groupId>javax.annotation</groupId>
<artifactId>javax.annotation-api</artifactId>
<version>1.3.2</version>
</dependency>
</dependencies>
<build>
<extensions>
<!-- 自动检测当前操作系统的平台分类 (os.detected.classifier) -->
<extension>
<groupId>kr.motd.maven</groupId>
<artifactId>os-maven-plugin</artifactId>
<version>1.7.1</version>
</extension>
</extensions>
<plugins>
<!-- protobuf 编译插件 -->
<plugin>
<groupId>org.xolstice.maven.plugins</groupId>
<artifactId>protobuf-maven-plugin</artifactId>
<version>0.6.1</version>
<configuration>
<protocArtifact>com.google.protobuf:protoc:${protoc.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>
</project>
② Gradle 配置 (build.gradle)¶
plugins {
id 'java'
id 'com.google.protobuf' version '0.9.4'
}
repositories {
mavenCentral()
}
def grpcVersion = '1.62.2'
def protobufVersion = '3.25.3'
dependencies {
implementation "io.grpc:grpc-protobuf:${grpcVersion}"
implementation "io.grpc:grpc-stub:${grpcVersion}"
runtimeOnly "io.grpc:grpc-netty-shaded:${grpcVersion}"
compileOnly 'javax.annotation:javax.annotation-api:1.3.2'
}
protobuf {
protoc {
artifact = "com.google.protobuf:protoc:${protobufVersion}"
}
plugins {
grpc {
artifact = "io.grpc:protoc-gen-grpc-java:${grpcVersion}"
}
}
generateProtoTasks {
all()*.plugins {
grpc {}
}
}
}
2. 接口定义 (src/main/proto/order.proto)¶
syntax = "proto3";
package commerce.order.v1;
// Java 专项编译选项
option java_multiple_files = true; // 每个消息生成独立的 .java 类文件
option java_package = "com.example.grpc.order.v1"; // Java 包名
option java_outer_classname = "OrderProto";
service OrderService {
// 一元调用
rpc CreateOrder(CreateOrderRequest) returns (CreateOrderResponse);
// 服务端流式调用
rpc TrackOrder(TrackOrderRequest) returns (stream TrackOrderResponse);
}
message CreateOrderRequest {
string user_id = 1;
double amount = 2;
repeated string item_ids = 3;
}
message CreateOrderResponse {
string order_id = 1;
string status = 2;
}
message TrackOrderRequest {
string order_id = 1;
}
message TrackOrderResponse {
string status = 1;
string message = 2;
int64 timestamp = 3;
}
执行生成指令:
3. 服务端生产级实现¶
服务端继承自动生成的抽象类 OrderServiceImplBase,所有异步与流式通信均围绕响应观察者 StreamObserver<T> 展开:
package com.example.grpc.server;
import com.example.grpc.order.v1.*;
import io.grpc.Server;
import io.grpc.ServerBuilder;
import io.grpc.Status;
import io.grpc.stub.StreamObserver;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
import java.util.logging.Logger;
public class OrderServer {
private static final Logger logger = Logger.getLogger(OrderServer.class.getName());
private Server server;
private void start() throws IOException {
int port = 50051;
server = ServerBuilder.forPort(port)
.addService(new OrderServiceImpl())
.build()
.start();
logger.info("gRPC Java 服务端已在端口 " + port + " 启动");
// 注册 JVM 关闭钩子实现平滑关机
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
logger.info("捕获 JVM 关机信号,正在停止 gRPC 服务...");
try {
OrderServer.this.stop();
} catch (InterruptedException e) {
logger.severe("停机异常: " + e.getMessage());
}
}));
}
private void stop() throws InterruptedException {
if (server != null) {
// 优雅关闭:拒绝新请求并等待活跃 RPC 执行完成
server.shutdown().awaitTermination(30, TimeUnit.SECONDS);
}
}
private void blockUntilShutdown() throws InterruptedException {
if (server != null) {
server.awaitTermination();
}
}
// 业务服务实现
static class OrderServiceImpl extends OrderServiceGrpc.OrderServiceImplBase {
@Override
public void createOrder(CreateOrderRequest request, StreamObserver<CreateOrderResponse> responseObserver) {
// 参数校验
if (request.getUserId().isEmpty()) {
responseObserver.onError(
Status.INVALID_ARGUMENT
.withDescription("user_id 不能为空")
.asRuntimeException()
);
return;
}
String orderId = "java_ord_" + System.currentTimeMillis();
CreateOrderResponse response = CreateOrderResponse.newBuilder()
.setOrderId(orderId)
.setStatus("CREATED")
.build();
// 发送响应并宣告结束
responseObserver.onNext(response);
responseObserver.onCompleted();
}
@Override
public void trackOrder(TrackOrderRequest request, StreamObserver<TrackOrderResponse> responseObserver) {
String[] steps = {"订单确认", "分拣打包", "顺丰出库", "同城派送", "签收完成"};
for (String step : steps) {
TrackOrderResponse update = TrackOrderResponse.newBuilder()
.setStatus(step)
.setMessage("订单 " + request.getOrderId() + " 状态变更: " + step)
.setTimestamp(System.currentTimeMillis() / 1000)
.build();
responseObserver.onNext(update);
try {
Thread.sleep(500); // 模拟状态推送间隔
} catch (InterruptedException e) {
responseObserver.onError(Status.CANCELLED.withDescription("传输中断").asRuntimeException());
return;
}
}
// 推流结束
responseObserver.onCompleted();
}
}
public static void main(String[] args) throws IOException, InterruptedException {
final OrderServer server = new OrderServer();
server.start();
server.blockUntilShutdown();
}
}
4. 客户端实现与三种桩(Stub)模型¶
gRPC-Java 官方生成了三种不同编程范式的客户端桩(Stub):
| 客户端桩类型 | 类名规则 | 核心特征与适用场景 |
|---|---|---|
| 阻塞式桩 (BlockingStub) | OrderServiceBlockingStub |
同步阻塞调用,简单直观,仅支持一元 RPC 与服务端流 |
| 异步非阻塞桩 (Stub) | OrderServiceStub |
基于回调(StreamObserver),全场景通用,支持客户端流与双向流 |
| Future 桩 (FutureStub) | OrderServiceFutureStub |
基于 Guava ListenableFuture,便于与响应式编排结合(仅一元) |
客户端生产级调用示例¶
package com.example.grpc.client;
import com.example.grpc.order.v1.*;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.StatusRuntimeException;
import io.grpc.stub.StreamObserver;
import java.util.Iterator;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.logging.Logger;
public class OrderClient {
private static final Logger logger = Logger.getLogger(OrderClient.class.getName());
public static void main(String[] args) throws InterruptedException {
String target = "localhost:50051";
// 1. 创建 ManagedChannel (全局单例复用)
ManagedChannel channel = ManagedChannelBuilder.forTarget(target)
.usePlaintext() // 本地开发使用明文,生产必须配合 useTransportSecurity()
.keepAliveTime(30, TimeUnit.SECONDS)
.build();
try {
// 2. 阻塞式一元调用 (必须指定超时时间 Deadline)
OrderServiceGrpc.OrderServiceBlockingStub blockingStub =
OrderServiceGrpc.newBlockingStub(channel)
.withDeadlineAfter(3, TimeUnit.SECONDS);
CreateOrderRequest createReq = CreateOrderRequest.newBuilder()
.setUserId("usr_8888")
.setAmount(299.0)
.addItemIds("item_101")
.addItemIds("item_102")
.build();
CreateOrderResponse createResp = blockingStub.createOrder(createReq);
logger.info("订单创建成功: ID=" + createResp.getOrderId() + ", 状态=" + createResp.getStatus());
// 3. 服务端流式消费 (阻塞式迭代器)
TrackOrderRequest trackReq = TrackOrderRequest.newBuilder()
.setOrderId(createResp.getOrderId())
.build();
Iterator<TrackOrderResponse> updates = blockingStub.trackOrder(trackReq);
while (updates.hasNext()) {
TrackOrderResponse update = updates.next();
logger.info("[流事件] " + update.getStatus() + " | " + update.getMessage());
}
// 4. 异步非阻塞调用 (使用 OrderServiceStub)
OrderServiceGrpc.OrderServiceStub asyncStub = OrderServiceGrpc.newStub(channel);
CountDownLatch latch = new CountDownLatch(1);
asyncStub.trackOrder(trackReq, new StreamObserver<TrackOrderResponse>() {
@Override
public void onNext(TrackOrderResponse value) {
logger.info("[异步流通知] " + value.getStatus());
}
@Override
public void onError(Throwable t) {
logger.severe("流异常: " + t.getMessage());
latch.countDown();
}
@Override
public void onCompleted() {
logger.info("异步流接收结束");
latch.countDown();
}
});
latch.await(5, TimeUnit.SECONDS);
} catch (StatusRuntimeException e) {
logger.severe("RPC 调用失败: " + e.getStatus().getCode() + " - " + e.getStatus().getDescription());
} finally {
// 优雅关闭通道
channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);
}
}
}
5. 服务端拦截器实现 (ServerInterceptor)¶
在 Java 中实现 AOP 切面只需实现 ServerInterceptor 接口:
package com.example.grpc.interceptor;
import io.grpc.*;
import java.util.logging.Logger;
public class LoggingServerInterceptor implements ServerInterceptor {
private static final Logger logger = Logger.getLogger(LoggingServerInterceptor.class.getName());
@Override
public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(
ServerCall<ReqT, RespT> call,
Metadata headers,
ServerCallHandler<ReqT, RespT> next) {
long startTime = System.currentTimeMillis();
String methodName = call.getMethodDescriptor().getFullMethodName();
logger.info("[RPC 收到] Method=" + methodName);
// 包装 ServerCall 回调以捕获关闭状态
ForwardingServerCall.SimpleForwardingServerCall<ReqT, RespT> wrappedCall =
new ForwardingServerCall.SimpleForwardingServerCall<>(call) {
@Override
public void close(Status status, Metadata trailers) {
long duration = System.currentTimeMillis() - startTime;
logger.info("[RPC 完成] Method=" + methodName
+ ", Code=" + status.getCode()
+ ", 耗时=" + duration + "ms");
super.close(status, trailers);
}
};
return next.startCall(wrappedCall, headers);
}
}
挂载拦截器:
Server server = ServerBuilder.forPort(50051)
.addService(new OrderServiceImpl())
.intercept(new LoggingServerInterceptor()) // 挂载拦截器
.build();
6. Spring Boot 生态集成推荐¶
在现代 Java 微服务开发中,推荐引入成熟的开源整合库:
grpc-spring-boot-starter:- 声明式服务发布:只需在类上标注
@GrpcService,即可自动托管为 Spring Bean 并注册至 gRPC 端口。 - 声明式客户端注入:使用
@GrpcClient("order-service")字段注入即可开箱即用 Channel 与 Stub。 - 自动集成注册中心:原生支持 Eureka、Nacos、Consul 与 Kubernetes。