如何让gRPC客户端保持监听模式接收服务器消息?
让C++ gRPC客户端保持监听服务器消息的解决方案
看起来你现在用的是一元RPC(也就是GetFeature那种一次请求一次响应的模式),这种模式下客户端拿到响应后就会结束连接,自然就退出了。要让客户端持续监听,咱们得换成流式RPC,下面一步步来实现:
1. 修改.proto文件定义流式RPC方法
首先得在你的Protocol Buffers定义里新增一个流式RPC接口,比如服务器流式(客户端发一次请求,服务器可以多次返回响应):
service RouteGuide { // 保留原来的一元RPC方法 rpc GetFeature(Point) returns (Feature) {} // 新增服务器流式RPC:客户端发一个Point请求,服务器持续返回Feature rpc ListenForFeatures(Point) returns (stream Feature) {} }
小贴士:
stream关键字放在returns后面,代表是服务器流式RPC;如果前后都加stream,就是双向流式RPC(双方都能持续发消息)。
2. 重新生成C++ gRPC代码
修改完.proto后,用protoc重新生成客户端和服务器的代码:
protoc --grpc_out=. --plugin=protoc-gen-grpc=`which grpc_cpp_plugin` your_proto_file.proto protoc --cpp_out=. your_proto_file.proto
3. 实现服务器端的流式方法
接下来在你的stream_server里实现新增的ListenForFeatures方法,让服务器可以持续发送消息:
Status RouteGuideImpl::ListenForFeatures(ServerContext* context, const Point* request, ServerWriter<Feature>* writer) { // 循环发送消息,直到客户端主动断开连接 while (!context->IsCancelled()) { Feature feature; // 这里可以从你的DB读取最新数据,或者生成测试内容 feature.set_name("Real-time Feature Update"); feature.mutable_location()->set_latitude(request->latitude()); feature.mutable_location()->set_longitude(request->longitude()); // 发送消息给客户端 writer->Write(feature); // 模拟间隔发送(比如每1秒发一次) std::this_thread::sleep_for(std::chrono::seconds(1)); } // 客户端断开后,返回成功状态 return Status::OK; }
这里context->IsCancelled()用来检测客户端是否断开,确保服务器不会一直死循环。
4. 修改客户端代码实现持续监听
最后修改你的stream_client,让它持续读取服务器发送的消息:
void RunContinuousListener(const std::string& server_addr) { // 创建gRPC通道和Stub std::unique_ptr<RouteGuide::Stub> stub = RouteGuide::NewStub( grpc::CreateChannel(server_addr, grpc::InsecureChannelCredentials()) ); ClientContext context; Point request; // 设置你要监听的位置(比如和之前GetFeature一样的坐标) request.set_latitude(407838); request.set_longitude(-746144); // 创建服务器流式调用的Reader std::unique_ptr<ClientReader<Feature>> reader(stub->ListenForFeatures(&context, request)); Feature response; // 循环读取服务器发送的每一条消息 while (reader->Read(&response)) { std::cout << "-------------- Received Update --------------" << std::endl; std::cout << "Found feature called " << response.name() << " at " << response.location().latitude()/1000000.0 << ", " << response.location().longitude()/1000000.0 << std::endl; } // 当流结束时,检查状态 Status status = reader->Finish(); if (status.ok()) { std::cout << "Listener stopped normally." << std::endl; } else { std::cout << "Listener failed: " << status.error_message() << std::endl; } } // 在main函数里调用这个方法,代替原来的GetFeature调用 int main(int argc, char** argv) { std::string server_address("0.0.0.0:50051"); RunContinuousListener(server_address); return 0; }
这样客户端就会一直处于监听状态,直到服务器关闭流或者连接出现错误。
额外说明
如果你需要客户端也能随时给服务器发消息(双向通信),可以把.proto里的方法改成双向流式:
rpc BidirectionalStream(stream Point) returns (stream Feature) {}
服务器端用ServerReaderWriter,客户端用ClientReaderWriter,双方都可以循环进行读写操作。
内容的提问来源于stack exchange,提问作者Vinay Shukla
相关产品推荐
相关产品推荐

