深入解析 gRPC 客户端与服务端 Polling Engine 的使用机制
【免费下载链接】grpcC++ based gRPC (C++, Python, Ruby, Objective-C, PHP, C#)项目地址: https://gitcode.com/GitHub_Trending/gr/grpc
导读
本文以 gRPC 核心文档 grpc-client-server-polling-engine-usage.md 为骨架,系统性讲解 gRPC 底层 I/O 多路复用(Polling Engine)在客户端与服务端两条代码路径上分别如何被使用:客户端如何把 Call 与 Channel/Completion Queue 关联、并通过grpc_pollset_set推进子通道(sub-channel)上的异步connect();服务端又如何把监听 fd 挂到每个 Completion Queue 的 pollset 上,借助SO_REUSEPORT与轮询(round-robin)分配新连接。读完本文,你将掌握 gRPC core 中grpc_fd、grpc_pollset、grpc_pollset_set、grpc_pollset_worker四类核心对象的关系,并能结合 polling engine 文档、completion queue 文档 与src/core/lib/iomgr下的真实实现,独立看懂任意一条 gRPC 连接在读写事件层面的完整流转链路。
适用说明:本文描述的机制位于 gRPC core 的 iomgr 层,源码路径以当前仓库(
src/core/lib/iomgr/)为准。文中出现的具体行号为当前仓库实现中的近似位置,不同版本(如文档撰写时期的 v1.15.1)行号与实现细节可能略有出入,但总体架构保持一致。
一、背景:gRPC 为什么需要 Polling Engine
在深入客户端与服务端的使用细节之前,先明确 Polling Engine 的存在意义(详见 grpc-polling-engines.md):
- gRPC core 需要同时监控大量文件描述符(fd)上的可读 / 可写 / 出错三类事件;
- 当事件发生时,gRPC 知道要执行的具体动作——例如
grpc_endpoint在 fd 可读时调用recvmsg、可写时调用sendmsg;tcp_client的 connect 代码在 fd 变为可写(即connect()真正完成)时结束客户端创建流程; - gRPC 需要一个组件"高效地"完成上述监控,并且复用应用自己提供的线程(绝不新建线程),同时针对时延与吞吐进行优化。
Polling Engine 正是承担这一职责的抽象层。它按操作系统提供多种实现(Linux 上的epoll1、无 epoll 时的poll、macOS 上的poll、Windows 上的 I/O completion port 实现),但它们全部暴露同一套接口,上层代码无需关心具体平台。
这套接口抽象出四个关键结构:
| 结构 | 含义 |
|---|---|
grpc_fd | 对一个文件描述符的封装(注意各实现内部定义不同,对外不透明) |
grpc_pollset | 一组被轮询可读/可写/出错事件的grpc_fd的集合;同一个 fd 可同时存在于多个 pollset |
grpc_pollset_worker | 一个"轮询线程"的抽象,特指调用grpc_pollset_work()的线程 |
grpc_pollset_set | 一个容器,可包含grpc_fd、grpc_pollset,甚至嵌套其他的grpc_pollset_set |
对象关系可参考文档 grpc-ps-pss-fd.png(worker 持有/驱动 pollset,pollset 聚合 fd)与 grpc-pss.png(pollset_set 的组合语义)。
二、gRPC 客户端:Call、Channel、Completion Queue 与 pollset 的绑定关系
2.1 生命周期内的三角绑定
在 gRPC 客户端,一次 RPC(Call)的完整生命周期内同时与三类实体绑定:
- Channel(具体来说是 sub-channel):Call 必须挂在某个 channel 上发起;
- Completion Queue:用户通过
grpc_completion_queue_next()/grpc_completion_queue_pluck()领取完成事件; grpc_pollset:依据 grpc-cq.md 的描述,一个 Completion Queue 默认自带一个 pollset。
关键规则如下:
- 一个 gRPC Call 在其整个生命周期内,与某个 channel(更确切地说是其中的 sub-channel)以及某个 Completion Queue 绑定;
- 一旦为 Call 选定 sub-channel,该连接的 fd(TCP channel 场景下即 socket fd)就会被加入该 Call 对应 Completion Queue 的 pollset。
这两条规则构成了客户端事件分发的基础:Call 所在 CQ 的 pollset 一旦被某个线程通过grpc_pollset_work()驱动起来,就能及时发现该 Call 的 socket 上发生的读写事件。下图直观展示了 Call、Channel(sub-channel)、Completion Queue 之间的关系:
从当前仓库源码可以印证"CQ 自带 pollset"这一设计:在 src/core/lib/surface/completion_queue_factory.cc 中,grpc_completion_queue_create_for_next()、grpc_completion_queue_create_for_pluck()、grpc_completion_queue_create_for_callback()创建出的 CQ 均携带GRPC_CQ_DEFAULT_POLLING属性——即调用next/pluck的线程会执行真实的事件轮询,而轮询的对象正是 CQ 内部关联的那个 pollset。
2.2 Completion Queue 如何驱动 pollset 工作
理解了绑定关系后,还需要回答:谁、在何时真正去 poll 这些 fd?答案藏在grpc_completion_queue_next()/grpc_completion_queue_pluck()的循环实现里(伪代码见 grpc-cq.md):
grpc_completion_queue_next(cq, deadline) / pluck(cq, deadline, tag) { while (true) { // 1. 若 CQ 中已有就绪事件,出队并返回 // (pluck 模式下仅返回与目标 tag 匹配的事件) // 2. 若 CQ 已 shutdown,直接返回 // 3. pluck 模式下,将 (tag, worker) 对登记到 CQ 的 tag<->worker 映射中 // 4. 调用 grpc_pollset_work(cq 的 pollset, deadline) 执行轮询 // ——若发现某些 fd 可读/可写/出错, // 会调度对应的 closure(这些 closure 可能向"某个"CQ 投递完成事件, // 注意不一定是当前这个 CQ) } }反过来,当一个 I/O 事件完成、需要向 CQ 投递完成事件(即把 tag 入队)时,走的是grpc_cq_end_op:
- 先把 tag 放入事件队列;
- 找到该 CQ 对应的 pollset 并"踢醒"等待线程:
- 若 CQ 类型为
GRPC_CQ_NEXT:调用grpc_pollset_kick(pollset, nullptr),踢醒任意一个worker; - 若 CQ 类型为
GRPC_CQ_PLUCK:先查 CQ 上的 tag↔worker 映射,找到等待该 tag 的那个特定 worker,再调用grpc_pollset_kick(pollset, worker)定向踢醒。
- 若 CQ 类型为
由此形成闭环:用户线程阻塞在pollset_work上等待 fd 事件 → 事件触发后 closure 入队完成事件并 kick worker →pollset_work返回 → 循环重新出队并返回事件给用户。这正是"用应用线程做轮询、不额外创建线程"的落地方式。
三、客户端推进子通道异步 connect:grpc_pollset_set的经典用例
3.1 问题:大量在途 connect 需要被"额外"监控
理解完 Call 与 CQ 的绑定后,还有一个独立于 Call 生命周期的问题需要处理:
- 一个 gRPC Channel 建立在客户端与某个"target"之间,而该 target 经域名解析后可能对应一个或多个后端服务器;
- sub-channel 正是客户端到某个后端服务器之间的那条"连接";
- 在向各后端建立 sub-channel(即发起连接)的过程中,gRPC 会发出异步
connect()(当前实现见 src/core/lib/iomgr/tcp_client_posix.cc 的grpc_tcp_client_create_from_prepared_fd),这类调用通常不会立刻完成; - 当
connect()最终成功时,对应的 socket fd 会变为"可写(writable)"。
于是出现一个难点:此刻"连接尚未建立",sub-channel 尚未与任何 Call/CQ 绑定,这些进行中的 connect fd 不属于任何 CQ 的 pollset。然而 Polling Engine 必须持续监控所有这些 sub-channel 的 fd 上的可写事件,并保证确实存在轮询线程在监控它们——否则连接永远无法建立,后续所有 RPC 都会卡死。
3.2 解决:把"感兴趣的各方"装进grpc_pollset_set
grpc_pollset_set正是为这种"一组 fd 需要被一组(动态变化的)pollset 轮询"的场景而设计。它语义上支持三层组合(API 定义见 grpc-polling-engines.md):
grpc_pollset_set_add_fd(pss, fd):把 fd 加入集合;grpc_pollset_set_add_pollset(pss, ps):关键语义——一旦 pollset 加入集合,对该 pollset 调用grpc_pollset_work()时,也会顺带轮询集合内所有的 fd(相当于把集合内所有 fd 逻辑上加进该 pollset);pollset 被移除后该保证失效;grpc_pollset_set_add_pollset_set(bag, item):把item内的所有 fd 并入bag,实现集合嵌套。
当前仓库的实现位于 src/core/lib/iomgr/pollset_set.cc,这些 add 操作最终都会以合并/解算的方式作用到底层各 pollset 的 fd 列表上。
在异步 connect 场景中的具体用法如下(示意图见下):
- 调用
connect(fd, ...),若返回EWOULDBLOCK/EINPROGRESS,说明连接仍在进行中; - 将该 fd 封装的
grpc_fd通过grpc_pollset_set_add_fd(interested_parties, fdobj)加入连接发起方的"感兴趣集合"(源码见 src/core/lib/iomgr/tcp_client_posix.cc); - 对该 fd 调用
grpc_fd_notify_on_write(fd, &write_closure)注册可写通知("armed"),见 tcp_client_posix.cc——当底层connect()完成、fd 变可写时,on_writableclosure 被触发,随后在finish:分支中执行grpc_pollset_set_del_fd(...)并把 fdorphan掉(tcp_client_posix.cc); - 连接期间还挂了一个 deadline 定时器,超时由
tc_on_alarmclosure 处理,避免连接永远悬挂。
需要说明的是:connect 使用的interested_parties集合属于 iomgr 的客户端连接发起层(src/core/lib/iomgr/tcp_client.cc 等),它会与真正驱动轮询的 pollset(例如与负载均衡/建连阶段绑定的 pollset_set)通过 add-pollset 关联,从而保证"有轮询线程在监控这些在途 connect fd"。一旦 fd 可写事件被处理、连接完成,fd 便从集合中摘除,子通道随后正式进入后续的 Call 关联流程。
四、gRPC 服务端:监听 fd、SO_REUSEPORT 与新连接分配
4.1 监听 fd 被加入每一个服务端 CQ 的 pollset
服务端的事件模型与客户端"一个 Call 只关心一个 fd"不同,它需要同时响应任意多个新连接请求,因此采用"广播"式注册:
- 服务端监听 fd(即对应监听端口号的 socket fd)会被加入每一个服务端 Completion Queue的 pollset 中;
- 注意 gRPC 使用
SO_REUSEPORT选项,可以创建多个监听 fd,但所有这些 fd 都映射到同一个监听端口——这相当于把内核的 accept 队列拆分成多个,供多个轮询线程并行消费。
在 src/core/lib/iomgr/tcp_server_posix.cc 中可以看到对应实现:启动监听时,代码遍历服务端所有 pollset,逐一执行grpc_pollset_add_fd((*pollsets)[i], sp->emfd)把监听 fd 加入每个 pollset,并通过grpc_fd_notify_on_read(sp->emfd, &sp->read_closure)注册可读通知;on_read回调则循环grpc_accept4()尽可能多地接受新连接。而当so_reuseport开启且 pollset 数量大于 1 时,会先经clone_port()克隆出与 pollset 数匹配的多个监听 fd(见 tcp_server_posix.cc),实现每个轮询线程各持一个可独立 accept 的监听 fd。
4.2 新连接以 round-robin 方式分配到服务端 CQ
一个新到来的 incoming channel(连接)会被分配到某个服务端 Completion Queue,分配策略当前是在服务端所有 CQ 之间轮询(round-robin),而不是永远塞给第一个 CQ——这样能保证多个轮询线程(通常与多核对应)之间的负载基本均衡。
当前仓库实现中可以看到这个计数器式轮询:在 src/core/lib/iomgr/tcp_server_posix.cc 附近,代码通过
gpr_atm_no_barrier_fetch_add(&s->next_pollset_to_assign, 1) % s->pollsets->size()原子地递增next_pollset_to_assign并对 pollset 总数取模,从而选出本次负责"通知读取就绪"的 pollset;该计数器在grpc_tcp_server_create中被初始化为 0(tcp_server_posix.cc)。被选中的 pollset 会负责调度后续对该连接的读事件通知,而该连接真正对应用户 side 的轮询则由其所属 CQ 驱动。
说明:文档撰写时(v1.15.1)round-robin 逻辑位于
tcp_server_posix.cc中分配 accept 通知的路径;当前仓库因引入 EventEngine(UseEventEngineListener)等改造,相关分配逻辑分散在 accept 回调与 external connection handler 两条路径中,但"按 pollset 总数取模轮询"的核心策略保持一致。
五、服务端整体事件流转一图串联
综合第二、三、四节,可以把服务端的事件模型归纳为三步:
- 注册:启动时,监听 fd 被复制/克隆进所有服务端 CQ 的 pollset(有
SO_REUSEPORT时每个 pollset 甚至各有独立的监听 fd); - 发现:任意一个用户线程调用
grpc_completion_queue_next()时,会进入grpc_pollset_work()阻塞轮询自己 CQ 的 pollset,从而能够感知监听 fd 的可读事件(有新连接到达)或已建立连接的读写事件; - 分配:
on_read接受新连接后,通过原子计数器 round-robin 选定一个服务端 CQ 作为该连接后续读事件的"通知方",把新连接逐步挂到对应 CQ 的 pollset 上,等待grpc_cq_end_op机制将对应事件通知给等待的 worker。
这个模型保证了:无论有多少个并发 RPC、多少个轮询线程,服务端都无需为每个连接单独创建线程,全部复用应用通过grpc_completion_queue_next()提供的用户线程。
六、纵深:pollset 与 pollset_set 的接口速查(客户端与服务端通用)
为了让读者能够直接对照 iomgr 层源码(ev_epoll1_linux.cc、ev_poll_posix.cc、ev_apple.cc、iocp_windows.cc、pollset.cc、pollset_set.cc等)继续深入,下面汇总本主题涉及的核心 API 语义(完整版见 grpc-polling-engines.md):
grpc_fd相关
| API | 关键语义 |
|---|---|
grpc_fd_notify_on_read/write/error(fd, closure) | 注册一次性通知("arming");每个事件闭包恰好触发一次,触发后 fd 处于"unarmed",需再次注册才能再次收到通知 |
grpc_fd_shutdown(fd) | 立即以 error 调度所有当前(及未来)已注册的读写/出错闭包 |
grpc_fd_orphan(fd, on_done, release_fd, reason) | 释放结构;release_fd == nullptr时顺带close()底层 fd;非空时把底层 fd 交给调用方(例如 C-Ares DNS 解析器这种 fd 非 gRPC 拥有的场景) |
grpc_pollset相关
| API | 关键语义 |
|---|---|
grpc_pollset_add_fd(ps, fd) | 把 fd 加入 pollset;注意没有grpc_pollset_remove_fd,因为grpc_fd_orphan()已隐含完成移除 |
grpc_pollset_work(ps, worker, deadline) | 调用前必须持有 pollset 的 mutex;阻塞轮询,直到 deadline 到期、发现 fd 事件并调度闭包、或 worker 被 kick |
grpc_pollset_kick(ps, worker) | 强制 worker 从grpc_pollset_work()返回;worker == nullptr表示踢醒该 pollset 上任一活跃 worker |
grpc_pollset_set相关
| API | 关键语义 |
|---|---|
grpc_pollset_set_add/del_fd(pss, fd) | 增删集合中的 fd |
grpc_pollset_set_add/del_pollset(pss, ps) | 加入后,对该 pollset 做grpc_pollset_work()等价于同时轮询集合内全部 fd |
grpc_pollset_set_add/del_pollset_set(bag, item) | 集合嵌套,等价于把item中所有 fd 并入bag |
客户端的"在途 connect 监控"依赖后两张表的能力(pollset_set 动态聚合 + pollset_work 驱动),而服务端的"广播监听 + round-robin 分配"则依赖第二张表 + 第三张表的 add_fd 语义。若需了解各平台轮询实现(epoll1 的 neighborhood/root worker 机制、poll 的 level-trigger 处理、Windows 的 IOCP 模型),可继续阅读 grpc-polling-engines.md 及 epoll-polling-engine.md。
七、总结:一条规则贯穿两端
把客户端与服务端放在一起看,Polling Engine 的使用其实遵循同一条底层规则:
谁负责等待某组 fd 的事件,谁就把这组 fd(通过 pollset 或 pollset_set)挂到自己将要调用的 pollset 上,然后用
grpc_pollset_work()阻塞等待,靠grpc_pollset_kick()精确唤醒。
- 客户端一侧:每个 Call 的 fd 挂在所属 CQ 的 pollset 上(Call 生命周期内);尚未成型的 sub-channel connect fd 通过
grpc_pollset_set挂到建连阶段"感兴趣的"轮询集合上(连接建立前); - 服务端一侧:监听 fd 广播挂到所有服务端 CQ 的 pollset(配合
SO_REUSEPORT克隆监听 fd 实现多线程并行 accept),新连接再按 round-robin 指配给某个服务端 CQ,交由该 CQ 的 pollset 持续轮询。
理解这套机制,是深入阅读 gRPC core 中 iomgr(src/core/lib/iomgr/)、surface(src/core/lib/surface/completion_queue.cc)乃至上层的 channel/load-balancing 代码的重要前提。建议配合 grpc-cq.md 与 grpc-polling-engines.md 两篇姊妹文档一起阅读,即可对 gRPC 的事件驱动内核形成完整认识。
【免费下载链接】grpcC++ based gRPC (C++, Python, Ruby, Objective-C, PHP, C#)项目地址: https://gitcode.com/GitHub_Trending/gr/grpc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考