You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

grpc-java无需.proto文件定义含两个客户端流式方法的服务可行吗?

不用.proto文件在gRPC-Java中定义客户端流式服务

当然可以!gRPC-Java提供了动态API,让你完全在Java代码里定义服务、消息和方法,不需要编写.proto文件。这种方式适合快速验证原型、简单服务场景,或者需要动态生成服务的场景。下面我会一步步给你示例,实现一个包含两个客户端流式方法的服务。

第一步:添加必要依赖

首先确保你的项目里引入了gRPC的核心依赖(以Maven为例):

<dependencies>
    <!-- gRPC核心依赖 -->
    <dependency>
        <groupId>io.grpc</groupId>
        <artifactId>grpc-netty-shaded</artifactId>
        <version>1.59.0</version>
    </dependency>
    <dependency>
        <groupId>io.grpc</groupId>
        <artifactId>grpc-protobuf</artifactId>
        <version>1.59.0</version>
    </dependency>
    <dependency>
        <groupId>io.grpc</groupId>
        <artifactId>grpc-stub</artifactId>
        <version>1.59.0</version>
    </dependency>
    <!-- Protobuf Java核心库 -->
    <dependency>
        <groupId>com.google.protobuf</groupId>
        <artifactId>protobuf-java</artifactId>
        <version>3.24.3</version>
    </dependency>
</dependencies>

第二步:定义消息类型

我们用Protobuf的DynamicMessage来动态定义消息,不需要编译.proto。比如我们定义两组消息:一组用于文件上传的请求/响应,另一组用于批量日志上报的请求/响应。

先构造消息的Descriptor(相当于.proto里的消息定义):

import com.google.protobuf.Descriptors;
import com.google.protobuf.DynamicMessage;
import com.google.protobuf.FieldDescriptor.Type;

// 构建文件上传请求的Descriptor
Descriptors.Descriptor uploadRequestDescriptor = Descriptors.Descriptor.newBuilder()
    .setName("UploadRequest")
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("file_chunk")
        .setNumber(1)
        .setType(Type.BYTES)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("file_name")
        .setNumber(2)
        .setType(Type.STRING)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .build();

// 构建文件上传响应的Descriptor
Descriptors.Descriptor uploadResponseDescriptor = Descriptors.Descriptor.newBuilder()
    .setName("UploadResponse")
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("success")
        .setNumber(1)
        .setType(Type.BOOL)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("message")
        .setNumber(2)
        .setType(Type.STRING)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .build();

// 构建批量日志请求的Descriptor
Descriptors.Descriptor logRequestDescriptor = Descriptors.Descriptor.newBuilder()
    .setName("LogRequest")
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("log_line")
        .setNumber(1)
        .setType(Type.STRING)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("service_name")
        .setNumber(2)
        .setType(Type.STRING)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .build();

// 构建批量日志响应的Descriptor
Descriptors.Descriptor logSummaryResponseDescriptor = Descriptors.Descriptor.newBuilder()
    .setName("LogSummaryResponse")
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("total_lines")
        .setNumber(1)
        .setType(Type.INT32)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .addField(Descriptors.FieldDescriptor.newBuilder()
        .setName("service")
        .setNumber(2)
        .setType(Type.STRING)
        .setLabel(Descriptors.FieldDescriptor.Label.OPTIONAL)
        .build())
    .build();

第三步:定义服务和方法

接下来我们用ServiceDescriptor定义服务,然后为每个客户端流式方法创建MethodDescriptor,指定方法类型为CLIENT_STREAMING:

import io.grpc.MethodDescriptor;
import io.grpc.ServiceDescriptor;
import io.grpc.protobuf.ProtoUtils;

// 定义第一个客户端流式方法:uploadFile
MethodDescriptor<DynamicMessage, DynamicMessage> uploadFileMethod = MethodDescriptor.newBuilder(
        ProtoUtils.marshaller(DynamicMessage.newBuilder(uploadRequestDescriptor).build()),
        ProtoUtils.marshaller(DynamicMessage.newBuilder(uploadResponseDescriptor).build()))
    .setFullMethodName(MethodDescriptor.generateFullMethodName("FileService", "uploadFile"))
    .setType(MethodDescriptor.MethodType.CLIENT_STREAMING)
    .build();

// 定义第二个客户端流式方法:batchLogs
MethodDescriptor<DynamicMessage, DynamicMessage> batchLogsMethod = MethodDescriptor.newBuilder(
        ProtoUtils.marshaller(DynamicMessage.newBuilder(logRequestDescriptor).build()),
        ProtoUtils.marshaller(DynamicMessage.newBuilder(logSummaryResponseDescriptor).build()))
    .setFullMethodName(MethodDescriptor.generateFullMethodName("FileService", "batchLogs"))
    .setType(MethodDescriptor.MethodType.CLIENT_STREAMING)
    .build();

// 定义服务描述符
ServiceDescriptor fileServiceDescriptor = ServiceDescriptor.newBuilder("FileService")
    .addMethod(uploadFileMethod)
    .addMethod(batchLogsMethod)
    .build();

第四步:实现服务端逻辑

现在我们需要实现服务端的处理逻辑,通过BindableService绑定我们定义的服务,并重写方法处理逻辑:

import io.grpc.Server;
import io.grpc.ServerBuilder;
import io.grpc.stub.ServerCallStreamObserver;
import io.grpc.stub.StreamObserver;

import java.io.IOException;
import java.util.concurrent.CountDownLatch;

public class DynamicFileServiceServer {
    public static void main(String[] args) throws IOException, InterruptedException {
        Server server = ServerBuilder.forPort(8080)
                .addService(new io.grpc.BindableService() {
                    @Override
                    public ServiceDescriptor getServiceDescriptor() {
                        return fileServiceDescriptor;
                    }

                    @Override
                    public io.grpc.ServerCallHandler<?, ?> bindService(io.grpc.ServerMethodDefinition<?, ?> method) {
                        // 处理文件上传方法
                        if (method.getMethodDescriptor().equals(uploadFileMethod)) {
                            return call -> (serverCall, headers) -> {
                                ServerCallStreamObserver<DynamicMessage> responseObserver = 
                                    (ServerCallStreamObserver<DynamicMessage>) serverCall.responseObserver();
                                CountDownLatch finishLatch = new CountDownLatch(1);
                                responseObserver.setOnCloseHandler(finishLatch::countDown);

                                return new StreamObserver<DynamicMessage>() {
                                    int chunkCount = 0;
                                    String fileName = "";

                                    @Override
                                    public void onNext(DynamicMessage request) {
                                        fileName = request.getField(uploadRequestDescriptor.findFieldByName("file_name")).toString();
                                        chunkCount++;
                                    }

                                    @Override
                                    public void onError(Throwable t) {
                                        System.err.println("上传失败:" + t.getMessage());
                                        finishLatch.countDown();
                                    }

                                    @Override
                                    public void onCompleted() {
                                        // 构建响应并返回
                                        DynamicMessage response = DynamicMessage.newBuilder(uploadResponseDescriptor)
                                                .setField(uploadResponseDescriptor.findFieldByName("success"), true)
                                                .setField(uploadResponseDescriptor.findFieldByName("message"), 
                                                          "已接收" + chunkCount + "个文件块,文件名:" + fileName)
                                                .build();
                                        responseObserver.onNext(response);
                                        responseObserver.onCompleted();
                                        finishLatch.countDown();
                                    }
                                };
                            };
                        } 
                        // 处理批量日志方法
                        else if (method.getMethodDescriptor().equals(batchLogsMethod)) {
                            return call -> (serverCall, headers) -> {
                                ServerCallStreamObserver<DynamicMessage> responseObserver = 
                                    (ServerCallStreamObserver<DynamicMessage>) serverCall.responseObserver();
                                CountDownLatch finishLatch = new CountDownLatch(1);
                                responseObserver.setOnCloseHandler(finishLatch::countDown);

                                return new StreamObserver<DynamicMessage>() {
                                    int logCount = 0;
                                    String serviceName = "";

                                    @Override
                                    public void onNext(DynamicMessage request) {
                                        serviceName = request.getField(logRequestDescriptor.findFieldByName("service_name")).toString();
                                        logCount++;
                                    }

                                    @Override
                                    public void onError(Throwable t) {
                                        System.err.println("日志上报失败:" + t.getMessage());
                                        finishLatch.countDown();
                                    }

                                    @Override
                                    public void onCompleted() {
                                        DynamicMessage response = DynamicMessage.newBuilder(logSummaryResponseDescriptor)
                                                .setField(logSummaryResponseDescriptor.findFieldByName("total_lines"), logCount)
                                                .setField(logSummaryResponseDescriptor.findFieldByName("service"), serviceName)
                                                .build();
                                        responseObserver.onNext(response);
                                        responseObserver.onCompleted();
                                        finishLatch.countDown();
                                    }
                                };
                            };
                        } else {
                            throw new IllegalArgumentException("未知方法:" + method.getMethodDescriptor().getFullMethodName());
                        }
                    }
                })
                .build();

        server.start();
        System.out.println("服务已启动,端口:8080");
        server.awaitTermination();
    }
}

第五步:实现客户端调用

最后写一个客户端来调用这两个流式方法,验证服务是否正常工作:

import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
import io.grpc.stub.StreamObserver;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

public class DynamicFileServiceClient {
    public static void main(String[] args) throws InterruptedException {
        ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 8080)
                .usePlaintext()
                .build();

        // 调用文件上传方法
        callUploadFile(channel);
        // 调用批量日志方法
        callBatchLogs(channel);

        channel.shutdown().awaitTermination(5, TimeUnit.SECONDS);
    }

    private static void callUploadFile(ManagedChannel channel) throws InterruptedException {
        CountDownLatch finishLatch = new CountDownLatch(1);

        StreamObserver<DynamicMessage> requestObserver = channel.newCall(uploadFileMethod, io.grpc.Metadata.newInstance())
                .start(new StreamObserver<DynamicMessage>() {
                    @Override
                    public void onNext(DynamicMessage response) {
                        boolean success = (boolean) response.getField(uploadResponseDescriptor.findFieldByName("success"));
                        String message = (String) response.getField(uploadResponseDescriptor.findFieldByName("message"));
                        System.out.println("上传响应:成功=" + success + ",消息=" + message);
                    }

                    @Override
                    public void onError(Throwable t) {
                        System.err.println("上传失败:" + t.getMessage());
                        finishLatch.countDown();
                    }

                    @Override
                    public void onCompleted() {
                        System.out.println("上传流程完成");
                        finishLatch.countDown();
                    }
                });

        // 发送5个文件块请求
        for (int i = 0; i < 5; i++) {
            DynamicMessage request = DynamicMessage.newBuilder(uploadRequestDescriptor)
                    .setField(uploadRequestDescriptor.findFieldByName("file_name"), "test.txt")
                    .setField(uploadRequestDescriptor.findFieldByName("file_chunk"), 
                              com.google.protobuf.ByteString.copyFrom(("Chunk " + i).getBytes()))
                    .build();
            requestObserver.onNext(request);
        }
        requestObserver.onCompleted();
        finishLatch.await(1, TimeUnit.MINUTES);
    }

    private static void callBatchLogs(ManagedChannel channel) throws InterruptedException {
        CountDownLatch finishLatch = new CountDownLatch(1);

        StreamObserver<DynamicMessage> requestObserver = channel.newCall(batchLogsMethod, io.grpc.Metadata.newInstance())
                .start(new StreamObserver<DynamicMessage>() {
                    @Override
                    public void onNext(DynamicMessage response) {
                        int totalLines = (int) response.getField(logSummaryResponseDescriptor.findFieldByName("total_lines"));
                        String service = (String) response.getField(logSummaryResponseDescriptor.findFieldByName("service"));
                        System.out.println("日志汇总:服务名=" + service + ",总日志数=" + totalLines);
                    }

                    @Override
                    public void onError(Throwable t) {
                        System.err.println("日志上报失败:" + t.getMessage());
                        finishLatch.countDown();
                    }

                    @Override
                    public void onCompleted() {
                        System.out.println("日志上报流程完成");
                        finishLatch.countDown();
                    }
                });

        // 发送10条日志请求
        for (int i = 0; i < 10; i++) {
            DynamicMessage request = DynamicMessage.newBuilder(logRequestDescriptor)
                    .setField(logRequestDescriptor.findFieldByName("service_name"), "user-service")
                    .setField(logRequestDescriptor.findFieldByName("log_line"), "INFO: 用户" + i + "登录成功")
                    .build();
            requestObserver.onNext(request);
        }
        requestObserver.onCompleted();
        finishLatch.await(1, TimeUnit.MINUTES);
    }
}

注意事项

  • 这种动态方式虽然灵活,但缺乏.proto文件带来的跨语言兼容性和代码生成的便利性,适合简单场景或动态服务需求。
  • 消息的字段类型、编号必须严格一致,否则会出现序列化/反序列化错误。
  • 复杂消息结构(比如嵌套消息、枚举)也可以用Descriptors来构建,但会比写.proto繁琐很多。

内容的提问来源于stack exchange,提问作者saikamesh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.09 00:47:39