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

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;
}

关键修复说明

  1. 全局共享CompletionQueue:将CompletionQueue改为类成员变量,确保所有RPC调用使用同一个队列,异步处理线程监听该队列,保证响应能被正确接收。
  2. 正确处理流式RPC流程:使用ClientAsyncReader对应服务端的流式RPC,依次调用StartCall()、Read()、Finish(),通过不同tag标记不同异步事件。
  3. 调整响应检查时机:仅在Read()事件触发(tag=2)时读取并检查response,此时服务端数据已返回,能正确判断是否为换行符。
  4. 线程生命周期管理:析构函数中关闭队列并等待线程结束,避免资源泄漏。

服务端代码说明

服务端代码无需大改,run_stream方法中sys.stdin.read(1)会阻塞直到输入一个字符,输入换行符时会正确返回\n,符合需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 01:47:19