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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 18:32:10