Grpc异步C++客户端问题:无法识别服务端输入信号
gRPC异步客户端无法接收服务端信号问题排查与修复
问题描述
开发基于gRPC的C++异步客户端,需求为:客户端每秒打印一个点“.”,接收到服务端发送的换行符“\n”信号时,打印收到信号的提示;未收到信号时持续循环打印点。当前客户端能正常打印点,但无法识别服务端发送的信号。
错误原因分析
原客户端代码存在以下核心问题:
- CompletionQueue 生命周期与关联错误:
run()方法中每次循环新建CompletionQueue,且AsyncCompleteRpc()里的队列是独立局部变量,两个队列无关联,导致RPC响应无法被正确处理。 - 异步流程顺序错误:调用
StartCall()后直接读取response,此时服务端还未返回数据,response为空;且调用Finish()后未等待队列处理完成就尝试判断响应,完全违背异步调用逻辑。 - RPC类型不匹配:服务端
run_stream是服务器端流式RPC,客户端应使用ClientAsyncReader而非ClientAsyncResponseReader,原代码的异步调用方式不符合流式RPC处理流程。
修复后的客户端代码
#include <stdlib.h> #include <grpcpp/grpcpp.h> #include <string> #include <unistd.h> #include <thread> #include <iostream> #include "dts.grpc.pb.h" using grpc::Channel; using grpc::ClientContext; using grpc::ClientAsyncReader; using grpc::Status; using grpc::CompletionQueue; class DTSClient{ public: explicit DTSClient(std::string socket); ~DTSClient() { cq_.Shutdown(); if(rpc_thread_.joinable()){ rpc_thread_.join(); } } void run() { while(true){ unsigned int microsecond = 1000000; usleep(1 * microsecond); // 每秒打印一个点 std::cout << "." << std::flush; // 发起服务器流式RPC调用 ClientContext context; EntryCreateRequest request; EntryResponse response; Status status; // 创建异步Reader,关联全局的CompletionQueue auto reader = stub_->PrepareAsyncrun_stream(&context, request, &cq_); reader->StartCall((void*)1); // 标记调用启动的tag // 请求读取响应,tag标记为2 reader->Read(&response, (void*)2); // 请求完成,tag标记为3 reader->Finish(&status, (void*)3); } } private: void AsyncCompleteRpc() { void* got_tag; bool ok = false; EntryResponse response; while (cq_.Next(&got_tag, &ok)) { if (!ok) { std::cout << "\nRPC调用失败" << std::endl; continue; } switch(reinterpret_cast<intptr_t>(got_tag)){ case 1: // RPC调用启动完成,无需额外操作 break; case 2: // 读取到服务端响应 if(!response.title().empty() && response.title()[0] == '\n'){ std::cout << "\n已收到服务端换行信号" << std::endl; } break; case 3: // RPC完成,检查状态 if(!status.ok()){ std::cout << "\nRPC状态异常: " << status.error_message() << std::endl; } break; default: std::cout << "\n未知tag: " << reinterpret_cast<intptr_t>(got_tag) << std::endl; break; } } } std::shared_ptr<grpc::Channel> channel_; std::unique_ptr<DTSService::Stub> stub_; CompletionQueue cq_; // 全局唯一的CompletionQueue std::thread rpc_thread_; // 处理队列的线程 }; DTSClient::DTSClient(std::string socket_file) { std::string full_socket("unix://"); full_socket += socket_file; channel_ = grpc::CreateChannel(full_socket, grpc::InsecureChannelCredentials()); stub_ = DTSService::NewStub(channel_); // 启动异步处理线程 rpc_thread_ = std::thread(&DTSClient::AsyncCompleteRpc, this); } int main(int argc, char** argv) { DTSClient dts_client("/tmp/grpc.sock"); std::cout << "Client start" << std::endl; dts_client.run(); return 0; }
关键修复说明
- 全局共享CompletionQueue:将
CompletionQueue改为类成员变量,确保所有RPC调用使用同一个队列,异步处理线程监听该队列,保证响应能被正确接收。 - 正确处理流式RPC流程:使用
ClientAsyncReader对应服务端的流式RPC,依次调用StartCall()、Read()、Finish(),通过不同tag标记不同异步事件。 - 调整响应检查时机:仅在
Read()事件触发(tag=2)时读取并检查response,此时服务端数据已返回,能正确判断是否为换行符。 - 线程生命周期管理:析构函数中关闭队列并等待线程结束,避免资源泄漏。
服务端代码说明
服务端代码无需大改,run_stream方法中sys.stdin.read(1)会阻塞直到输入一个字符,输入换行符时会正确返回\n,符合需求。
内容的提问来源于stack exchange,提问作者Ron Domnitz
相关产品推荐
相关产品推荐

