注:该文用于个人学习记录和知识交流,如有不足,欢迎指点。

一、服务端异步API

服务端 gRPC 异步编程的核心模型:异步操作注册到 CQ → 底层完成操作后放入事件 → CQ 线程 Next () 取事件 → 调用 Proceed 处理

1. CompletionQueue

CompletionQueue(CQ)本质是 “异步操作完成事件的队列”

  • 服务端调用 RequestServerStreamWriteFinish 等异步方法时,并不会立即执行逻辑,而是把 “操作完成后需要触发的事件” 注册到 CQ 中;
  • gRPC 底层会在异步操作完成后(比如监听到客户端请求、数据写完、流结束),把对应的 “完成事件” 放入 CQ;
  • 你需要启动独立的 CQ 线程,循环调用 cq_->Next(&tag, &ok) 从队列中取出事件,然后触发对应的业务逻辑(Proceed)。

2. cq_->Next(&tag, &ok)

2.1 &tag:事件的 “唯一标识”(定位到具体的请求处理对象)

  • 类型void**(最终是 void* 类型);
  • 核心作用:用来标识 “这个完成事件属于哪个 CallData 实例”;
  • 你的代码场景:你调用 service_->RequestServerStream(&ctx_, &req_, &responder_, cq_, cq_, this) 时,最后一个参数 this(当前 ServerStreamCallData 对象的指针)就是这个 tag——gRPC 会把这个 this 和 “监听客户端服务端流请求” 的异步操作绑定;当 Next() 取出事件时,tag 就会被赋值为这个 this,你可以把它强转回 ServerStreamCallData*,从而找到对应的请求处理对象,调用它的 Proceed(ok)

2.2 &ok:异步操作的 “结果状态”(告诉业务逻辑操作是否成功)

  • 类型bool*
  • 核心作用:表示 “触发这个事件的异步操作是否成功完成”,也是你 Proceed(bool ok)ok 参数的来源;
  • 不同异步操作的 ok 含义(关键)
    异步操作类型 触发 Proceed 的时机 ok=true 含义 ok=false 含义
    1. RequestXXX 监听到客户端发起 RPC 请求时 成功监听到一个客户端的 RPC 请求 监听失败(如服务端关闭、CQ 关闭)
    2. Read() 尝试从客户端读取数据完成后 成功读到了客户端发送的一条数据 客户端已关闭写流(调用 WritesDone())/ 客户端调用 Finish() / 流异常中断
    3. Write() 尝试向客户端发送数据完成后 数据成功投递到 gRPC 发送缓冲区(不代表客户端已接收) 写数据失败(如流已关闭、网络异常、缓冲区满)
    4. Finish() 流关闭操作完成后 流终结操作成功投递到 gRPC 底层缓冲区(gRPC 会兜底发送完缓冲区中所有数据。注意此时客户端可能没接收到) 流关闭失败(极端底层错误,如服务端崩溃)

2.3 Next() 本身的行为:阻塞取事件

  • Next()阻塞调用:如果 CQ 中没有完成事件,调用 Next() 的线程会一直阻塞,直到有事件到来;
  • Next() 是 CQ 线程的核心:你必须在独立线程中循环调用 Next(),否则无法处理异步事件(比如:
    // CQ 事件循环线程函数
    void HandleCQ(ServerCompletionQueue* cq) {
      void* tag;
      bool ok;
      // 循环取事件,直到 CQ 被 Shutdown
      while (cq->Next(&tag, &ok)) {
        // 强转回 CallData,调用 Proceed 处理事件
        static_cast<CallData*>(tag)->Proceed(ok);
      }
    }
    

3. 结合 RequestServerStream 的完整流程

service_->RequestServerStream(&ctx_, &req_, &responder_, cq_, cq_, this) 是 “监听客户端服务端流请求” 的异步操作,结合 Next() 的完整执行链路:

  1. 初始化:服务端启动时,创建 ServerStreamCallData 实例,调用 RequestServerStream,把 “监听请求” 的异步操作注册到 cq_,绑定 tag=this
  2. 阻塞等待:CQ 线程调用 cq_->Next(&tag, &ok),阻塞等待事件;
  3. 客户端发起请求:客户端调用 stub->ServerStream(...) 发起服务端流请求;
  4. gRPC 触发事件:gRPC 底层完成 “监听请求” 操作,把事件(tag=this + ok=true)放入 CQ;
  5. 取出事件并处理:CQ 线程的 Next() 解除阻塞,tag 被赋值为 thisok 被赋值为 true
  6. 触发业务逻辑:把 tag 强转回 ServerStreamCallData*,调用 Proceed(true),进入状态机处理(比如创建新的 CallData 监听下一个请求,然后开始写数据)。

4. 关键补充

  1. CQ 关闭后 Next() 的行为:调用 cq->Shutdown() 后,Next() 会返回 false,事件循环线程可以退出(这也是上面 while (cq->Next(...)) 循环终止的条件);
  2. tag 必须有效tagthis(CallData 指针),因此必须保证 CallData 实例在 Next() 触发前不被释放,否则会导致野指针崩溃(这也是为什么 CallData 要在 FINISH 状态才 delete this);
  3. 两个 cq_ 参数的含义RequestServerStream 的第 4、5 个参数都是 cq_,分别表示 “监听请求的 CQ” 和 “后续操作(Write/Finish)的 CQ”,你这里用同一个 CQ,也可以用不同 CQ 做负载均衡(之前讲的多 CQ 场景)。

5. 总结

  • cq_->Next(&tag, &ok) 是 CQ 线程的核心:阻塞取出异步操作的完成事件,是异步逻辑的 “驱动引擎”;
  • tag 是 “事件归属标识”:通过 this 定位到具体的 CallData 实例,保证事件能精准分发;
  • ok 是 “操作结果”:告诉业务逻辑异步操作是否成功,是 Proceed 中判断逻辑分支的核心依据;
  • 结合 RequestServerStreamNext() 取出的第一个事件就是 “是否监听到客户端请求”,触发服务端流的首次 Proceed

二、客户端异步API

1. 先明确核心关系

这三行代码是客户端发起异步服务端流的 “三步曲”,缺一不可:

// 第一步:准备异步请求对象(仅创建对象,不发网络请求)
reader_ = stub_->PrepareAsyncServerStream(&ctx_, req_, cq_);
// 第二步:真正发起异步RPC请求(注册事件到CQ,不阻塞)
reader_->StartCall(this);
// 第三步:CQ线程阻塞等待“请求发起完成”的事件(驱动后续逻辑)
cq_.Next(&tag, &ok);

2. 逐行解析(结合参数 / 语义 / 和服务端对比)

2.1 reader_ = stub_->PrepareAsyncServerStream(&ctx_, req_, cq_);

  • 核心作用创建异步服务端流的 “请求对象”ClientAsyncReader<Response>),但不发起任何网络请求,仅完成 “请求上下文、请求数据、CQ” 的绑定。
  • 你可以把它理解为:客户端先 “打包好请求数据(req_)、RPC 上下文(ctx_)、事件队列(cq_)”,准备好 “发起请求的工具(reader_)”。
  • 参数解析:
    • &ctx_:RPC 上下文(可设置超时、元数据等);
    • req_:客户端发给服务端的唯一请求数据(服务端流是 “客户端单请求→服务端多响应”,所以只有这一次请求数据);
    • cq_:绑定的 CompletionQueue,后续所有异步事件(请求发起成功、读取服务端响应)都通过这个 CQ 触发;
  • 返回值:ClientAsyncReader<Response>*(即 reader_)—— 客户端用来读取服务端推送数据的核心对象(只有 Read() 方法,没有 Write(),因为服务端流是单向的)。

2.2 reader_->StartCall(this);

  • 核心作用真正发起异步 RPC 网络请求,并把 “请求发起完成事件” 注册到 cq_,绑定 tag=this(当前客户端 CallData 实例)。
  • 关键特性:
    • 异步非阻塞:调用后立即返回,不会等服务端响应,底层由 gRPC IO 线程处理网络请求;
    • this 的作用:和服务端 RequestServerStream(..., this) 中的 this 完全一样 —— 作为事件的 tag,标识 “这个请求发起事件属于当前 CallData 实例”;
    • 仅发起请求:不会触发数据读取,后续需要调用 reader_->Read(&res_, this) 才能读取服务端推送的响应。

2.3 cq_.Next(&tag, &ok);

  • 核心作用:和服务端 cq_->Next 逻辑完全一致,是客户端 CQ 线程的 “事件驱动引擎”——阻塞等待 cq_ 中的异步事件,取出后触发业务逻辑。
  • 针对 StartCall 对应的事件,参数含义:

    参数 具体含义(客户端服务端流场景) 和服务端的对比
    tag 取出的是 StartCall(this) 绑定的 this(当前 CallData 指针) 和服务端 RequestServerStreamtag 作用一致:定位事件归属的 CallData
    ok 表示 “发起 RPC 请求” 的操作是否成功:
    - ok=true:请求成功发到服务端,服务端已收到;
    - ok=false:请求发起失败(如服务端不可达、网络断连、CQ 已关闭)
    对应服务端 RequestServerStreamok(服务端是 “监听请求是否成功”,客户端是 “发起请求是否成功”)
  • 行为特性:
    • 阻塞:CQ 中无事件时,调用 Next 的线程会一直阻塞;
    • 循环调用:客户端必须在独立线程中循环调用 Next,否则无法处理后续的 Read 事件(服务端推送数据后触发的事件)。

2.4 ok的含义

异步操作类型 触发 Proceed 的时机 ok=true 含义 ok=false 含义
1. StartCall() 发起 RPC 请求完成后 成功向服务端发起 RPC 请求 请求发起失败(如服务端不可用、网络异常)
2. Write() 尝试向服务端发送数据完成后 数据成功投递到 gRPC 发送缓冲区(不代表服务端已接收) 写数据失败(如流已关闭、网络异常)
3. Read() 尝试从服务端读取数据完成后 成功读到了服务端发送的一条数据 服务端已关闭写流 / 服务端调用 Finish() / 流异常中断
4. WritesDone() 标记写流结束完成后 成功标记 “客户端不再发送数据” 标记写流结束失败(极端底层错误)
5. Finish() 收到服务端最终回执后 成功收到服务端的 Status(无论服务端传的是 OK 还是异常状态) 未收到服务端回执(如网络断连、服务端崩溃)

3. 客户端 vs 服务端

客户端(发起服务端流) 服务端(监听服务端流)
PrepareAsyncServerStream:创建请求对象 RequestServerStream:注册监听事件
StartCall(this):发起异步请求,绑定 tag 无需额外调用,RequestServerStream 直接注册监听事件
CQ.Next:等待 “请求发起完成” 事件 CQ.Next:等待 “监听到客户端请求” 事件
reader_->Read:读取服务端推送的数据 responder_.Write:向客户端推送数据

4. 总结

  • PrepareAsyncServerStream准备工具(创建 Reader 对象,不发请求);
  • StartCall(this)发起请求(异步发网络请求,注册事件到 CQ,绑定 tag=this);
  • cq_.Next(&tag, &ok)等待结果(阻塞取 “请求发起” 事件,tag 定位 CallData,ok 表示请求是否成功);
  • 客户端 CQ.Next 和服务端逻辑完全一致:tag 是事件归属标识,ok 是操作结果,阻塞取事件是异步逻辑的驱动核心。

三、grpc结束流程

1.  客户端

有writedown和finish这两个API

客户端在客户端流(ClientStream)/ 双向流(BidiStream) 场景下

WritesDone() = 「优雅告知对方:我主动结束发送,且已发完所有数据」;

直接 Finish() = 「强制终止整个 RPC,我的发送流被动关闭(对方会认为是异常)」。


1.1 客户端流(ClientStream)

客户端流的核心逻辑是「客户端多写 → 服务端单读 → 服务端单响应」,不调 WritesDone() 直接 Finish() 的后果最严重:

  • 服务端感知异常:服务端循环调用 reader->Read(&req) 时,会提前返回 false(而非读完所有数据后返回),且最终通过 Finish(&res) 拿到的 rpc_status 会是 CANCELLED/UNKNOWN(而非 OK),服务端会误以为「客户端断连 / 异常终止」,而非「正常发完数据」。
  • 数据丢失:如果客户端有未发送的缓存数据(比如 Write() 了但还没刷到网络),直接 Finish() 会导致这些数据被 gRPC 底层丢弃,服务端完全收不到。
  • 业务逻辑错误:服务端本应根据客户端发送的所有数据做计算 / 处理,现在因提前中断,会得到不完整的数据,返回错误的响应。

示例(客户端流错误写法)

// 客户端流 - 错误:直接Finish,不调WritesDone
class ClientStreamCallData {
  void Proceed(bool ok) override {
    if (status_ == WRITING) {
      if (send_idx_ < send_msgs_.size()) {
        // 发送部分数据
        req_.set_data(send_msgs_[send_idx_++]);
        writer_->Write(req_, this);
      } else {
        // ❌ 直接Finish,不调WritesDone
        status_ = FINISH;
        writer_->Finish(res_, this); // 强制关闭
      }
    }
  }
};

// 服务端感知异常
Status ExampleService::ClientStream(ServerContext* ctx,
                                   ServerReader<Request>* reader,
                                   Response* res) {
  Request req;
  while (reader->Read(&req)) { // 提前返回false,只收到部分数据
    std::cout << "收到数据:" << req.data() << std::endl;
  }
  // 服务端拿到的状态是异常(而非OK)
  std::cout << "客户端流结束,状态:" << ctx->status().error_message() << std::endl;
  return Status::OK;
}

1.2 双向流(BidiStream):状态感知错位

双向流客户端有「发送流 + 接收流」,不调 WritesDone() 直接 Finish() 的后果:

  • 服务端误判发送流状态:服务端循环 stream->Read(&req) 会提前返回 false,认为「客户端发送流异常关闭」(比如断连),而非「客户端主动结束发送」;
    但服务端仍能正常接收自己的 Write() 响应,只是日志会记录 “客户端发送流异常”。
  • 客户端未发送数据丢失:客户端已 Write() 但未刷到网络的缓存数据,会被丢弃,服务端收不到。
  • 上下游状态不一致:客户端 Finish() 后拿到的 rpc_status 可能是 OK(如果服务端没报错),但服务端感知到的是 “客户端发送流异常”,导致日志 / 监控中出现 “客户端认为成功,服务端认为异常” 的矛盾。

1.3「优雅关闭」vs「强制关闭」对比表

操作方式 发送流结束语义 未发送缓存数据 服务端 Read() 行为 服务端感知状态
WritesDone() + Finish() 主动、正常结束 全部发送 读完所有数据后返回 false 正常(OK
直接 Finish() 被动、异常终止 直接丢弃 提前返回 false 异常(CANCELLED/UNKNOWN

1.4 例外场景:直接 Finish() 是合理的

只有一种情况可以直接 Finish()客户端主动想取消 RPC(比如用户中断操作、超时),此时需要传异常状态,明确告知服务端 “我主动取消,而非正常结束”:

// 合理场景:客户端主动取消RPC
status_ = FINISH;
writer_->Finish(res_, Status(grpc::StatusCode::CANCELLED, "用户主动中断"), this);

1.5 总结

  • 语法上:不调 WritesDone() 直接 Finish() 完全可行,不会编译 / 运行报错;
  • 逻辑上:会导致「数据丢失 + 服务端误判异常」,生产环境(尤其是需要数据完整收发的场景,如批量传输、聊天)严禁这么做
  • 正确做法

    正常结束发送:先 WritesDone()(告知对方 “发完了”),再 Finish()(等待最终响应);

    主动取消 RPC:直接 Finish(),但要传 CANCELLED 等异常状态,让服务端感知真实原因。

这也是 gRPC 设计 WritesDone() 的核心目的 —— 区分「正常结束发送」和「异常终止 RPC」,避免上下游状态感知错位。

2. 服务端

服务端只有finish接口,没有writedown接口

3. 两者finish的差别

客户端调用 stream_->Finish(&rpc_status_, this) 后,rpc_status_ 的最终值完全由服务端调用 stream_.Finish(xxx, this) 时传入的 Status 参数决定—— 这是 gRPC 双向流中 “服务端→客户端” 传递最终调用状态的核心机制。

维度 客户端 Finish() 服务端 Finish()
核心语义 等待服务端的最终回执 发送最终回执给客户端,终结 RPC
触发 Proceed 时机 必须等服务端 Finish() 并传回 Status 调用后立即触发(本地操作,无需等待)
ok 参数核心含义 是否收到服务端的 Status(网络 / 服务端状态) 本地 Finish 操作是否成功(几乎永远 true)
角色 被动等待方 主动终结方

注意:服务端 Finish() 触发 Proceed 的唯一条件是 “本地调用 Finish() 成功”,和数据是否发送完毕、客户端是否收到数据完全无关—— 这是 gRPC 异步编程的核心特性,服务端只需关注 “业务逻辑是否完成”,数据发送的细节全由底层兜底。

场景 服务端 Proceed 触发时机 数据发送状态
调用 Finish() 立即触发(本地操作) 数据可能还在缓冲区 / 网络中
Proceed 进入 FINISH 无需等数据发送完毕 gRPC 底层兜底发送剩余数据
数据最终是否到客户端 不影响服务端 Proceed 执行 失败则客户端收到异常 Status

4. 想要请求rpc结束

RPC 类型 客户端要求 服务端要求
服务端流 无需调用WritesDone(),仅通过Read()接收直到返回false 调用Finish()
(注:同步 return Status::OK 默认发送结束符)
客户端流 发送完数据后调用WritesDone(),最后Finish() 接收Read()返回false后,调用Finish()响应
双向流 结束发送时调用WritesDone(),最终Finish() 结束发送调用Finish()

5. 双向流完整流程(Finish 触发链路)

5.1 客户端正常关闭(正常逻辑)

  • 客户端发送所有数据 → 调用 WritesDone()(告知服务端 “我发完了”),发送成功后调用Finish();
  • 服务端 Read() 返回 ok=false(感知客户端发完)→ 调用 Finish(Status::OK, this)
  • 服务端 Finish() 立即触发自身 Proceed 进入 FINISH(释放资源),同时把 Status::OK 发给客户端;
  • 客户端收到服务端的 Status::OK → 触发自身 Finish() 对应的 Proceed 进入 FINISH(释放资源);
  • 整个 RPC 终结。

5.2 客户端被动关闭

  • 服务端先 Finish()
  • 客户端收到 ok=false
  • 客户端调用 Finish(&rpc_status_,this) 。这里仍能rpc_status_.ok()判断是否为正常还是异常

四、补充:Finsih的细节

gRPC 中 Finish 不是 “一个方法”,而是 “一类终结 RPC 的行为”—— 核心分为两大阵营:

  • 服务端 Finish:「主动决策者」→ 宣告 RPC 终结,发送最终状态给客户端,是 RPC 终结的 “总开关”;
  • 客户端 Finish:「被动等待者」→ 等待服务端的终结指令,拿到最终状态后释放自身资源,是 RPC 收尾的 “确认操作”。

同步 / 异步场景的差异仅在于「是否显式调用方法」(同步是框架隐式处理,异步是手动显式调用),但核心语义不变。

1. 服务端的 Finish(主动终结 RPC)

服务端的 Finish 是「终结本次 RPC 的唯一操作」,无论一元 / 流、同步 / 异步,核心都是:✅ 发最终状态(Status)给客户端 + 兜底发完缓冲区未发送的数据(流 RPC 有效) + 释放 RPC 资源。

1.1 同步服务端:隐式 Finish(return Status 就是 Finish)

同步场景下没有 Finish() 方法,框架通过函数返回值隐式完成 Finish 的所有逻辑 —— 你可以理解为 “框架替你调用了 Finish”。

特性 一元 RPC 具体说明 流 RPC(服务端流 / 客户端流 / 双向流)具体说明
触发方式 函数执行到 return Status 时自动触发 函数执行到 return Status 时自动触发
阻塞性 return 前仅处理单次请求 / 响应,无阻塞写 / 读(一元 RPC 无流) return 前的 Write()/Read() 都是阻塞操作(直到数据进入 gRPC 缓冲区)
数据兜底 无缓冲区数据(一元 RPC 仅单次响应) return 前所有已写未发的缓存数据,gRPC 会兜底发给客户端
资源释放 框架自动释放:RPC 上下文(ServerContext)、请求 / 响应对象等 框架自动释放:RPC 上下文(ServerContext)、流对象(ServerWriter/ServerReaderWriter)等
状态传递 return 的 Status 直接传给客户端,客户端通过 stub 调用的返回值获取 return 的 Status 直接传给客户端,客户端通过 stub 调用 /ctx 获取

示例 1:同步一元 RPC(服务端)

// 同步一元 RPC:return Status::OK 就是“隐式 Finish”
Status ExampleService::UnaryRPC(ServerContext* ctx, 
                                const Request* req, 
                                Response* res) {
  // 处理单次请求,返回单次响应
  res->set_data("Sync Unary Echo: " + req->data());
  
  // ✅ 隐式 Finish:发送 Status::OK + 释放资源(一元 RPC 无缓冲区数据)
  return Status::OK; 
  // 若返回异常:return Status(grpc::StatusCode::INVALID_ARGUMENT, "数据非法");
}

示例 2:同步服务端流(服务端)

// 同步服务端流:return Status::OK 就是“隐式 Finish”
Status ExampleService::ServerStream(ServerContext* ctx, 
                                    const Request* req, 
                                    ServerWriter<Response>* writer) {
  // 循环推送数据(阻塞写,确保数据进缓冲区)
  for (int i = 0; i < 3; ++i) {
    Response res;
    res.set_data("Sync Echo: " + std::to_string(i));
    writer->Write(res); 
  }
  
  // ✅ 隐式 Finish:发送 Status::OK + 兜底发数据 + 释放资源
  return Status::OK; 
}

1.2 异步服务端:显式 Finish(手动调用 Finish () 方法)

异步场景下必须手动调用 Finish(),主动触发 RPC 终结,需配合状态机使用(一元 RPC 和流 RPC 的调用形式略有差异)。

特性 一元 RPC 具体说明 流 RPC(服务端流 / 客户端流 / 双向流)具体说明
触发方式 显式调用 responder.Finish(响应, Status, this),绑定当前 CallData 作为事件 tag 显式调用 responder_.Finish(Status, this)/stream_.Finish(Status, this),绑定当前 CallData 作为事件 tag
阻塞性 非阻塞:调用后立即返回,触发 Proceed 进入 FINISH 分支 非阻塞:调用后立即返回,触发 Proceed 进入 FINISH 分支
数据兜底 无缓冲区数据(一元 RPC 仅单次响应) Finish 前的缓存数据会被兜底发送
资源释放 手动释放:需在 FINISH 分支 delete this(CallData 实例) 手动释放:需在 FINISH 分支 delete this(CallData 实例)
状态传递 Finish 传入的 Status 发给客户端,客户端通过自身 Finish 拿到 Finish 传入的 Status 发给客户端,客户端通过自身 Finish 拿到

示例 1:异步一元 RPC(服务端)

void Proceed(bool ok) override {
  if (status_ == PROCESSING) {
    // 处理单次请求,构造响应
    Response res;
    res.set_data("Async Unary Echo: " + req_.data());
    
    // ✅ 显式 Finish:返回响应 + 发送 Status::OK + 触发 Proceed 进入 FINISH
    status_ = FINISH;
    responder_.Finish(res, Status::OK, this); 
    // 若传异常:responder_.Finish(res, Status(grpc::StatusCode::INVALID_ARGUMENT, "数据非法"), this);
  } else if (status_ == FINISH) {
    delete this; // 手动释放资源
  }
}

示例 2:异步服务端流(服务端)

void Proceed(bool ok) override {
  if (status_ == WRITING) {
    if (send_count_ < 3) {
      // 异步写(非阻塞,数据进缓冲区即返回)
      res_.set_data("Async Echo: " + std::to_string(send_count_++));
      responder_.Write(res_, this); 
    } else {
      // ✅ 显式 Finish:等价于同步的 return Status::OK
      status_ = FINISH;
      responder_.Finish(Status::OK, this); 
    }
  } else if (status_ == FINISH) {
    delete this; // 手动释放资源(同步由框架处理)
  }
}

2.客户端的 Finish(被动等待回执)

客户端的 Finish 是「等待服务端终结指令的确认操作」,核心是:✅ 等服务端的最终 Status + 释放客户端侧的 RPC 资源,永远无法主动终结 RPC(只有服务端能)。

2.1 同步客户端:含隐式 Finish

特性 一元 RPC 具体说明 服务端流具体说明 客户端流 / 双向流具体说明
触发方式 调用 stub->UnaryRPC(&ctx, req, &res) 后,直接阻塞等待服务端 Finish,函数返回 Status

-手动reader->Finish

-若未手动调用finish,reader 析构时也会触发 Finish,阻塞拿 Status 并填充到 ctx

调用 writer->Finish() / stream->Finish() 时触发 Finish,阻塞等待服务端响应并返回 Status
阻塞性 全程阻塞:直到服务端 Finish 并传回 Status + 响应,函数才返回 读数据时不阻塞 Status 获取,仅调用reader->Finish时会阻塞 调用 Finish() 时阻塞:直到服务端 Finish 并传回 Status + 响应,函数才返回
状态获取 直接通过 stub 调用返回值拿到服务端的最终 Status 通过reader->Finish() 直接通过 writer->Finish() / stream->Finish() 的返回值拿到服务端的最终 Status
资源释放 框架自动释放:stub、上下文、请求 / 响应对象等 框架自动释放:reader 析构时释放流资源,上下文自动释放 框架自动释放:writer / stream 析构时释放流资源,上下文自动释放

示例 1:同步一元 RPC(客户端)

// 同步一元 RPC:stub 调用直接返回 Status(阻塞等服务端 Finish)
std::unique_ptr<ExampleService::Stub> stub = ExampleService::NewStub(channel);
ClientContext ctx;
Request req;
Response res;
req.set_data("Hello");

// ✅ 隐式 Finish:阻塞等服务端 Finish,直接拿到最终 Status
Status status = stub->UnaryRPC(&ctx, req, &res); 
if (status.ok()) {
  std::cout << "RPC正常终结,响应:" << res.data() << std::endl;
} else {
  std::cerr << "RPC异常:" << status.error_message() << std::endl;
}

示例 2:同步服务端流

Request req;
req.set_data(msg);

ClientContext ctx;

std::unique_ptr<ClientReader<Response>> reader = stub_->ServerStream(&ctx, req);
Response res;
while (reader->Read(&res))
{ // Read()阻塞等待响应,收到数据时返回 true,收到结束标记时返回 false,
   std::cout << "Server Stream: " << res.data() << std::endl;
}
Status status = reader->Finish(); // 如果不手动调用,reader析构时也会自动调用Finish获取状态
    

2.2 异步客户端:显式 Finish(手动调用 Finish () 方法)

异步客户端必须手动调用 Finish(),等待服务端的 Finish 回执,是异步状态机的最后一步(一元 RPC 和流 RPC 逻辑一致)。

特性 一元 RPC 具体说明 流 RPC(服务端流 / 客户端流 / 双向流) 具体说明
触发方式 显式调用 Finish(&rpc_status_, this),绑定当前 CallData 作为事件 tag 显式调用 Finish(&rpc_status_, this),绑定当前 CallData 作为事件 tag
阻塞性 异步阻塞:调用后 CQ 线程阻塞,直到服务端 Finish 并传回 Status 异步阻塞:调用后 CQ 线程阻塞,直到服务端 Finish 并传回 Status
状态获取 通过 &rpc_status_ 拿到服务端的最终 Status 通过 &rpc_status_ 拿到服务端的最终 Status
资源释放 手动释放:需在 FINISH 分支 delete this(CallData 实例) 手动释放:需在 FINISH 分支 delete this(CallData 实例)

示例 1:异步一元 RPC(客户端)

void Proceed(bool ok) override {
  if (status_ == SENDING) {
    // 发送一元请求
    stub_->async()->UnaryRPC(&ctx_, &req_, &res_, this);
    status_ = WAITING;
    StartCall();
  } else if (status_ == WAITING) {
    // ✅ 服务端已 Finish,调用 Finish 拿最终状态
    status_ = FINISH;
    stub_->async()->Finish(&rpc_status_, this); 
  } else if (status_ == FINISH) {
    // 拿到服务端的 Finish 状态
    if (rpc_status_.ok()) {
      std::cout << "一元 RPC 正常终结,响应:" << res_.data() << std::endl;
    } else {
      std::cerr << "一元 RPC 异常:" << rpc_status_.error_message() << std::endl;
    }
    delete this; // 手动释放资源
  }
}

示例 2:异步服务端流(客户端)

void Proceed(bool ok) override {
  if (status_ == READING) {
    if (ok) {
      // 成功读到服务端数据,继续读
      std::cout << "收到数据:" << res_.data() << std::endl;
      reader_->Read(&res_, this);
    } else {
      // ✅ 服务端已 Finish,客户端调用 Finish 拿最终状态
      status_ = FINISH;
      reader_->Finish(&rpc_status_, this); 
    }
  } else if (status_ == FINISH) {
    // 拿到服务端的 Finish 状态
    if (rpc_status_.ok()) {
      std::cout << "RPC正常终结" << std::endl;
    } else {
      std::cerr << "RPC异常:" << rpc_status_.error_message() << std::endl;
    }
    delete this; // 手动释放资源
  }
}

3. 核心对比表

3.1 服务端 Finish 操作表

RPC 类型 场景 Finish 形式 核心语义 资源释放 阻塞性
一元 同步 / 异步 同步:return Status异步:responder.Finish(响应, Status, this) 主动终结 RPC,发最终状态 同步:框架自动异步:手动 delete 同步:无异步:非阻塞
服务端流 / 客户端流 / 双向流 同步 / 异步 同步:return Status异步:调用 Finish() 主动终结 RPC,发最终状态 同步:框架自动异步:手动 delete 同步:阻塞(写数据时)异步:非阻塞

3.2 客户端 Finish 操作表(终极压缩版)

RPC 类型 场景 Finish 形式 核心语义 资源释放 阻塞性
一元 同步 / 异步 同步:stub 调用返回 Status异步:调用 Finish() 被动等服务端终结,拿最终状态 同步:框架自动异步:手动 delete 同步:全程阻塞异步:异步阻塞
服务端流 同步 / 异步 同步:输出参数重载返回 Status / unique_ptr 析构触发异步:调用 Finish() 被动等服务端终结,拿最终状态 同步:框架自动异步:手动 delete 同步:析构 / 返回时阻塞异步:异步阻塞
客户端流 / 双向流 同步 / 异步

同步:writer->Finish() 返回 Status

异步:调用 Finish()

被动等服务端终结,拿最终状态 同步:框架自动异步:手动 delete 同步:调用 Finish 时阻塞异步:异步阻塞

4. 总结(关键记忆点)

  • Finish 并非仅流 RPC 特有:一元 RPC 同样有 Finish 逻辑,服务端通过 return Status(同步)/responder.Finish(异步)终结 RPC,客户端通过 stub 返回值(同步)/Finish()(异步)拿状态;
  • 服务端 Finish:RPC 终结的 “总开关”,同步隐式(return)、异步显式(调用方法),核心是 “主动终结 + 发状态”;
  • 客户端 Finish:RPC 收尾的 “确认键”,同步隐式(一元:stub 返回值;流:reader 析构 /stub 返回值)、异步显式(调用方法),核心是 “被动等状态 + 释放资源”;
  • 核心规则:永远是「服务端先 Finish,客户端才能 Finish」—— 客户端 Finish 只是 “等服务端指令”,无法主动终结 RPC;。

五、异步编写需遵循的原则

  1. 状态机是骨架:状态与操作一一对应,流转靠ok回调;

    即,执行了read,就令status = Reading。
  2. 单 Pending 准则:

    同类型
    的异步操作,同一时间只能有 1 个 处于 “待处理(pending)” 状态;不同类型的操作(读 vs 写)是完全独立的,可以同时各有 1 个 pending 操作。

    简单将就是:可以同时发起单个Read和单个Write,但是不能同时发起两个Read(或两个Write)
  3. CQ 管理要规范:多 CQ 原子分配,生命周期先 Shutdown 后 Join;
  4. 资源自销毁:CallData 仅在FINISH状态delete
  5. 异常必处理ok=false时主动Finish(),不遗漏错误;
  6. 流结束要规范WritesDone()/Finish()按类型调用;
  7. 线程安全不忽视:共享变量用原子 / 加锁。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐