如何用C++实现无状态通用gRPC代理?找不到对应API求解
C++ 实现服务无关gRPC代理的方案
核心思路
gRPC C++ 没有直接对应Java中ServerCallHandler、ServerCall.Listener这类的API,但可以基于异步通用服务API结合双向数据流组件实现服务无关的代理,核心是在代理与客户端、代理与后端服务器之间搭建透明的ByteBuffer级别的消息转发,无需解析具体的Protobuf内容。
实现步骤
- 注册通用异步服务
使用ServerBuilder::RegisterAsyncGenericService注册通用服务,这样代理可以接收任意服务、任意方法的gRPC请求,无需提前绑定具体的服务定义。 - 建立后端连接
针对每个客户端请求,通过请求上下文获取目标服务方法名,创建异步客户端通道连接到后端服务器,使用GenericStub发起通用调用。 - 双向数据流转发
利用ServerAsyncReaderWriter<grpc::ByteBuffer, grpc::ByteBuffer>处理客户端到代理的数据流,ClientAsyncReaderWriter<grpc::ByteBuffer, grpc::ByteBuffer>处理代理到后端的数据流,实现两端消息的直接转发。 - 生命周期管理
通过CompletionQueue处理异步事件(消息读写、请求完成等),在事件回调中完成状态流转,确保资源正确清理。
关键代码示例
#include <grpcpp/grpcpp.h> #include <grpcpp/generic/generic_stub.h> #include <memory> #include <iostream> class ProxyCallData { public: ProxyCallData(grpc::GenericService* service, grpc::ServerCompletionQueue* cq, std::shared_ptr<grpc::Channel> backend_channel) : service_(service), cq_(cq), responder_(&client_ctx_), backend_stub_(backend_channel), status_(CREATE) { Proceed(); } void Proceed() { switch (status_) { case CREATE: { status_ = PROCESS; // 注册接收客户端的通用请求 service_->RequestCall(&client_ctx_, &responder_, cq_, cq_, this); break; } case PROCESS: { // 提前创建新的CallData处理下一个请求,避免阻塞 new ProxyCallData(service_, cq_, backend_stub_.get_channel()); // 从客户端上下文获取请求方法名,构造后端调用 const std::string& method = client_ctx_.method(); backend_call_ = backend_stub_.PrepareCall(&backend_ctx_, method, cq_); backend_call_->StartCall(&backend_writer_, this); // 开始读取客户端发送的消息 responder_.Read(&client_buffer_, this); status_ = FORWARD_TO_BACKEND; break; } case FORWARD_TO_BACKEND: { if (client_buffer_.Valid()) { // 将客户端消息转发到后端 backend_writer_.Write(client_buffer_, this); // 继续读取下一条客户端消息 responder_.Read(&client_buffer_, this); } else { // 客户端发送完毕,通知后端 backend_writer_.WritesDone(this); status_ = FORWARD_TO_CLIENT; // 开始读取后端返回的消息 backend_writer_.Read(&backend_buffer_, this); } break; } case FORWARD_TO_CLIENT: { if (backend_buffer_.Valid()) { // 将后端消息转发给客户端 responder_.Write(backend_buffer_, this); // 继续读取下一条后端消息 backend_writer_.Read(&backend_buffer_, this); } else { // 后端返回完毕,结束客户端请求 responder_.WritesDone(this); status_ = FINISH; responder_.Finish(grpc::Status::OK, this); } break; } case FINISH: { // 清理当前CallData资源 delete this; break; } } } private: enum CallStatus { CREATE, PROCESS, FORWARD_TO_BACKEND, FORWARD_TO_CLIENT, FINISH }; grpc::GenericService* service_; grpc::ServerCompletionQueue* cq_; grpc::GenericServerContext client_ctx_; grpc::ServerAsyncReaderWriter<grpc::ByteBuffer, grpc::ByteBuffer> responder_; grpc::ByteBuffer client_buffer_; grpc::GenericStub backend_stub_; grpc::ClientContext backend_ctx_; std::unique_ptr<grpc::ClientAsyncReaderWriter<grpc::ByteBuffer, grpc::ByteBuffer>> backend_call_; grpc::ClientAsyncReaderWriter<grpc::ByteBuffer, grpc::ByteBuffer>& backend_writer_ = *backend_call_; grpc::ByteBuffer backend_buffer_; CallStatus status_; }; void RunProxy(const std::string& proxy_addr, const std::string& backend_addr) { grpc::ServerBuilder builder; builder.AddListeningPort(proxy_addr, grpc::InsecureServerCredentials()); auto generic_service = std::make_unique<grpc::GenericService>(); builder.RegisterAsyncGenericService(generic_service.get()); auto cq = builder.AddCompletionQueue(); auto server = builder.BuildAndStart(); std::cout << "Proxy listening on " << proxy_addr << std::endl; // 创建后端通道 auto backend_channel = grpc::CreateChannel(backend_addr, grpc::InsecureChannelCredentials()); // 启动事件循环处理请求 new ProxyCallData(generic_service.get(), cq.get(), backend_channel); void* tag; bool ok; while (true) { GPR_ASSERT(cq->Next(&tag, &ok)); GPR_ASSERT(ok); static_cast<ProxyCallData*>(tag)->Proceed(); } } int main(int argc, char** argv) { RunProxy("0.0.0.0:50051", "localhost:50052"); return 0; }
注意要点
- 元数据转发:需要将客户端请求的元数据(如Headers、Trailers)同步到后端请求上下文,同时将后端返回的元数据回传给客户端,可通过
client_ctx_.client_metadata()和backend_ctx_.AddMetadata()实现。 - 错误处理:需处理后端连接失败、请求超时等异常场景,及时向客户端返回对应的gRPC错误状态。
- 并发优化:可通过多线程处理多个CompletionQueue,提升代理的并发处理能力,避免单线程瓶颈。
内容的提问来源于stack exchange,提问作者Anil Kumar
相关产品推荐
相关产品推荐

