本文目录导读:

我来详细说明 gRPC Java 的使用方法。
环境准备
添加 Maven 依赖
<dependencies>
<!-- gRPC 核心依赖 -->
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-netty-shaded</artifactId>
<version>1.60.0</version>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-protobuf</artifactId>
<version>1.60.0</version>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-stub</artifactId>
<version>1.60.0</version>
</dependency>
<!-- Protocol Buffers -->
<dependency>
<groupId>com.google.protobuf</groupId>
<artifactId>protobuf-java</artifactId>
<version>3.25.1</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:3.25.1:exe:${os.detected.classifier}
</protocArtifact>
<pluginId>grpc-java</pluginId>
<pluginArtifact>
io.grpc:protoc-gen-grpc-java:1.60.0:exe:${os.detected.classifier}
</pluginArtifact>
</configuration>
<executions>
<execution>
<goals>
<goal>compile</goal>
<goal>compile-custom</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
定义 Protobuf 服务
创建 src/main/proto/helloworld.proto:
syntax = "proto3";
package helloworld;
option java_package = "com.example.grpc";
option java_multiple_files = true;
// 定义服务
service Greeter {
// 简单 RPC
rpc SayHello (HelloRequest) returns (HelloReply);
// 服务端流式 RPC
rpc SayHelloServerStream (HelloRequest) returns (stream HelloReply);
// 客户端流式 RPC
rpc SayHelloClientStream (stream HelloRequest) returns (HelloReply);
// 双向流式 RPC
rpc SayHelloBidiStream (stream HelloRequest) returns (stream HelloReply);
}
// 请求消息
message HelloRequest {
string name = 1;
}
// 响应消息
message HelloReply {
string message = 1;
}
运行 mvn compile 生成 Java 代码。
实现服务端
package com.example.grpc;
import io.grpc.Server;
import io.grpc.ServerBuilder;
import io.grpc.stub.StreamObserver;
import java.io.IOException;
import java.util.concurrent.TimeUnit;
import java.util.logging.Logger;
public class HelloWorldServer {
private static final Logger logger = Logger.getLogger(HelloWorldServer.class.getName());
private Server server;
private void start() throws IOException {
// 启动 gRPC 服务器,监听端口 50051
int port = 50051;
server = ServerBuilder.forPort(port)
.addService(new GreeterImpl())
.build()
.start();
logger.info("Server started, listening on " + port);
// 添加 JVM 关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
System.err.println("*** shutting down gRPC server since JVM is shutting down");
try {
HelloWorldServer.this.stop();
} catch (InterruptedException e) {
e.printStackTrace(System.err);
}
System.err.println("*** server shut down");
}));
}
private void stop() throws InterruptedException {
if (server != null) {
server.shutdown().awaitTermination(30, TimeUnit.SECONDS);
}
}
// 阻塞等待服务器关闭
private void blockUntilShutdown() throws InterruptedException {
if (server != null) {
server.awaitTermination();
}
}
// 实现服务接口
static class GreeterImpl extends GreeterGrpc.GreeterImplBase {
// 简单 RPC
@Override
public void sayHello(HelloRequest req, StreamObserver<HelloReply> responseObserver) {
HelloReply reply = HelloReply.newBuilder()
.setMessage("Hello " + req.getName())
.build();
responseObserver.onNext(reply);
responseObserver.onCompleted();
}
// 服务端流式 RPC
@Override
public void sayHelloServerStream(HelloRequest req,
StreamObserver<HelloReply> responseObserver) {
for (int i = 0; i < 5; i++) {
HelloReply reply = HelloReply.newBuilder()
.setMessage("Hello " + req.getName() + " - Message " + i)
.build();
responseObserver.onNext(reply);
try {
Thread.sleep(1000); // 模拟延迟
} catch (InterruptedException e) {
e.printStackTrace();
}
}
responseObserver.onCompleted();
}
// 客户端流式 RPC
@Override
public StreamObserver<HelloRequest> sayHelloClientStream(
StreamObserver<HelloReply> responseObserver) {
return new StreamObserver<HelloRequest>() {
StringBuilder names = new StringBuilder();
@Override
public void onNext(HelloRequest request) {
names.append(request.getName()).append(", ");
}
@Override
public void onError(Throwable t) {
logger.warning("Error: " + t.getMessage());
}
@Override
public void onCompleted() {
HelloReply reply = HelloReply.newBuilder()
.setMessage("Hello " + names.toString())
.build();
responseObserver.onNext(reply);
responseObserver.onCompleted();
}
};
}
// 双向流式 RPC
@Override
public StreamObserver<HelloRequest> sayHelloBidiStream(
StreamObserver<HelloReply> responseObserver) {
return new StreamObserver<HelloRequest>() {
@Override
public void onNext(HelloRequest request) {
HelloReply reply = HelloReply.newBuilder()
.setMessage("Hello " + request.getName() + " (Bidi)")
.build();
responseObserver.onNext(reply);
}
@Override
public void onError(Throwable t) {
logger.warning("Error: " + t.getMessage());
}
@Override
public void onCompleted() {
responseObserver.onCompleted();
}
};
}
}
public static void main(String[] args) throws IOException, InterruptedException {
final HelloWorldServer server = new HelloWorldServer();
server.start();
server.blockUntilShutdown();
}
}
实现客户端
package com.example.grpc;
import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.logging.Logger;
public class HelloWorldClient {
private static final Logger logger = Logger.getLogger(HelloWorldClient.class.getName());
private final GreeterGrpc.GreeterBlockingStub blockingStub;
private final GreeterGrpc.GreeterStub asyncStub;
public HelloWorldClient(String host, int port) {
// 创建 gRPC 通道
ManagedChannel channel = ManagedChannelBuilder.forAddress(host, port)
.usePlaintext() // 开发环境使用明文,生产环境应使用 TLS
.build();
// 创建两种类型的 stub
blockingStub = GreeterGrpc.newBlockingStub(channel);
asyncStub = GreeterGrpc.newStub(channel);
}
// 简单的 RPC 调用
public void sayHello(String name) {
HelloRequest request = HelloRequest.newBuilder()
.setName(name)
.build();
HelloReply response = blockingStub.sayHello(request);
logger.info("Response: " + response.getMessage());
}
// 服务端流式 RPC
public void sayHelloServerStream(String name) throws InterruptedException {
HelloRequest request = HelloRequest.newBuilder()
.setName(name)
.build();
CountDownLatch latch = new CountDownLatch(1);
asyncStub.sayHelloServerStream(request, new StreamObserver<HelloReply>() {
@Override
public void onNext(HelloReply response) {
logger.info("Received: " + response.getMessage());
}
@Override
public void onError(Throwable t) {
logger.warning("Error: " + t.getMessage());
latch.countDown();
}
@Override
public void onCompleted() {
logger.info("Server streaming completed");
latch.countDown();
}
});
latch.await(10, TimeUnit.SECONDS);
}
// 客户端流式 RPC
public void sayHelloClientStream() throws InterruptedException {
CountDownLatch latch = new CountDownLatch(1);
StreamObserver<HelloRequest> requestObserver =
asyncStub.sayHelloClientStream(new StreamObserver<HelloReply>() {
@Override
public void onNext(HelloReply response) {
logger.info("Received: " + response.getMessage());
}
@Override
public void onError(Throwable t) {
logger.warning("Error: " + t.getMessage());
latch.countDown();
}
@Override
public void onCompleted() {
logger.info("Client streaming completed");
latch.countDown();
}
});
// 发送多个请求
for (int i = 0; i < 3; i++) {
HelloRequest request = HelloRequest.newBuilder()
.setName("User " + i)
.build();
requestObserver.onNext(request);
}
requestObserver.onCompleted();
latch.await(10, TimeUnit.SECONDS);
}
public static void main(String[] args) throws InterruptedException {
HelloWorldClient client = new HelloWorldClient("localhost", 50051);
// 测试不同的 RPC 类型
System.out.println("=== Simple RPC ===");
client.sayHello("World");
System.out.println("\n=== Server Streaming ===");
client.sayHelloServerStream("World");
System.out.println("\n=== Client Streaming ===");
client.sayHelloClientStream();
}
}
运行示例
- 编译项目:
mvn clean compile - 启动服务端:运行
HelloWorldServer.main() - 启动客户端:运行
HelloWorldClient.main()
SSL/TLS 配置(生产环境)
服务端
// 生成自签名证书
Server server = ServerBuilder.forPort(port)
.useTransportSecurity(
new File("server.crt"), // 证书文件
new File("server.pem") // 私钥文件
)
.addService(new GreeterImpl())
.build();
客户端
ManagedChannel channel = ManagedChannelBuilder.forAddress(host, port)
.useTransportSecurity()
.sslContext(GrpcSslContexts.forClient()
.trustManager(new File("ca.crt"))
.build())
.build();
- 定义 proto 文件:描述服务和消息结构
- 生成代码:使用 protobuf-maven-plugin 自动生成
- 实现服务端:继承生成的抽象类,实现业务逻辑
- 创建客户端:使用 stub 调用远程方法
- 错误处理:实现适当的异常处理机制
- 生产配置:必须使用 SSL/TLS 加密通信
这就是 gRPC Java 的基本使用方法,根据你的具体需求,还可以添加拦截器、负载均衡、健康检查等高级功能。