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

如何用C++实现支持异步多RPC(含双向流)的gRPC服务器

C++ gRPC异步服务器实现(含服务器流与双向流)

核心逻辑说明

gRPC异步服务器完全依赖Completion Queue(CQ) 驱动所有异步操作,每个RPC调用的生命周期需要通过状态机式的异步步骤完成。不同类型RPC的处理流程差异主要体现在数据交互的阶段:

服务器流RPC(Get接口)

和你参考的官方一元RPC异步逻辑类似,流程为:

  1. 注册Get方法的请求监听,等待客户端发起调用
  2. 接收到GetRequest后,异步批量发送GetResponse
  3. 所有响应发送完毕后,结束RPC并回收资源

双向流RPC(Modify接口)

双向流需要同时处理客户端流式请求接收和服务器流式响应发送,核心是通过CQ标签(Tag)区分不同操作阶段:

  1. 注册Modify方法的流请求监听
  2. 接收到客户端请求后,立即发起下一次请求的异步读取(保持流的持续监听)
  3. 处理当前请求,异步发送对应响应
  4. 重复步骤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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 14:52:53