brpc 的 IO 模型深度解析:从 EventDispatcher 收消息、wait-free 发消息到 Socket 生命周期管理
2026/9/13 17:35:59 网站建设 项目流程

brpc 的 IO 模型深度解析:从 EventDispatcher 收消息、wait-free 发消息到 Socket 生命周期管理

【免费下载链接】brpcbrpc is an Industrial-grade RPC framework using C++ Language, which is often used in high performance system such as Search, Storage, Machine learning, Advertisement, Recommendation etc. "brpc" means "better RPC".项目地址: https://gitcode.com/GitHub_Trending/brpc/brpc

brpc 是一个工业级的 C++ RPC 框架,其 IO 层设计直接决定了高并发下的吞吐与延迟表现。本篇技术指南以 docs/cn/io.md 为骨架,结合 event_dispatcher.h、input_messenger.h、socket.h 与 socket.cpp 的源码实现,系统讲解 brpc 收消息、发消息与 Socket 生命周期的完整链路。读完本文,你将掌握 brpc 选择 non-blocking IO 的工程动机、EDISP + bthread 的收包并发模型、wait-free MPSC 链表的发包原理,以及 SocketId/SocketUniquePtr 的内存管理设计。

三种 IO 方式:为什么 brpc 选择 non-blocking

计算机系统里操作 IO 的方式通常有三种:

  • blocking IO(阻塞 IO):发起 IO 操作后阻塞当前线程,直到 IO 结束。这是标准的同步 IO,例如默认行为下的 posixread/write系统调用。其实现完全由内核负责,read/write这类系统调用经过高度优化,在 IO 并发度很低时效率甚至高于需要多线程协作的 non-blocking IO。
  • non-blocking IO(非阻塞 IO):发起 IO 操作后不阻塞,用户可以阻塞等待多个 IO 操作同时结束。它本质上也是一种同步 IO,可以理解为"批量的同步"。典型代表是 Linux 下的pollselectepoll,以及 BSD 下的kqueue
  • asynchronous IO(异步 IO):发起 IO 操作后不阻塞,用户需要递一个回调,待 IO 结束后回调被调用。典型代表是 Windows 下的OVERLAPPED+IOCP。需要注意的是,Linux 的 native AIO 只对文件有效,对网络 socket 并不适用。

Linux 上通常使用 non-blocking IO 来提高 IO 并发度。为什么?当 IO 并发度很低时,blocking IO 完全由内核负责,read/write已被高度优化,效率高于多线程协作的 non-blocking IO。但当 IO 并发度提高后,blocking IO 阻塞一个线程的弊端就暴露出来:

  • 内核不得不持续在线程间切换才能完成有效的工作,一个 CPU core 上可能只做了一点点事情就马上切换到另一个线程,CPU cache 得不到充分利用;
  • 大量线程会使依赖 thread-local 加速的代码性能明显下降,例如 tcmalloc——一旦 malloc 变慢,程序整体性能往往随之下降。

而 non-blocking IO 一般由少量 event dispatching 线程一些运行用户逻辑的 worker 线程组成。这些线程往往会被复用(调度工作转移到了用户态),event dispatching 和 worker 可以同时在不同核上运行(流水线化),内核不用频繁切换就能完成有效工作;线程总量也不用很多,对 thread-local 的使用比较充分。此时 non-blocking IO 往往比 blocking IO 更快。

不过 non-blocking IO 也有自己的代价:

  • 需要调用更多系统调用,比如epoll_ctl。由于 epoll 内部实现为一棵红黑树,epoll_ctl并不是一个很快的操作,特别是在多核环境下,依赖epoll_ctl的实现往往会面临棘手的扩展性问题;
  • non-blocking 需要更大的缓冲,否则会触发更多的事件而影响效率;
  • 还得解决不少多线程问题,代码比 blocking 复杂很多。

brpc 正是在这种权衡下,选择了 non-blocking IO 作为网络层的基础,并围绕它设计了一整套消息收发机制。

收消息:EventDispatcher 与 bthread 的协作

EDISP 是什么

"消息"指从连接读入的有边界的二进制串,可能是来自上游 client 的 request,或来自下游 server 的 response。brpc 使用一个或多个EventDispatcher(简称EDISP)等待任一 fd 发生事件。

与常见的"IO 线程"不同,EDISP 不负责读取。IO 线程的问题在于:一个线程同时只能读一个 fd,当多个繁忙的 fd 聚集在一个 IO 线程中时,一些读取就被延迟了。多租户、复杂分流算法、Streaming RPC 等功能会加重这个问题;高负载下常见的某次读取卡顿会拖慢一个 IO 线程中所有 fd 的读取,对可用性的影响幅度较大。

在 event_dispatcher.h 中可以看到EventDispatcher的核心接口:AddConsumer把 fd 挂到内部 epoll 上(注释明确说明"Dispatch edge-triggered events of file descriptors to consumers",即分发edge-triggered(边沿触发)事件),RegisterEvent/UnregisterEvent用于动态增删 EPOLLOUT 监听,Start则把 dispatcher 本身作为一个 bthread 启动(event_dispatcher.h)。全局 dispatcher 的数量由FLAGS_event_dispatcher_num控制,并按 bthread tag 分组,见 event_dispatcher.cpp。

Edge triggered 与 wait-free 的事件消费

brpc 选择 Edge triggered(边沿触发)模式,原因有二:

  1. 规避 epoll 的一个 历史 bug(开发 brpc 时仍存在);
  2. 减少epoll_ctl带来的较大开销。

当收到事件时,EDISP 给一个原子变量加 1,只有当加 1 前的值是 0 时才启动一个 bthread 处理对应 fd 上的数据。在背后,EDISP 把所在的 pthread 让给了新建的 bthread,使其有更好的 cache locality,可以尽快地读取 fd 上的数据;而 EDISP 所在的 bthread 会被偷到另外一个 pthread 继续执行,这个过程就是 bthread 的work stealing 调度

要准确理解那个原子变量的工作方式,可以先阅读 atomic_instructions.md,再看Socket::StartInputEvent(位于 socket.cpp)。这些方法使得 brpc 读取同一个 fd 时产生的竞争是wait-free的(关于 wait-free 的定义,可参见非阻塞算法领域的经典分类)。

在当前实现里,Transport::ProcessEvent会按EventDispatcherUnsched()选择启动方式:

  • 返回false时走bthread_start_urgent(前台调度,先于调用者继续执行);
  • 返回true时走bthread_start_background(后台调度,允许被调度出去)。

EventDispatcherUnsched()直接读取 gflags 开关event_dispatcher_edisp_unsched(默认false),定义见 event_dispatcher.cpp,该 flag 的注释为 "Disable event dispatcher schedule",用户可通过命令行-event_dispatcher_edisp_unsched控制这一行为。

此外,RDMA 在轮询模式与事件模式下对last_msg的处理不同:rdma_use_polling=false时不会在RdmaTransport::QueueMessage里处理last_msg,轮询模式下会继续处理。并且在EventDispatcherUnsched()返回true时,last_msg不会在当前执行流里直接处理,而是在新的 bthread 中执行。RDMA 相关的整体背景可参考 rdma.md。

InputMessenger:从 fd 上切割并处理消息

InputMessenger(input_messenger.h)负责从 fd 上切割和处理消息,它通过用户回调函数理解不同的格式。回调在InputMessageHandler中定义(input_messenger.h):

  • Parse:把消息从二进制流上切割下来,运行时间较固定;
  • Process:进一步解析消息(比如反序列化为 protobuf)后调用用户回调,时间不确定;
  • Verify:仅在该 socket 收到的第一条消息上调用,用于鉴权,可空。

若一次从某个 fd 读取出 n 个消息(n > 1),InputMessenger 会启动n-1 个 bthread分别处理前 n-1 个消息,最后一个消息则会在原地被Process。这一行为在 input_messenger.cpp 的OnNewMessages中有明确注释:所有消息都在当前 bthread 中被 Parse(即从butil::IOBuf上切下来,protobuf 反序列化属于 "process" 阶段),除最后一条外的消息放入独立 bthread 处理,且为了最小化开销,调度是批量的(使用BTHREAD_NOSIGNALbthread_flush)。在 input_messenger_processor.cpp 中可以看到特殊处理:RDMA / UBRING 模式下last_msg也会通过QueueMessage放入新 bthread 执行,因为处理消息的方法可能调用同步原语,导致轮询 bthread 被调度出去。

InputMessenger 会逐一尝试多种协议。由于一个连接上往往只有一种消息格式,它会记录下上次的选择,避免每次都重复尝试(FindProtocolIndex/NameOfProtocol维护协议与 handler 的映射,见 input_messenger.h)。

可以看到,fd 之间和 fd 内部的消息都会在 brpc 中获得并发,这使 brpc 非常擅长大消息的读取,在高负载时仍能及时处理不同来源的消息,减少长尾的存在。整个端到端流程(Client 侧 Channel → LB → Socket → 网络 → Server 侧 Acceptor → Socket → Service)可参考下图的完整链路示意:

发消息:wait-free MPSC 链表与 KeepWrite 线程

"消息"指向连接写出的有边界的二进制串,可能是发向上游 client 的 response 或下游 server 的 request。多个线程可能会同时向一个 fd 发送消息,而写 fd 又是非原子的,所以如何高效率地排队不同线程写出的数据包是这里的关键。

brpc 使用一种wait-free MPSC(多生产者单消费者)链表来实现这个功能,核心代码在 socket.cpp 的StartWrite(约 L1700-L1712):

  • 所有待写出的数据都放在一个单链表节点中,next 指针初始化为一个特殊值Socket::WriteRequest::UNCONNECTED
  • 当一个线程想写出数据前,它先尝试和对应的链表头Socket::_write_head原子交换,返回值是交换前的链表头;
  • 如果返回值为空,说明它获得了写出的权利,它会在原地写一次数据;
  • 否则说明有另一个线程在写,它把 next 指针指向返回的头以让链表连通。正在写的线程之后会看到新的头并写出这块数据。

源码中req->next = WriteRequest::UNCONNECTED在每次入队时被设置(socket.cpp),随后_write_head.exchange(req, butil::memory_order_release)完成入队;当prev_head != nullptr时,把req->next = prev_head接回链表并立即返回(socket.cpp)。

这套方法可以让写竞争是 wait-free 的。而获得写权利的线程虽然在原理上不是 wait-free 也不是 lock-free——它可能会被一个值仍为UNCONNECTED的节点锁定(这需要发起写的线程正好在原子交换后、设置 next 指针前、仅仅一条指令的时间内被 OS 换出)——但在实践中很少出现(源码注释明确指出在高竞争测试中几乎观察不到自旋)。

在当前的实现中,如果获得写权利的线程一下子无法写出所有的数据,会启动一个KeepWrite 线程继续写,直到所有的数据都被写出(socket.cpp 中bthread_start_background以 "KeepWrite" 命名启动)。这套逻辑非常复杂,大致原理如下图所示,细节可阅读 socket.cpp:

由于 brpc 的写出总能很快地返回,调用线程可以更快地处理新任务;后台 KeepWrite 写线程每次拿到一批任务批量写出,在大吞吐时容易形成流水线效应而提高 IO 效率

Socket:用 64 位 SocketId 管理 fd 的一生

和 fd 相关的数据均在Socket(socket.h)中,是 RPC 最复杂的结构之一。这个结构的独特之处在于:用 64 位的 SocketId 指代 Socket 对象,以方便在多线程环境下使用 fd。SocketId定义于 socket_id.h,本质是VRefId(带版本号的引用 ID),SocketUniquePtr则是其对应的自动释放指针。

常用的三个方法:

  • Create:创建 Socket,并返回其 SocketId。
  • Address:取得 id 对应的 Socket,包装在一个会自动释放的 unique_ptr 中(SocketUniquePtr)。当 Socket 被SetFailed后,返回指针为空。只要 Address 返回了非空指针,其内容保证不会变化,直到指针自动析构。这个函数是wait-free的。
  • SetFailed:标记一个 Socket 为失败,之后所有对那个 SocketId 的 Address 会返回空指针(直到健康检查成功)。当 Socket 对象没人使用后会被回收。这个函数是lock-free的。

可以看到 Socket 类似shared_ptr,SocketId 类似weak_ptr,但 Socket 独有的SetFailed可以在需要时确保 Socket 不能被继续 Address 而最终引用计数归 0。单纯使用shared_ptr/weak_ptr则无法保证这点——当一个 server 需要退出时,如果请求仍频繁地到来,对应 Socket 的引用计数可能迟迟无法清 0 而导致 server 无法退出。另外weak_ptr无法直接作为 epoll 的 data,而 SocketId 可以(epoll data 里可以直接存放 64 位整数)。这些因素促使 brpc 设计了 Socket 这个类,其核心部分自 2014 年完成后很少改动,非常稳定。

存储SocketUniquePtr还是SocketId取决于是否需要强引用

  • Controller贯穿了 RPC 的整个流程,和 Socket 中的数据有大量交互,它存放的是SocketUniquePtr
  • epoll 主要是提醒对应 fd 上发生了事件,如果 Socket 回收了,那这个事件是可有可无的,所以它存放的是SocketId

由于SocketUniquePtr只要有效,其中的数据就不会变,这个机制使用户不用关心麻烦的 race condition 和 ABA problem,可以放心地对共享的 fd 进行操作。这种方法也规避了隐式的引用计数,内存的 ownership 明确,程序的质量有很好的保证。brpc 中有大量的SocketUniquePtrSocketId,它们确实简化了开发。

值得强调的是,Socket 不仅仅用于管理原生的 fd,它也被用来管理其他资源:

  • SelectiveChannel中的每个 Sub Channel 都被置入了一个 Socket 中,这样 SelectiveChannel 可以像普通 channel 选择下游 server 那样选择一个 Sub Channel 进行发送,这个"假 Socket"甚至还实现了健康检查;
  • Streaming RPC 也使用了 Socket,以复用 wait-free 的写出过程。

全链路视图与进一步阅读

把收消息、发消息与 Socket 生命周期串起来,brpc 的一次 RPC 在 IO 层面经历了:Client 侧 Channel 通过命名服务与负载均衡选中目标 Socket → 经 wait-free 链表与 KeepWrite 写出请求 → Server 侧 Acceptor 接受连接、EventDispatcher 边沿触发分发事件、InputMessenger 切割并并发处理消息 → Service 执行业务 → 响应再沿同样的链路返回 Client,全程 fd 间与 fd 内均保持并发,这正是 brpc 在搜索、存储、机器学习等高吞吐场景下表现出色的 IO 基础。

如果想继续深入 IO 相关的其他主题,建议按以下路径阅读:

  • bthread.md:收消息链路依赖的调度原语,理解 work stealing 与 bthread 切换;
  • atomic_instructions.md:理解 EDISP 原子变量与 wait-free 写链表的基础;
  • iobuf.md:消息切割与缓冲的基础数据结构;
  • streaming_rpc.md:复用 Socket 与 wait-free 写出的流式 RPC;
  • rdma.md:RDMA 模式下last_msg与轮询/事件模式的行为差异;
  • 核心实现源码:event_dispatcher.h、event_dispatcher.cpp、input_messenger.h、input_messenger.cpp、socket.h、socket.cpp、socket_id.h。

【免费下载链接】brpcbrpc is an Industrial-grade RPC framework using C++ Language, which is often used in high performance system such as Search, Storage, Machine learning, Advertisement, Recommendation etc. "brpc" means "better RPC".项目地址: https://gitcode.com/GitHub_Trending/brpc/brpc

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询