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

如何用gRPC C++从描述符集运行时创建流式RPC服务?

问题描述

我正在为基于gRPC通信的C++应用开发测试驱动,多数测试驱动用grpcurl向被测应用发送消息并验证响应。但部分应用会连接流式RPC,虽然编写专用测试驱动不难,但我想实现一个通用工具:该工具可接收描述符集(通过protoc的--descriptor_set_out参数生成)、要提供服务的流式方法名,以及定义流式RPC连接时返回消息的JSON文件。

目前我已完成描述符集解析、获取流式方法的服务和方法描述符、从JSON加载返回消息的步骤,但卡在从描述符创建服务的环节。以下是我的概念验证代码(无错误检查、路径硬编码,仅用于验证可行性):

#include "google/protobuf/descriptor.pb.h"
#include "google/protobuf/dynamic_message.h"
#include "google/protobuf/util/json_util.h"

#include <fstream>
#include <sstream>

int main(int argc, char** argv)
{
    google::protobuf::FileDescriptorSet desc;
    std::stringstream sst;
    {
        std::ifstream i("/tmp/test.protoset");
        sst << i.rdbuf();
    }
    desc.ParseFromString(sst.str());
    google::protobuf::DescriptorPool desc_pool;

    for (const auto& fdesc : desc.file())
    {
        desc_pool.BuildFile(fdesc);
    }

    auto sdesc = desc_pool.FindServiceByName("TestService");
    auto mdesc = sdesc->FindMethodByName("connect");
    auto resp_type = mdesc->output_type();
    google::protobuf::DynamicMessageFactory dmf(&desc_pool);
    sst.str("");
    sst.clear();
    auto out_message = std::shared_ptr<google::protobuf::Message>(dmf.GetPrototype(resp_type)->New());
    {
        std::ifstream i("/tmp/test_message.json");
        sst << i.rdbuf();
    }
    auto stat = google::protobuf::util::JsonStringToMessage(sst.str(), out_message.get());
    std::cout << "READ " << stat << " " << out_message->DebugString() << std::endl;
}

请问现在是否可以创建TestService/connect流式RPC,等待连接并返回out_message中构建的消息?

解决方案

可以实现,你需要基于gRPC的动态服务框架搭建流式RPC服务,核心是利用ServerBuilder配合动态生成的消息和通用服务处理逻辑,具体步骤如下:

  1. 引入gRPC核心头文件
    在原有代码基础上添加gRPC相关依赖头文件:

    #include <grpcpp/grpcpp.h>
    #include <grpcpp/server_builder.h>
    
  2. 实现通用流式RPC服务类
    继承grpc::Service,实现动态绑定方法和流式请求处理逻辑(以下示例针对服务器端流式RPC,即客户端发单请求、服务器返回流式响应;若为双向流式或客户端流式,可调整读写逻辑):

    class DynamicStreamingServiceImpl : public grpc::Service {
    public:
        DynamicStreamingServiceImpl(const google::protobuf::ServiceDescriptor* sdesc,
                                    const google::protobuf::MethodDescriptor* mdesc,
                                    std::shared_ptr<google::protobuf::Message> resp_msg)
            : service_desc_(sdesc), method_desc_(mdesc), response_msg_(std::move(resp_msg)) {
            // 向gRPC服务注册目标流式方法
            AddMethod(new grpc::GenericAsyncService::GenericMethod(
                method_desc_->full_name(),
                grpc::RpcMethod::SERVER_STREAMING,
                nullptr));
        }
    
        // 通用流式请求处理逻辑
        grpc::Status HandleStreamingCall(grpc::ServerContext* context,
                                         grpc::GenericServerAsyncReaderWriter* reader_writer) override {
            // 将预加载的响应消息序列化为ByteBuffer
            grpc::ByteBuffer resp_buf;
            grpc::Status serialize_status = grpc::SerializationTraits<google::protobuf::Message>::Serialize(
                *response_msg_, &resp_buf);
            if (!serialize_status.ok()) {
                return serialize_status;
            }
    
            // 发送响应消息
            grpc::Status write_status = reader_writer->Write(resp_buf, grpc::WriteOptions());
            if (!write_status.ok()) {
                return write_status;
            }
    
            // 结束流式响应
            return reader_writer->Finish(grpc::Status::OK, nullptr);
        }
    
    private:
        const google::protobuf::ServiceDescriptor* service_desc_;
        const google::protobuf::MethodDescriptor* method_desc_;
        std::shared_ptr<google::protobuf::Message> response_msg_;
    };
    
  3. 搭建并启动gRPC服务器
    在main函数末尾添加服务器初始化和启动逻辑:

    int main(int argc, char** argv) {
        // 原有的描述符解析、消息加载代码保持不变...
    
        // 创建动态流式服务实例
        DynamicStreamingServiceImpl dynamic_service(sdesc, mdesc, out_message);
    
        // 配置并启动gRPC服务器
        std::string server_addr("0.0.0.0:50051");
        grpc::ServerBuilder builder;
        builder.AddListeningPort(server_addr, grpc::InsecureServerCredentials());
        builder.RegisterService(&dynamic_service);
    
        std::unique_ptr<grpc::Server> server(builder.BuildAndStart());
        std::cout << "Server listening on " << server_addr << std::endl;
    
        // 阻塞等待服务器终止
        server->Wait();
    
        return 0;
    }
    
  4. 关键注意事项

    • 区分流式RPC类型:服务器端流式、客户端流式、双向流式的处理逻辑有差异,需根据实际方法类型调整读写逻辑。
    • 编译时需链接gRPC和Protobuf相关库(如grpc++、protobuf)。
    • 建议将硬编码的文件路径改为命令行参数,提升工具通用性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 21:40:35