gRPC Java怎么用?

wen python案例 4

本文目录导读:

gRPC Java怎么用?

  1. 环境准备
  2. 定义 Protobuf 服务
  3. 实现服务端
  4. 实现客户端
  5. 运行示例
  6. SSL/TLS 配置(生产环境)

我来详细说明 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();
    }
}

运行示例

  1. 编译项目:mvn clean compile
  2. 启动服务端:运行 HelloWorldServer.main()
  3. 启动客户端:运行 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();
  1. 定义 proto 文件:描述服务和消息结构
  2. 生成代码:使用 protobuf-maven-plugin 自动生成
  3. 实现服务端:继承生成的抽象类,实现业务逻辑
  4. 创建客户端:使用 stub 调用远程方法
  5. 错误处理:实现适当的异常处理机制
  6. 生产配置:必须使用 SSL/TLS 加密通信

这就是 gRPC Java 的基本使用方法,根据你的具体需求,还可以添加拦截器、负载均衡、健康检查等高级功能。

抱歉,评论功能暂时关闭!