如何在gRPC双向流式RPC中取消读取而不终止整个流?
在gRPC双向流式RPC中取消单次Read调用而不终止整个流
在gRPC服务端的双向流式RPC调用中,需要实现取消当前stream->Read(..)的阻塞等待,让调用返回但不终止整个RPC流的逻辑。取消时机由外部事件(如用户操作)触发,无法提前预知,该需求同样适用于客户端的读取场景。
原代码示例(存在阻塞问题)
#include <chrono> #include <thread> #include <atomic> #include "my_service.h" class MyServiceImpl : public MyService::Service { public: Status my_rpc_call(ServerContext* context, ServerReaderWriter<RPCResponse, RPCRequest>* stream) override { RPCRequest request; while (stream->Read(&request)) { // 问题:如果客户端不发送请求,会永久阻塞 // 处理请求逻辑 } /* 期望实现: // 监听用户输入的线程 std::atomic_bool cancel_read; std::jthread cancel_periodically_thread([&cancel_read] { for (int i(0); true; ++i) { std::this_thread::sleep_for(std::chrono::seconds(i)); // 模拟用户触发取消的时机 cancel_read = true; } }); while (true) { cancel_read = false; while (stream->Read(&request, cancel_read)) { // 问题:gRPC无此重载,且原生Read会永久阻塞 // 处理请求逻辑 } stream->Write("please input something"); // 仅示意语法,实际需构造RPCResponse } */ } };
解决方案
gRPC的同步Read调用本身不支持直接中断,因此需要通过异步读取+超时检查的方式实现可中断的读等待,既可以响应外部取消信号,又不会终止整个RPC流。
服务端实现示例
#include <chrono> #include <thread> #include <atomic> #include "my_service.h" class MyServiceImpl : public MyService::Service { public: Status my_rpc_call(ServerContext* context, ServerReaderWriter<RPCResponse, RPCRequest>* stream) override { std::atomic_bool cancel_read{false}; // 模拟外部触发取消的线程(实际可替换为用户操作监听逻辑) std::jthread cancel_trigger_thread([&cancel_read, &context] { while (!context->IsCancelled()) { // 模拟5秒后触发取消 std::this_thread::sleep_for(std::chrono::seconds(5)); cancel_read = true; } }); RPCRequest request; RPCResponse response; while (!context->IsCancelled()) { cancel_read = false; bool read_success = false; // 异步读取+超时检查,实现可中断的等待 grpc::CompletionQueue cq; stream->AsyncRead(&request, &cq, nullptr); while (!cancel_read && !context->IsCancelled()) { // 每次等待1秒,超时后检查取消标志 auto cq_status = cq.AsyncNext(nullptr, nullptr, std::chrono::seconds(1)); if (cq_status == grpc::CompletionQueue::GOT_EVENT) { // 成功读取到请求,处理逻辑 response.set_message("已收到请求"); stream->Write(response); read_success = true; break; } else if (cq_status == grpc::CompletionQueue::SHUTDOWN) { // 队列关闭,退出循环 break; } // 超时则继续循环,检查取消标志 } if (cancel_read && !read_success) { // 触发了取消,发送提示信息 response.set_message("请输入内容"); stream->Write(response); } } return grpc::Status::OK; } };
核心逻辑说明
- 异步读取替代同步阻塞:使用
AsyncRead将读取操作放入CompletionQueue,避免同步Read的永久阻塞。 - 定期检查取消标志:通过
CompletionQueue::AsyncNext的超时参数(这里设为1秒),定期跳出等待并检查外部取消信号cancel_read。 - 不终止RPC流:仅中断当前的读等待,RPC上下文
ServerContext并未被取消,因此整个双向流保持活跃,可继续后续的读写操作。
客户端场景适配
客户端的实现逻辑完全类似,将ServerReaderWriter替换为ClientReaderWriter,同样通过AsyncRead+CompletionQueue超时检查的方式,实现可中断的读等待,同时不终止客户端的流式连接。
内容的提问来源于stack exchange,提问作者Maximilian Mordig
相关产品推荐
相关产品推荐

