求助:gRPC异步服务器如何维持流打开并等待数据就绪
问题:gRPC异步服务器如何保持流打开以等待数据就绪?
我需要实现一个gRPC异步服务器,让Python客户端在调用next()从流中获取下一个元素时,如果数据耗尽且暂无新数据(比如数据正在生成),就进入阻塞状态。
Python客户端代码
my_service = MyServiceStub(channel) create_stream_request = MyService.MyServiceCreateStreamRequest() stream = my_service.CreateStream(create_stream_request) # 打开流对象 # 触发服务器生成数据(这段代码会在其他地方调用,需要通过编程方式通知服务器生成数据) produce_data_request = MyService.MyServiceProduceDataRequest() my_service.ProduceData(produce_data_request) # 从流中获取生成的数据 - 无数据时阻塞! new_data = stream.next() print(new_data)
C++服务器问题代码
我尝试编写了C++异步服务器代码,但遇到了问题:用condition_variable会导致整个gRPC处理线程挂起,无法处理后续请求;直接从Proceed方法返回会导致服务器崩溃,报错初始元数据无效。以下是我的服务器代码:
class StreamingManager { StreamingManager(grpc::ServerBuilder& builder) { // 构建并启动异步服务 builder.RegisterService(&grpcAsyncStreamingService_); completionQueue_ = builder.AddCompletionQueue(); grpcServer_ = builder.BuildAndStart(); // 在独立线程中启动HandleRpcs函数 std::thread(&StreamingManager::HandleRpcs, this).detach(); } ~StreamingManager() { shutdownRequested_ = true; grpcServer_->Shutdown(); // 优雅停止gRPC服务器 completionQueue_->Shutdown(); // 关闭完成队列以解除线程阻塞 } std::unique_ptr<grpc::Server> grpcServer_; std::unique_ptr<grpc::ServerCompletionQueue> completionQueue_; ::my::proto::MyService::AsyncService grpcAsyncStreamingService_; std::atomic<bool> shutdownRequested_; // 公共基类,保存请求状态和上下文,以及父类指针 enum CallStatus { CREATE, PROCESS, FINISH }; class CallDataBase { protected: CallStatus status_ = CREATE; grpc::ServerContext serverContext_; StreamingManager* streamingService_ = nullptr; public: virtual void Proceed() = 0; virtual ~CallDataBase() {} }; class CallDataProduceData : public CallDataBase { protected: ::my::proto::MyServiceProduceDataRequest request_; ::my::proto::MyServiceProduceDataResponse reply_; grpc::ServerAsyncResponseWriter<decltype(reply_)> responder_; // 非流式响应器 public: CallDataProduceData(StreamingManager* streamingService) : responder_(&serverContext_) { streamingService_ = streamingService; Proceed(); } public: void Proceed() override { if (status_ == CREATE) { status_ = PROCESS; // 原代码请求注册错误:应该是RequestProduceData streamingService_->grpcAsyncStreamingService_.RequestCreateStream(&serverContext_, &request_, &responder_, streamingService_->completionQueue_.get(), streamingService_->completionQueue_.get(), this); } else if (status_ == PROCESS) { new CallDataCreateStream(streamingService_); // 生成数据的处理代码开始 // 不知道如何保持流打开,无法处理此处逻辑 status_ = FINISH; responder_.Finish(reply_, ::grpc::Status::OK, this); } else { // 处理完成,自销毁 if (status_ != FINISH) { // 记录错误日志 } delete this; } } }; std::condition_variable produceDataCv_; std::mutex produceDataMutex_; std::atomic<bool> produceDataRequested_; class CallDataCreateStream : public CallDataBase { protected: ::my::proto::MyServiceCreateStreamRequest request_; ::my::proto::MyServiceCreateStreamResponse reply_; grpc::ServerAsyncWriter<decltype(reply_)> streaming_responder_; // 流式响应器 public: CallDataCreateStream(StreamingManager* streamingService) : streaming_responder_(&serverContext_) { streamingService_ = streamingService; Proceed(); } public: void Proceed() override { if (status_ == CREATE) { status_ = PROCESS; // 原代码请求注册错误:应该是RequestCreateStream streamingService_->grpcAsyncStreamingService_.RequestProduceData(&serverContext_, &request_, &streaming_responder_, streamingService_->completionQueue_.get(), streamingService_->completionQueue_.get(), this); } else if (status_ == PROCESS) { new CallDataCreateStream(streamingService_); // 创建流的处理代码开始 streaming_responder_.SendInitialMetadata(this); // 问题所在:如何保持流打开(不调用streaming_responder_.Finish()) // 同时等待新数据就绪?用condition_variable会让线程休眠,导致其他ProduceData请求无法处理 // 直接返回会崩溃,报错初始元数据无效 while (!streamingService_->shutdownRequested_) { reply_ = populate_reply_with_random_data(); streaming_responder_.Write(reply_, this); } status_ = FINISH; streaming_responder_.Finish(::grpc::Status::OK, this); } else { // 处理完成,自销毁 if (status_ != FINISH) { // 记录错误日志 } delete this; } } }; void HandleRpcs() { // 注册CreateStream和ProduceData处理器 new CallDataCreateStream(this); new CallDataProduceData(this); void* tag; bool ok; while (!shutdownRequested_) { bool ret = completionQueue_->Next(&tag, &ok); if (ok == false || ret == false) { return; } static_cast<CallDataBase*>(tag)->Proceed(); } } };
解决方案
核心思路
gRPC异步模型的核心是绝对不能阻塞Completion Queue的处理线程,所有等待操作必须通过异步回调或独立线程+队列实现,不能在Proceed方法内用循环或条件变量直接阻塞。
步骤1:维护活跃流实例
在StreamingManager中添加线程安全的容器,保存所有正在运行的流式请求实例,让ProduceData请求能找到对应的流并发送数据:
// StreamingManager成员变量 std::mutex activeStreamsMutex_; std::list<CallDataCreateStream*> activeStreams_;
步骤2:重构流式请求的状态机
将原来的循环等待改为状态驱动的异步回调,新增WAITING_FOR_DATA状态,避免阻塞处理线程:
// 扩展CallStatus枚举 enum CallStatus { CREATE, PROCESS, WAITING_FOR_DATA, FINISH }; class CallDataCreateStream : public CallDataBase { protected: ::my::proto::MyServiceCreateStreamRequest request_; ::my::proto::MyServiceCreateStreamResponse reply_; grpc::ServerAsyncWriter<decltype(reply_)> streaming_responder_; bool shouldFinish_ = false; public: CallDataCreateStream(StreamingManager* streamingService) : streaming_responder_(&serverContext_) { streamingService_ = streamingService; Proceed(); } void Proceed() override { if (status_ == CREATE) { status_ = PROCESS; // 修正请求注册:对应CreateStream方法 streamingService_->grpcAsyncStreamingService_.RequestCreateStream( &serverContext_, &request_, &streaming_responder_, streamingService_->completionQueue_.get(), streamingService_->completionQueue_.get(), this); } else if (status_ == PROCESS) { // 立即注册新的CreateStream处理器,避免漏接请求 new CallDataCreateStream(streamingService_); // 将当前流加入活跃容器 { std::lock_guard<std::mutex> lock(streamingService_->activeStreamsMutex_); streamingService_->activeStreams_.push_back(this); } // 发送初始元数据,完成后进入等待数据状态 streaming_responder_.SendInitialMetadata(this); status_ = WAITING_FOR_DATA; } else if (status_ == WAITING_FOR_DATA) { if (shouldFinish_) { // 关闭流 status_ = FINISH; streaming_responder_.Finish(::grpc::Status::OK, this); return; } // 无数据则直接返回,等待外部触发 if (!reply_.IsInitialized()) { return; } // 发送数据,完成后重置数据并回到等待状态 streaming_responder_.Write(reply_, this); reply_.Clear(); } else if (status_ == FINISH) { // 从活跃容器移除当前流 { std::lock_guard<std::mutex> lock(streamingService_->activeStreamsMutex_); streamingService_->activeStreams_.remove(this); } delete this; } } // 供外部调用,设置要发送的数据并触发回调 void SetData(const ::my::proto::MyServiceCreateStreamResponse& data) { reply_ = data; // 通知Completion Queue触发Proceed streamingService_->completionQueue_->Pluck(this); } // 关闭流 void FinishStream() { shouldFinish_ = true; streamingService_->completionQueue_->Pluck(this); } };
步骤3:修正ProduceData请求的处理逻辑
当收到ProduceData请求时,生成数据并分发到活跃流:
class CallDataProduceData : public CallDataBase { protected: ::my::proto::MyServiceProduceDataRequest request_; ::my::proto::MyServiceProduceDataResponse reply_; grpc::ServerAsyncResponseWriter<decltype(reply_)> responder_; public: CallDataProduceData(StreamingManager* streamingService) : responder_(&serverContext_) { streamingService_ = streamingService; Proceed(); } void Proceed() override { if (status_ == CREATE) { status_ = PROCESS; // 修正请求注册:对应ProduceData方法 streamingService_->grpcAsyncStreamingService_.RequestProduceData( &serverContext_, &request_, &responder_, streamingService_->completionQueue_.get(), streamingService_->completionQueue_.get(), this); } else if (status_ == PROCESS) { // 注册新的ProduceData处理器 new CallDataProduceData(streamingService_); // 生成数据 ::my::proto::MyServiceCreateStreamResponse data; data = populate_reply_with_random_data(); // 分发数据到所有活跃流 std::lock_guard<std::mutex> lock(streamingService_->activeStreamsMutex_); for (auto stream : streamingService_->activeStreams_) { stream->SetData(data); } // 完成ProduceData请求 status_ = FINISH; responder_.Finish(reply_, ::grpc::Status::OK, this); } else { delete this; } } };
关键注意事项
- 禁止阻塞Completion Queue线程:所有等待逻辑必须通过异步回调实现,否则服务器会失去响应。
- 修正请求注册错误:原代码中
CallDataCreateStream和CallDataProduceData的请求注册完全颠倒,必须对应正确的gRPC方法。 - 线程安全:活跃流容器的访问必须加锁,避免多线程竞争。
- 触发回调:通过
completionQueue_->Pluck(this)主动触发Proceed方法,确保数据就绪时能立即发送。
内容的提问来源于stack exchange,提问作者Dean
相关产品推荐
相关产品推荐

