如何用C++实现支持异步多RPC(含双向流)的gRPC服务器
C++ gRPC异步服务器实现(含服务器流与双向流)
核心逻辑说明
gRPC异步服务器完全依赖Completion Queue(CQ) 驱动所有异步操作,每个RPC调用的生命周期需要通过状态机式的异步步骤完成。不同类型RPC的处理流程差异主要体现在数据交互的阶段:
服务器流RPC(Get接口)
和你参考的官方一元RPC异步逻辑类似,流程为:
- 注册
Get方法的请求监听,等待客户端发起调用 - 接收到
GetRequest后,异步批量发送GetResponse - 所有响应发送完毕后,结束RPC并回收资源
双向流RPC(Modify接口)
双向流需要同时处理客户端流式请求接收和服务器流式响应发送,核心是通过CQ标签(Tag)区分不同操作阶段:
- 注册
Modify方法的流请求监听 - 接收到客户端请求后,立即发起下一次请求的异步读取(保持流的持续监听)
- 处理当前请求,异步发送对应响应
- 重复步骤2-3,直到客户端关闭请求流,服务器发送完所有响应后结束RPC
关键代码实现示例
1. RPC处理类(封装状态与逻辑)
#include <grpcpp/grpcpp.h> #include "test.grpc.pb.h" using grpc::Server; using grpc::ServerAsyncResponseWriter; using grpc::ServerBuilder; using grpc::ServerCompletionQueue; using grpc::ServerContext; using grpc::Status; using test::test; using test::GetRequest; using test::GetResponse; using test::ModifyRequest; using test::ModifyResponse; // 通用RPC标签基类,用于CQ回调区分RPC类型 class AsyncRpcTag { public: enum class Type { GET, MODIFY }; AsyncRpcTag(Type type) : type_(type) {} virtual ~AsyncRpcTag() = default; virtual void Process(bool ok) = 0; Type type() const { return type_; } private: Type type_; }; // 服务器流Get RPC处理类 class GetRpc : public AsyncRpcTag { public: GetRpc(test::AsyncService* service, ServerCompletionQueue* cq) : AsyncRpcTag(Type::GET), service_(service), cq_(cq), responder_(&ctx_), status_(CREATE) { Proceed(true); } void Proceed(bool ok) override { if (status_ == CREATE) { // 注册Get方法的请求监听 status_ = PROCESS; service_->RequestGet(&ctx_, &request_, &responder_, cq_, cq_, this); } else if (status_ == PROCESS) { // 立即创建新实例处理下一个Get请求,支持并发 new GetRpc(service_, cq_); // 模拟生成3条流式响应 for (int i = 0; i < 3; ++i) { GetResponse response; response.set_getresponsemessage("Response for " + request_.getrequestmessage() + " - " + std::to_string(i)); responder_.Write(response, this); } // 发送结束状态,完成RPC responder_.Finish(Status::OK, this); status_ = FINISH; } else { // 回收当前实例资源 delete this; } } private: test::AsyncService* service_; ServerCompletionQueue* cq_; ServerContext ctx_; GetRequest request_; ServerAsyncResponseWriter<GetResponse> responder_; enum CallStatus { CREATE, PROCESS, FINISH }; CallStatus status_; }; // 双向流Modify RPC处理类 class ModifyRpc : public AsyncRpcTag { public: ModifyRpc(test::AsyncService* service, ServerCompletionQueue* cq) : AsyncRpcTag(Type::MODIFY), service_(service), cq_(cq), responder_(&ctx_), status_(CREATE) { Proceed(true); } void Proceed(bool ok) override { if (status_ == CREATE) { // 注册Modify方法的流请求监听 status_ = READ; service_->RequestModify(&ctx_, &stream_, &responder_, cq_, cq_, this); // 发起第一次请求的异步读取 stream_.Read(&request_, this); } else if (status_ == READ) { if (!ok) { // 客户端关闭请求流,结束RPC responder_.Finish(Status::OK, this); status_ = FINISH; return; } // 立即注册下一次请求读取,避免丢失客户端后续数据 stream_.Read(&request_, this); // 处理请求并异步发送响应 ModifyResponse response; response.set_modifyresponsemessage("Processed: " + request_.modifyrequestmessage()); responder_.Write(response, this); } else if (status_ == FINISH) { // 回收当前实例资源 delete this; } } private: test::AsyncService* service_; ServerCompletionQueue* cq_; ServerContext ctx_; grpc::ServerAsyncReaderWriter<ModifyResponse, ModifyRequest> stream_; ModifyRequest request_; enum CallStatus { CREATE, READ, FINISH }; CallStatus status_; };
2. 服务器主运行逻辑
class AsyncServer { public: ~AsyncServer() { server_->Shutdown(); // 关闭CQ前处理完所有待执行操作 cq_->Shutdown(); } void Run(const std::string& server_address) { ServerBuilder builder; builder.AddListeningPort(server_address, grpc::InsecureServerCredentials()); builder.RegisterAsyncService(&service_); cq_ = builder.AddCompletionQueue(); server_ = builder.BuildAndStart(); std::cout << "Server listening on " << server_address << std::endl; // 初始化两种RPC的监听实例 new GetRpc(&service_, cq_.get()); new ModifyRpc(&service_, cq_.get()); // 驱动CQ处理所有异步事件 void* tag; bool ok; while (cq_->Next(&tag, &ok)) { static_cast<AsyncRpcTag*>(tag)->Process(ok); } } private: std::unique_ptr<ServerCompletionQueue> cq_; test::AsyncService service_; std::unique_ptr<Server> server_; }; int main(int argc, char** argv) { AsyncServer server; server.Run("0.0.0.0:50051"); return 0; }
关键细节解析
- CQ标签(Tag):每个异步操作绑定一个Tag实例,CQ触发时通过Tag定位到对应RPC的处理逻辑,驱动状态机流转。
- 双向流并发读写:必须在接收到当前请求后立即发起下一次
Read,否则会丢失客户端后续发送的数据;Write操作可异步执行,无需等待前一个响应发送完成。 - 资源自动回收:每个RPC处理实例在完成所有操作后自行销毁,避免内存泄漏。
官方示例扩展说明
你参考的官方异步示例是一元RPC的基础实现,核心是通过CQ驱动请求接收与响应发送。扩展到流式RPC只需替换对应的处理类:
- 服务器流使用
ServerAsyncResponseWriter - 双向流使用
ServerAsyncReaderWriter - 每种RPC类型需单独实现状态机,维护自身的生命周期阶段
内容的提问来源于stack exchange,提问作者user3159857
相关产品推荐
相关产品推荐

