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

关闭grpc::CompletionQueue触发断言的原因及优雅退出方案

如何优雅退出监听grpc::Channel状态变化的CompletionQueue线程?

我之前在Stack Overflow提问过如何解除grpc::Channel::NotifyOnStateChange(无限超时)导致的grpc::CompletionQueue::Next()阻塞,但问题未解决。后来改用带非无限截止时间的NotifyOnStateChange作为替代方案,代码如下:

// main.cpp
#include <chrono>
#include <iostream>
#include <memory>
#include <thread>
#include <atomic>
#include <grpcpp/grpcpp.h>
#include <unistd.h>

using namespace std;
using namespace grpc;

void threadFunc(shared_ptr<Channel> ch, CompletionQueue* cq, atomic<bool>& stop_flag) {
  void* tag = NULL;
  bool ok = false;
  int i = 1;
  grpc_connectivity_state state = ch->GetState(false);
  std::chrono::time_point<std::chrono::system_clock> now =
        std::chrono::system_clock::now();
  std::chrono::time_point<std::chrono::system_clock> deadline =
    now + std::chrono::seconds(2);

  cout << "state " << i++ << " = " << (int)state << endl;
  // 首次注册状态监听前检查退出标志
  if (!stop_flag) {
    ch->NotifyOnStateChange(state,
                            deadline,
                            cq,
                            (void*)1);
  }

  while (cq->Next(&tag, &ok)) {
    if (stop_flag) {
      break;
    }
    state = ch->GetState(false);
    cout << "state " << i++ << " = " << (int)state << endl;
    now = std::chrono::system_clock::now();
    deadline = now + std::chrono::seconds(2);
    // 注册新监听前检查退出标志
    if (!stop_flag) {
      ch->NotifyOnStateChange(state,
                              deadline,
                              cq,
                              (void*)1);
    }
  }

  cout << "thread end" << endl;
}

int main(int argc, char* argv[]) {
  ChannelArguments channel_args;
  CompletionQueue cq;
  atomic<bool> stop_flag{false};

  channel_args.SetInt(GRPC_ARG_HTTP2_MAX_PINGS_WITHOUT_DATA, 0);
  channel_args.SetInt(GRPC_ARG_MIN_RECONNECT_BACKOFF_MS, 2000);
  channel_args.SetInt(GRPC_ARG_MAX_RECONNECT_BACKOFF_MS, 2000);
  channel_args.SetInt(GRPC_ARG_HTTP2_BDP_PROBE, 0);
  channel_args.SetInt(GRPC_ARG_KEEPALIVE_TIME_MS, 60000);
  channel_args.SetInt(GRPC_ARG_KEEPALIVE_TIMEOUT_MS, 30000);
  channel_args.SetInt(GRPC_ARG_HTTP2_MIN_SENT_PING_INTERVAL_WITHOUT_DATA_MS,
                      60000);

  {
    shared_ptr<Channel> ch(CreateCustomChannel("my_grpc_server:50051",
                                               InsecureChannelCredentials(),
                                               channel_args));
    std::thread my_thread(&threadFunc, ch, &cq, ref(stop_flag));
    cout << "sleeping" << endl;
    sleep(5);
    cout << "slept" << endl;
    // 先标记线程退出,再关闭队列
    stop_flag = true;
    cq.Shutdown();
    cout << "shut down cq" << endl;
    my_thread.join();
  }
}

原程序运行输出如下:

$ ./a.out
sleeping
state 1 = 0
state 2 = 0
state 3 = 0
slept
shut down cq
state 4 = 0
E1012 15:29:07.677225824      54 channel_connectivity.cc:234] assertion failed: grpc_cq_begin_op(cq, tag)
Aborted (core dumped)

问题原因

断言触发的核心原因是:调用cq.Shutdown()之后,线程仍在执行ch->NotifyOnStateChange尝试向已关闭的CompletionQueue添加新操作,违反了gRPC的断言检查。

解决思路

  1. 新增原子类型的线程退出标志std::atomic<bool>,用于主线程向监听线程发送退出信号。
  2. 主线程在调用cq.Shutdown()前先设置退出标志,让监听线程提前停止注册新的状态监听操作。
  3. 监听线程每次准备注册NotifyOnStateChange前,先检查退出标志,若已标记退出则不再执行注册逻辑。

修改后效果

修改后的程序运行时,主线程会先通知监听线程停止注册新操作,再关闭CompletionQueue,线程会在处理完最后一次Next()返回后优雅退出,不会触发断言崩溃。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 14:20:36