跳转至

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;
}

执行生成指令:

# Maven
mvn clean compile


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。