Linux C/C++ 学习日记(71):grpc(四):基于grpc编写异步server和client(1):常用api前置知识
注:该文用于个人学习记录和知识交流,如有不足,欢迎指点。
一、服务端异步API
服务端 gRPC 异步编程的核心模型:异步操作注册到 CQ → 底层完成操作后放入事件 → CQ 线程 Next () 取事件 → 调用 Proceed 处理。
1. CompletionQueue
CompletionQueue(CQ)本质是 “异步操作完成事件的队列”:
- 服务端调用
RequestServerStream、Write、Finish等异步方法时,并不会立即执行逻辑,而是把 “操作完成后需要触发的事件” 注册到 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() 的完整执行链路:
- 初始化:服务端启动时,创建
ServerStreamCallData实例,调用RequestServerStream,把 “监听请求” 的异步操作注册到cq_,绑定tag=this; - 阻塞等待:CQ 线程调用
cq_->Next(&tag, &ok),阻塞等待事件; - 客户端发起请求:客户端调用
stub->ServerStream(...)发起服务端流请求; - gRPC 触发事件:gRPC 底层完成 “监听请求” 操作,把事件(
tag=this+ok=true)放入 CQ; - 取出事件并处理:CQ 线程的
Next()解除阻塞,tag被赋值为this,ok被赋值为true; - 触发业务逻辑:把
tag强转回ServerStreamCallData*,调用Proceed(true),进入状态机处理(比如创建新的 CallData 监听下一个请求,然后开始写数据)。
4. 关键补充
- CQ 关闭后
Next()的行为:调用cq->Shutdown()后,Next()会返回false,事件循环线程可以退出(这也是上面while (cq->Next(...))循环终止的条件); tag必须有效:tag是this(CallData 指针),因此必须保证 CallData 实例在Next()触发前不被释放,否则会导致野指针崩溃(这也是为什么 CallData 要在FINISH状态才delete this);- 两个
cq_参数的含义:RequestServerStream的第 4、5 个参数都是cq_,分别表示 “监听请求的 CQ” 和 “后续操作(Write/Finish)的 CQ”,你这里用同一个 CQ,也可以用不同 CQ 做负载均衡(之前讲的多 CQ 场景)。
5. 总结
cq_->Next(&tag, &ok)是 CQ 线程的核心:阻塞取出异步操作的完成事件,是异步逻辑的 “驱动引擎”;tag是 “事件归属标识”:通过this定位到具体的 CallData 实例,保证事件能精准分发;ok是 “操作结果”:告诉业务逻辑异步操作是否成功,是Proceed中判断逻辑分支的核心依据;- 结合
RequestServerStream:Next()取出的第一个事件就是 “是否监听到客户端请求”,触发服务端流的首次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 指针)和服务端 RequestServerStream的tag作用一致:定位事件归属的 CallDataok表示 “发起 RPC 请求” 的操作是否成功:
-ok=true:请求成功发到服务端,服务端已收到;
-ok=false:请求发起失败(如服务端不可达、网络断连、CQ 已关闭)对应服务端 RequestServerStream的ok(服务端是 “监听请求是否成功”,客户端是 “发起请求是否成功”) - 行为特性:
- 阻塞:CQ 中无事件时,调用
Next的线程会一直阻塞; - 循环调用:客户端必须在独立线程中循环调用
Next,否则无法处理后续的Read事件(服务端推送数据后触发的事件)。
- 阻塞:CQ 中无事件时,调用
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 |
同步:析构 / 返回时阻塞异步:异步阻塞 |
| 客户端流 / 双向流 | 同步 / 异步 |
同步: 异步:调用 |
被动等服务端终结,拿最终状态 | 同步:框架自动异步:手动 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;。
五、异步编写需遵循的原则
- 状态机是骨架:状态与操作一一对应,流转靠
ok回调;
即,执行了read,就令status = Reading。 - 单 Pending 准则:
同类型的异步操作,同一时间只能有 1 个 处于 “待处理(pending)” 状态;不同类型的操作(读 vs 写)是完全独立的,可以同时各有 1 个 pending 操作。
简单将就是:可以同时发起单个Read和单个Write,但是不能同时发起两个Read(或两个Write) - CQ 管理要规范:多 CQ 原子分配,生命周期先 Shutdown 后 Join;
- 资源自销毁:CallData 仅在
FINISH状态delete; - 异常必处理:
ok=false时主动Finish(),不遗漏错误; - 流结束要规范:
WritesDone()/Finish()按类型调用; - 线程安全不忽视:共享变量用原子 / 加锁。
更多推荐




所有评论(0)