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

求助: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;
        }
    }
};

关键注意事项

  1. 禁止阻塞Completion Queue线程:所有等待逻辑必须通过异步回调实现,否则服务器会失去响应。
  2. 修正请求注册错误:原代码中CallDataCreateStream和CallDataProduceData的请求注册完全颠倒,必须对应正确的gRPC方法。
  3. 线程安全:活跃流容器的访问必须加锁,避免多线程竞争。
  4. 触发回调:通过completionQueue_->Pluck(this)主动触发Proceed方法,确保数据就绪时能立即发送。

内容的提问来源于stack exchange,提问作者Dean

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 07:40:55