如何用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配合动态生成的消息和通用服务处理逻辑,具体步骤如下:
引入gRPC核心头文件
在原有代码基础上添加gRPC相关依赖头文件:#include <grpcpp/grpcpp.h> #include <grpcpp/server_builder.h>实现通用流式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_; };搭建并启动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; }关键注意事项
- 区分流式RPC类型:服务器端流式、客户端流式、双向流式的处理逻辑有差异,需根据实际方法类型调整读写逻辑。
- 编译时需链接gRPC和Protobuf相关库(如
grpc++、protobuf)。 - 建议将硬编码的文件路径改为命令行参数,提升工具通用性。
内容的提问来源于stack exchange,提问作者Kevin
相关产品推荐
相关产品推荐

