
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/brpcbrpc 是一个工业级的 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 下的poll、select、epoll以及 BSD 下的kqueue。asynchronous IO异步 IO发起 IO 操作后不阻塞用户需要递一个回调待 IO 结束后回调被调用。典型代表是 Windows 下的OVERLAPPEDIOCP。需要注意的是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边沿触发模式原因有二规避 epoll 的一个 历史 bug开发 brpc 时仍存在减少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_pollingfalse时不会在RdmaTransport::QueueMessage里处理last_msg轮询模式下会继续处理。并且在EventDispatcherUnsched()返回true时last_msg不会在当前执行流里直接处理而是在新的 bthread 中执行。RDMA 相关的整体背景可参考 rdma.md。InputMessenger从 fd 上切割并处理消息InputMessengerinput_messenger.h负责从 fd 上切割和处理消息它通过用户回调函数理解不同的格式。回调在InputMessageHandler中定义input_messenger.hParse把消息从二进制流上切割下来运行时间较固定Process进一步解析消息比如反序列化为 protobuf后调用用户回调时间不确定Verify仅在该 socket 收到的第一条消息上调用用于鉴权可空。若一次从某个 fd 读取出 n 个消息n 1InputMessenger 会启动n-1 个 bthread分别处理前 n-1 个消息最后一个消息则会在原地被Process。这一行为在 input_messenger.cpp 的OnNewMessages中有明确注释所有消息都在当前 bthread 中被 Parse即从butil::IOBuf上切下来protobuf 反序列化属于 process 阶段除最后一条外的消息放入独立 bthread 处理且为了最小化开销调度是批量的使用BTHREAD_NOSIGNAL与bthread_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 相关的数据均在Socketsocket.h中是 RPC 最复杂的结构之一。这个结构的独特之处在于用 64 位的 SocketId 指代 Socket 对象以方便在多线程环境下使用 fd。SocketId定义于 socket_id.h本质是VRefId带版本号的引用 IDSocketUniquePtr则是其对应的自动释放指针。常用的三个方法Create创建 Socket并返回其 SocketId。Address取得 id 对应的 Socket包装在一个会自动释放的 unique_ptr 中SocketUniquePtr。当 Socket 被SetFailed后返回指针为空。只要 Address 返回了非空指针其内容保证不会变化直到指针自动析构。这个函数是wait-free的。SetFailed标记一个 Socket 为失败之后所有对那个 SocketId 的 Address 会返回空指针直到健康检查成功。当 Socket 对象没人使用后会被回收。这个函数是lock-free的。可以看到 Socket 类似shared_ptrSocketId 类似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 中的数据有大量交互它存放的是SocketUniquePtrepoll 主要是提醒对应 fd 上发生了事件如果 Socket 回收了那这个事件是可有可无的所以它存放的是SocketId。由于SocketUniquePtr只要有效其中的数据就不会变这个机制使用户不用关心麻烦的 race condition 和 ABA problem可以放心地对共享的 fd 进行操作。这种方法也规避了隐式的引用计数内存的 ownership 明确程序的质量有很好的保证。brpc 中有大量的SocketUniquePtr和SocketId它们确实简化了开发。值得强调的是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 写出的流式 RPCrdma.mdRDMA 模式下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),仅供参考