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

外部反应器流场景下gRPC Reactor的取消与生命周期管理问询

gRPC流式RPC取消与生命周期管理问题

我正在编写gRPC API,需实现用户发送RPC启动测试、流式返回测试进度直至结束的功能。核心要求是用户取消RPC后,测试仍需继续运行以便后续查看结果。我已采用callback API与outside-reactor flow实现该功能,伪代码如下:

实现代码

Reactor类定义

class Reactor : public grpc::ServerWriteReactor<protos::TestResults> {
public:
  explicit Reactor(const std::vector<std::string> &testNames) {
    // 省略:初始化m_response
  }

  void OnWriteDone(bool ok) override {
  }

  void OnDone() override {
    delete this;
  }

  void OnCancel() override {
    m_cancelled = true;
  }

  void OnTestStatusChange(const std::string &testName,
                          const TestStatus newStatus) override {
    auto guard = std::lock_guard{m_lock};

    if (m_cancelled) {
      return;
    }

    // 省略:找到对应测试并更新状态

    StartWrite(&m_response);
  }

  void OnAllTestsFinish() override {
    auto guard = std::lock_guard{m_lock};

    if (m_cancelled) {
      Finish(grpc::Status{grpc::StatusCode::CANCELLED, "Client cancelled RPC"});
    } else {
      Finish(grpc::Status::OK);
    }
  }

private:
  protos::TestResults m_response;
  std::mutex m_lock;
  bool m_cancelled{false};
};

服务方法实现

grpc::ServerWriteReactor<protos::TestResults> *
TestRunnerService::RunTests(grpc::CallbackServerContext * /*context*/,
                            const google::protobuf::Empty * /*request*/) {
  auto testNames = std::vector<std::string>{"test_1", "test_2", "test_3"};
  auto reactor = new Reactor(testNames);

  // 后续应改为成员变量存储线程
  auto t = std::thread(
      [reactor, testNames]() { TestController::RunTests(reactor, testNames); });
  t.detach();

  return reactor;
}

测试控制器逻辑

void TestController::RunTests(Reactor *subscriber,
                              const std::vector<std::string> &testNames) {
  for (const auto &name : testNames) {
     std::cout << "Running test: " << name << "\n";

     subscriber->OnTestStatusChange(name, TestStatus::RUNNING);
     std::this_thread::sleep_for(std::chrono::seconds(2));

     std::cout << "Test passed: " << name << "\n";

     subscriber->OnTestStatusChange(name, TestStatus::PASSED);
     std::this_thread::sleep_for(std::chrono::seconds(2));
  }

  subscriber->OnAllTestsFinish();
  std::cout << "Tests completed\n";
}

当前疑问

  1. 当前在OnCancel中设置标记,使后续写入无操作,待测试全部结束后以取消状态终止RPC,这是否会导致流长期处于开启状态,属于不良实践?
  2. 若需在OnCancel中直接调用Finish(),该如何管理Reactor的生命周期?若保持OnDone()中delete this的逻辑,TestController中的方法调用Reactor会触发未定义行为。我曾考虑用shared_ptr<Reactor>存储并定期检查m_done标记,但想了解更惯用的处理方式。

问题解答

针对第一个疑问

这种做法属于不良实践,原因如下:

  • gRPC服务器的连接、内存等资源有限,长期维持已取消的流会占用资源,高并发场景下可能引发资源耗尽问题。
  • 客户端已取消RPC,服务器端继续维持流没有实际业务意义,反而增加状态维护成本。
  • 符合gRPC设计语义的做法是:客户端取消RPC后,服务器端应尽快终止对应流,释放相关资源。

针对第二个疑问

推荐采用智能指针+原子状态标记的组合方案来管理Reactor生命周期,既满足测试继续运行的需求,又避免悬空指针问题,具体实现思路如下:

1. 改造Reactor类

让Reactor继承std::enable_shared_from_this,添加原子状态标记m_finished,调整取消与完成逻辑:

class Reactor : public grpc::ServerWriteReactor<protos::TestResults>,
                public std::enable_shared_from_this<Reactor> {
public:
  explicit Reactor(const std::vector<std::string> &testNames) {
    // 初始化逻辑
  }

  void OnWriteDone(bool ok) override {
  }

  void OnDone() override {
    m_finished.store(true, std::memory_order_relaxed);
    // 不再手动delete this,由shared_ptr自动管理生命周期
  }

  void OnCancel() override {
    auto guard = std::lock_guard{m_lock};
    if (m_finished.load(std::memory_order_relaxed)) {
      return;
    }
    m_cancelled = true;
    m_finished.store(true, std::memory_order_relaxed);
    Finish(grpc::Status{grpc::StatusCode::CANCELLED, "Client cancelled RPC"});
  }

  void OnTestStatusChange(const std::string &testName,
                          const TestStatus newStatus) override {
    if (m_finished.load(std::memory_order_relaxed)) {
      return;
    }
    auto guard = std::lock_guard{m_lock};
    if (m_cancelled) {
      return;
    }
    // 更新测试状态并写入流
    StartWrite(&m_response);
  }

  void OnAllTestsFinish() override {
    if (m_finished.load(std::memory_order_relaxed)) {
      return;
    }
    auto guard = std::lock_guard{m_lock};
    if (m_cancelled) {
      return; // 已在OnCancel中处理过Finish
    }
    m_finished.store(true, std::memory_order_relaxed);
    Finish(grpc::Status::OK);
  }

  // 暴露状态供测试线程检查(可选)
  bool IsFinished() const {
    return m_finished.load(std::memory_order_relaxed);
  }

private:
  protos::TestResults m_response;
  std::mutex m_lock;
  bool m_cancelled{false};
  std::atomic<bool> m_finished{false};
};

2. 调整服务方法与测试线程

使用shared_ptr创建Reactor实例,传递给测试线程确保生命周期安全:

grpc::ServerWriteReactor<protos::TestResults> *
TestRunnerService::RunTests(grpc::CallbackServerContext * /*context*/,
                            const google::protobuf::Empty * /*request*/) {
  auto testNames = std::vector<std::string>{"test_1", "test_2", "test_3"};
  auto reactor = std::make_shared<Reactor>(testNames);
  
  // 传递shared_ptr到测试线程,保证测试过程中Reactor不会被销毁
  auto t = std::thread(
      [reactor, testNames]() { TestController::RunTests(reactor, testNames); });
  t.detach();
  
  // 返回裸指针给gRPC框架,框架通过OnDone触发后续处理,实际生命周期由shared_ptr管理
  return reactor.get();
}

3. 测试线程逻辑优化

测试线程持有shared_ptr,每次回调前检查Reactor状态:

void TestController::RunTests(std::shared_ptr<Reactor> subscriber,
                              const std::vector<std::string> &testNames) {
  for (const auto &name : testNames) {
     // 可选:如果不需要强制测试跑完,可在此提前终止;根据需求测试需继续运行则跳过此判断
     // if (subscriber->IsFinished()) break;
     
     std::cout << "Running test: " << name << "\n";
     subscriber->OnTestStatusChange(name, TestStatus::RUNNING);
     std::this_thread::sleep_for(std::chrono::seconds(2));

     std::cout << "Test passed: " << name << "\n";
     subscriber->OnTestStatusChange(name, TestStatus::PASSED);
     std::this_thread::sleep_for(std::chrono::seconds(2));
  }

  subscriber->OnAllTestsFinish();
  std::cout << "Tests completed\n";
}

如果不需要测试线程持有Reactor的强引用,也可以使用weak_ptr:每次回调前先尝试将weak_ptr锁定为shared_ptr,锁定失败则说明Reactor已销毁,停止回调操作即可。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 03:24:51