在学完 select、poll、epoll 之后,很多同学会陷入一种"手里拿着锤子,看什么都是钉子"的错觉:既然 epoll 能同时监听几百上千个文件描述符,那服务器是不是就写完了?答案当然没那么简单。
epoll 只解决了一半问题——它帮你"检测"哪些文件描述符就绪了。但"检测到了之后干什么、谁来干、怎么组织代码让它既高效又不容易写崩",这是另一门手艺。而这门手艺的核心,就叫 Reactor 反应堆模式。
这篇文章我会从"传统服务器为什么不行"讲起,一路走到:什么是事件驱动、什么是事件循环、Reactor 模式的定义和四个核心角色、如何用 epoll 亲手实现一个 Reactor(连接事件、读事件、写事件、异常事件的分发与回调)、事件回调如何解耦,再到单线程与多线程 Reactor 的各种变体(包括 OTOL 和 eventfd)。老规矩,全程干货,代码逐行注释,可编译可运行。
你应该有的知识准备
往下走之前,先确认你有这几块地基,看不懂哪条可以回去翻之前的章节:
epoll的三件套:epoll_create(建内核事件表)、epoll_ctl(增删改关注的事件)、epoll_wait(等待事件就绪)。这块不懂,Reactor 就是空中楼阁。- 非阻塞 IO:
fcntl把一个 socket 设成非阻塞。Reactor 之所以能"监而不抢",正是因为读写都是非阻塞的。 std::function与std::bind:C++ 里把"一段逻辑"打包成可调用对象的手法,回调机制的地基。- TCP 的流式特性与缓冲区:字节流没有天然边界,这是"半包/粘包"问题的源头。
- 多线程与同步:理解单线程 Reactor 之后扩展多线程时会用到,重点是"临界区"和"唤醒"。
从"一个连接一个线程"说起:为什么需要 Reactor
先看一个最朴素的服务端写法:阻塞式多线程服务器。
// thread_per_conn.c —— 传统"一连接一线程",注意观察它的问题
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <pthread.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <arpa/inet.h>
// 每个连接一个线程的处理函数:死循环读、回声发回去
static void *handle_conn(void *arg)
{
int cfd = *(int *)arg; // 取出客户端 socket
free(arg); // 释放拷贝过来的参数(指针指向堆上空间)
char buf[1024];
ssize_t n;
while ((n = read(cfd, buf, sizeof(buf))) > 0) // 阻塞读
{
write(cfd, buf, (size_t)n); // 原样写回
}
close(cfd); // 客户端断开,读返回 0,关闭连接
return NULL;
}
int main(void)
{
int lfd = socket(AF_INET, SOCK_STREAM, 0);
int opt = 1;
setsockopt(lfd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt));
struct sockaddr_in addr;
memset(&addr, 0, sizeof(addr));
addr.sin_family = AF_INET;
addr.sin_port = htons(8888);
addr.sin_addr.s_addr = INADDR_ANY;
bind(lfd, (struct sockaddr *)&addr, sizeof(addr));
listen(lfd, 128);
printf("server listen on 8888 ...\n");
while (1)
{
int cfd = accept(lfd, NULL, NULL); // 阻塞等待新连接
int *p = malloc(sizeof(int)); // 参数要拷贝到堆上,避免传栈上指针
*p = cfd;
pthread_t tid;
pthread_create(&tid, NULL, handle_conn, p); // 每个连接开一个线程
pthread_detach(tid); // 分离线程,用完自动回收
}
close(lfd);
return 0;
}这段代码的问题是"结构上的":
- 线程是稀缺且昂贵的资源。一个线程默认栈空间 8MB(可配),创建和切换都有开销。一万个连接就要一万个线程,光内存就得上百 GB,操作系统根本调不动。
- 大量线程阻塞在
read上。绝大多数 TCP 连接是"大部分时间闲着的",一个线程却为了等这点数据白白占着一个内核线程在调度队列里转。 - 同步复杂度爆炸。多线程要共享全局数据,就得加锁,锁一多,性能不升反降。
于是大家想到一个反直觉的思路:既然连接大多数时候都不活跃,那我们不要给每个连接配一个执行流,反正它们是"等待事件"的人。让一个(或少数几个)执行流专心去"监听『谁就绪了』这个事件",谁就绪了再临时去处理谁。 这就是事件驱动(Event-Driven)编程的出发点。
事件驱动的口号是:别让 CPU 去轮询每个连接"你有没有数据"(那是浪费),而是让内核说"某个连接有数据了",执行流再被唤醒去处理。 谁来"监听和分发"这些就绪通知,就是下一节的事件循环。
事件驱动编程与事件循环
事件驱动编程(Event-Driven Programming)是一种程序控制流模式:程序不是"主动按顺序执行每一步",而是"被动地等待事件发生,谁来了就处理谁"。事件可以来自网络数据到达、定时器到期、信号触发、管道可读等等。
而事件循环(Event Loop),就是驱动这套机制的那个"永不落幕的主循环"。它大致长这样:
┌──────────────────────────────┐
│ │
│ ┌──────────────────┐ │
│ │ 事件集合(epoll) │ │
│ └──────────────────┘ │
│ │ │
│ 阻塞等待事件就绪 │
│ │ │
│ ▼ │
│ ┌────────────────────┐ │
│ │ 拿到就绪的 fd 集合 │ │
│ └────────────────────┘ │
│ │ │
│ ▼ │
│ ┌────────────────────┐ │
│ │ 分发给对应处理函数 │ │
│ └────────────────────┘ │
│ │ │
│ └───────────────┘ ← 处理完回到等待
│
└──── 程序结束或被停止 ────
用代码抽象,事件循环的骨架是:
// event_loop_skeleton.c —— 事件循环的通用骨架(伪结构化示意)
while (1)
{
// 1. 阻塞等待:把"关注的事件"交给内核,内核就绪了才返回
int n = epoll_wait(epfd, events, MAX_EV, -1);
// 2. 逐个分发:每一个就绪的 fd,找到它的"处理函数"去调用
for (int i = 0; i < n; i++)
{
Handler *h = fd_to_handler(events[i].data.fd); // 由 fd 找回处理者
h->handle(events[i].events); // 把事件类型交给它
}
// 3. 一圈结束,继续回到等待——这就是"永动的循环"
}打个比方:事件循环像一个"前台接线员"。他不用给自己配一万个分身去伺候一万个电话(那每个电话都要占一个人),他只需要:坐在工位上(阻塞等待),哪个电话响了(事件就绪),他接起来(取出 fd),转给对应的部门(分发到处理函数),然后挂断接着等下一个。接线员只有一个人,却能同时"管着"成千上万个座机。
关键点在于:事件循环天生是单线程的——至少内心是。它在一个线程里"监听 → 分发 → 处理 → 再监听",就好像用一个人的精力驱动了整个服务器。这个"以少数几个线程承载海量连接"的模型,正是 Reactor 的灵魂。
这里先讲清楚站在门口就要懂的术语:
- 事件(Event):程序关心的某种状态变化,比如"某个 socket 可读了""某个 socket 可写了""到了一分钟"。在 Linux 网络场景下,事件大多来自
epoll的事件位,像EPOLLIN(可读)、EPOLLOUT(可写)、EPOLLERR(出错)。 - 事件循环(Event Loop):上面那个"监听-分发-处理"的永动循环。它是 Reactor 的中枢神经,通常也叫 Handler / EventLoop。在 libevent、libuv、Nginx、Redis、Netty 这些框架里,它都有名字(
event_base、EventLoop、event_base),但本质都是同一种东西。
什么是 Reactor 反应堆模式
Reactor 模式(Reactor Pattern,也叫反应堆模式 / 反应器模式)是一种事件驱动的服务器设计范式。它的核心思想是:I/O 多路复用负责"检测谁就绪",Reactor 负责"把就绪事件分发给对应的处理组件去执行",从而用少数执行流处理海量并发连接,且无需为每个连接创建一个线程或进程。
"反应堆"这个名字起得贴切:核反应堆里,铀原子裂变释放中子,中子再去触发别的铀原子裂变,形成链式反应。这里的"事件"就像中子——一个事件到来,经过分发,触发相应的处理逻辑,处理逻辑又可能产生新事件(比如收到请求要写回响应),形成一条"反应链"。
Reactor 模式涉及四个经典角色,我逐个讲清楚,这是理解全文的钥匙:
| 角色 | 名字 | 职责 |
|---|---|---|
| 事件源 | Handle(这些 fd/socket) | 产生了 I/O 事件的东西 |
| 事件多路复用器 | Synchronous Event Demultiplexer | 阻塞地等待事件就绪并返回哪些就绪了(Linux 下就是 epoll) |
| Reactor(反应器) | 也叫 EventLoop | 事件循环,管理事件注册/注销,并把就绪事件分发给对应 Handler |
| 事件处理器 | Handler | 处理具体某类事件(连接到来、可读、可写、错误) |
它们的关系可以画成这样一张图:
┌────────────────────────────────────────────────┐
│ Reactor (反应器 / EventLoop) │
│ │
│ 注册事件 ┌──────────────┐ │
│ ───────────▶ │ 事件表(epoll)│ │
│ └──────────────┘ │
│ │ ▲ │
│ epoll │ │ 就绪事件的队列/fd集合 │
│ 等待 │ │ │
│ ▼ │ │
│ Synchronous Event Demultiplexer (epoll) │
│ └───────┘ │
└───────────────────┬────────────────────────────┘
│ 分发(根据就绪的 fd 找到对应 Handler)
▼
┌────────────────────────────────────────────┐
│ HandlerA ┌─────────┐ │
│ (Accepter)│ (连接事件)│ → accept 新连接 │
│ HandlerB ┌─────────┐ │
│ (Recver) │ (可读事件)│ → 收数据、解析 │
│ HandlerC ┌─────────┐ │
│ (Sender) │ (写事件) │ → 发数据 │
│ HandlerD ┌─────────┐ │
│ (Excepter)│ (异常) │ → 清理连接 │
└────────────────────────────────────────────┘
注意区分两个容易混淆的词:
- Demultiplexer(多路解复用器 / 事件分离器):它负责"从一大堆东西里挑出就绪的那几个"。
select、poll、epoll都是它。它把"一大堆 fd"和多路复用(multiplex,多路合一路)反过来——一个执行流同时听很多路,谁就绪就把谁"解"出来。 - Reactor(反应器):收到 Demultiplexer 返回的就绪列表后,负责"找对应的 Handler 并调用"。它并不真正做 I/O,它只做"分发 + 组织"。
一句话记忆:epoll 负责"谁就绪了",Reactor 负责"把『你』喊过来干活"。
用 epoll 实现 Reactor 的心脏:事件循环
纸上谈兵结束,我们动手。先写 Reactor 的心脏——事件循环。它需要一个核心数据结构:把 fd 对应到它该由谁来处理的一张映射表。因为 epoll_wait 只告诉你"哪个 fd 就绪了",你必须有法子从 fd 反查回来"它的处理器是谁"。
epoll_event 里有个 data 联合体,我们常用 data.ptr 或 data.fd 来存这个"反查线索"。这里我用一个 std::unordered_map<int, Connection*> 来存"fd → 连接对象"的关系,既直观又安全。
先写一个可以独立编译运行的单线程 Reactor 回声服务器(echo,收什么回什么)。我会把它拆开讲,最后给完整可编译版本。
连接对象:把 socket、事件、回调绑在一起
在 C++ 里,我们很自然地用一个类来代表"一个连接",把所有跟这个连接绑定的东西(socket、收/发缓冲区、关心的事件、处理回调)封在一起。这是 Reactor 里"回调解耦"的第一步——数据的搬运工是 Reactor,而"这段数据该怎么处理"是回调函数,二者互不关心。
// Connection.hpp —— 一个连接对象:socket + 缓冲区 + 事件 + 回调
#pragma once
#include <functional> // std::function:用来装"处理函数"
#include <string>
#include <sys/socket.h> // socket 相关
#include <netinet/in.h> // sockaddr_in
#include <unistd.h> // close
class Connection; // 前置声明,下面回调签名要引用它
class TcpServer; // 前置声明,Connection 需要回指服务器
// func_t:一个"接收 Connection* 并处理它"的函数类型
// std::function 让我们能把成员函数 / lambda / 绑定器都存成一个对象
using func_t = std::function<void(Connection *)>;
class Connection
{
public:
// 构造:记住这个连接的 socket、关注的事件、以及它属于哪个服务器
Connection(int sockfd, uint32_t events, TcpServer *R)
: _sockfd(sockfd), _events(events), _R(R) {}
// 注册三个回调:收到数据怎么处理(recver)、可写时怎么发(sender)、异常/断开怎么办(excepter)
void RegisterCallback(func_t recver, func_t sender, func_t excepter)
{
_recver = recver;
_sender = sender;
_excepter = excepter;
}
// 把读到的字节追加进接收缓冲区。追加而非覆盖,是为了攒"半包"。
void AddInBuffer(const std::string &buffer) { _inbuffer += buffer; }
// 把要发的数据追加进发送缓冲区。这里是"攒着等 EPOLLOUT",是 Reactor 写事件的核心。
void AddOutBuffer(const std::string &buffer) { _outbuffer += buffer; }
bool OutBufferEmpty() const { return _outbuffer.empty(); } // 发送缓冲区是否空
int SockFd() const { return _sockfd; } // 取出这个连接的 socket
uint32_t Events() const { return _events; } // 取出关注的事件
void SetEvents(uint32_t events) { _events = events; } // 更新关注的事件
std::string &InBuffer() { return _inbuffer; } // 取出接收缓冲区(引用,允许外部改)
std::string &OutBuffer() { return _outbuffer; } // 取出发送缓冲区(引用)
void Close() { ::close(_sockfd); } // 关闭 socket
private:
int _sockfd; // 这个连接对应的 socket fd
std::string _inbuffer; // 接收缓冲区(先用 string 顶替,工程上会用更精细的 buffer)
std::string _outbuffer; // 发送缓冲区
uint32_t _events; // 这个连接当前关注的事件(EPOLLIN/EPOLLOUT/...)
struct sockaddr_in _client; // 对端 ip/port,便于日志
public:
func_t _recver; // 收数据时的回调
func_t _sender; // 可写时的回调
func_t _excepter; // 出错/断开时的回调
TcpServer *_R; // 回指服务器:方便调用服务器提供的"注册/移除事件"接口
};这里有个很关键的"职责分离":Connection 自己只管数据和状态,不写具体的业务逻辑。业务逻辑全在回调里(_recver 等),这个回调从外面"注册"进来。这就把**"数据怎么拿"(Reactor 的活)和"拿到数据以后干什么"(业务方的活)彻底解耦了——这就是回调解耦(Callback Decoupling)**。
连接工厂:监听连接和普通连接
Connection 有两种用途:监听 socket(只关注连接到来)和普通数据连接(关注读写)。工程上用工厂模式(一个专门负责"怎么造对象"的类)来区分:
// ConnectionFactory —— 工厂:屏蔽"怎么构造不同类型 Connection"的细节
class ConnectionFactory
{
public:
// 造一个"监听连接":只有 recver(连接到来时用),sender/excepter 暂时为空
static Connection *BuildListenConnection(int listensock, func_t recver,
uint32_t events, TcpServer *R)
{
Connection *conn = new Connection(listensock, events, R);
conn->RegisterCallback(recver, nullptr, nullptr);
return conn;
}
// 造一个"普通数据连接":三个回调都有
static Connection *BuildNormalConnection(int sockfd,
func_t recver, func_t sender, func_t excepter,
uint32_t events, TcpServer *R)
{
Connection *conn = new Connection(sockfd, events, R);
conn->RegisterCallback(recver, sender, excepter);
return conn;
}
};事件循环与 TcpServer:注册、分发、管理
接下来是核心的 TcpServer(它同时扮演 Reactor 和 Demultiplexer 的使用者)。它负责:建 epoll、把连接注册进 epoll、维护"fd → Connection"映射、运行事件循环、在就绪事件到来时"分发"给对应连接的回调。
// TcpServer.hpp —— Reactor 的主控类
#pragma once
#include "Connection.hpp"
#include <sys/epoll.h>
#include <unordered_map>
class TcpServer
{
public:
TcpServer(); // 建 socket、bind、listen、建 epoll,注册监听连接
~TcpServer();
void InitServer(); // 初始化 + 启动事件循环
void Loop(); // 事件循环本体
// 对外接口:添加/移除一个连接;开启/关闭某 fd 对读写事件的关注
void AddConnection(Connection *conn); // 新的数据连接进来
void RemoveConnection(int sockfd); // 从 epoll 与映射中移除
void EnableReadWrite(int sockfd, bool rd, bool wr); // 调整某 fd 关注的事件位
private:
int _listensock; // 监听 socket
int _epfd; // epoll 实例句柄
std::unordered_map<int, Connection *> _conns; // fd → 连接对象 的映射
void SetNonBlock(int fd); // 把 fd 设为非阻塞,Reactor 的读写必须在非阻塞下进行
void BuildListen(); // 创建监听 socket 并注册到 epoll
};关键实现 Loop 就是事件循环本体,也是"分发"发生的现场:
void TcpServer::Loop() // 事件循环:监听 → 分发 → 再监听
{
const int MAX_EVENTS = 64; // 每次 epoll_wait 最多取回的就绪事件数
struct epoll_event events[MAX_EVENTS]; // 就绪事件数组
while (1) // 永动循环
{
// 阻塞等待就绪事件。-1 表示永久等待,直到有事件返回
int n = epoll_wait(_epfd, events, MAX_EVENTS, -1);
if (n < 0) // 出错
{
if (errno == EINTR) break; // 被信号打断:通常 continue,此处简化直接退出
continue;
}
for (int i = 0; i < n; i++) // 逐个分发就绪的 fd
{
int fd = events[i].data.fd; // 取出就绪的 fd(我们注册时塞的是 data.fd)
uint32_t ev = events[i].events; // 取出就绪的事件位
auto it = _conns.find(fd); // 从映射反查这一个 fd 对应哪个连接
if (it == _conns.end()) continue; // 找不到(可能已被移除),跳过
Connection *conn = it->second; // 取到连接对象
// ---- 分发阶段:根据事件位调用对应的回调 ----
if (ev & EPOLLIN) // 可读:包含"新连接到来"和"数据到达"两种情况
{
if (conn->Events() & EPOLLIN) // 看它注册时是不是带了可读关注
conn->_recver(conn); // 调用它的"可读处理回调"(Accepter 或 Recver)
}
if (ev & EPOLLOUT) // 可写:一般意味着发送缓冲区可以被写入
{
if (conn->Events() & EPOLLOUT)
conn->_sender(conn); // 调用"可写处理回调"(Sender)
}
// 出错或对端关闭:EPOLLERR 一定会带;EPOLLHUP 表示对方半关闭
if ((ev & EPOLLERR) || (ev & EPOLLHUP))
{
if (conn->_excepter) // 有异常回调才调用(监听连接没有 excepter)
conn->_excepter(conn);
continue; // 异常后连接该清理或跳过后续处理
}
}
}
}注意上面分发里有个非常重要的"坑":监听 socket 和普通数据 socket,在处理时可读事件时走的是同一个 EPOLLIN 位。怎么区分"这是要 accept 新连接"还是"这是要 recv 数据"?——靠的是它们各自注册进去的回调不同。监听 socket 的 _recver 是 Accepter(去 accept),普通 socket 的 _recver 是 Recver(去收数据)。这就是 Reactor 的妙处:用不同的回调,把"同一类事件"分叉成"不同的行为",而分发者(事件循环)根本不需要知道里面是什么。这一点在后面"回调解耦"会再强调。
Accepter:连接事件的到来
监听 socket 可读时,说明有新的 TCP 连接请求来了,我们要把它 accept 出来。这里有个 epoll 边缘触发(ET)模式的经典问题:ET 模式下事件只通知一次,你怎么保证把"这一批"连接都收干净? 答案是在 EPOLLIN 出现后,用 while 循环 accept 直到返回 EAGAIN(没新连接了)。
// Accepter —— 连接管理器:处理"新连接到来"这一事件
class Accepter
{
public:
Accepter() {}
// conn 是监听 socket 对应的那个 Connection
void AccepterConnection(Connection *conn)
{
while (1) // 得一直 accept 到没有为止
{
struct sockaddr_in peer; // 对端地址
socklen_t len = sizeof(peer);
int sockfd = ::accept(conn->SockFd(),
(struct sockaddr *)&peer, &len);
if (sockfd > 0) // 成功 accept 出一个新连接
{
SetNonBlock(sockfd); // 立刻设为非阻塞,后面 epoll 才能照管它
// 为这个普通连接准备三个回调(绑定到 HandlerConnection 的静态方法)
auto recver = std::bind(&HandlerConnection::Recver, _1);
auto sender = std::bind(&HandlerConnection::Sender, _1);
auto excepter = std::bind(&HandlerConnection::Excepter, _1);
// 用工厂造出一个普通连接对象:关注 可读 + ET(边缘触发)
Connection *normal_conn =
ConnectionFactory::BuildNormalConnection(sockfd, recver, sender, excepter,
EPOLLIN | EPOLLET, conn->_R);
conn->_R->AddConnection(normal_conn); // 交给服务器:注册进 epoll 和映射表
}
else
{
// accept 失败:根据 errno 判断该怎么处理
if (errno == EAGAIN) break; // 没有更多连接了,正常退出循环
else if (errno == EINTR) continue; // 被信号打断,重试
else break; // 其它真正的错误,退出
}
}
}
};std::bind(&HandlerConnection::Recver, _1) 的意思是:把成员函数(这里是静态的)Recver 绑定成一个 func_t 形状的对象,调用时往 _1 这个位置塞一个 Connection*。这样 conn->_recver(conn) 就等价于 HandlerConnection::Recver(conn)。这里 Recver 用静态方法,_1 占位符来自 std::placeholders。
AddConnection 会把新连接注册进 epoll 并加入映射表:
void TcpServer::AddConnection(Connection *conn)
{
struct epoll_event ev;
ev.events = conn->Events(); // 关注的事件(EPOLLIN | EPOLLET)
ev.data.fd = conn->SockFd(); // 关键:把 fd 塞进 data.fd,epoll_wait 后靠它反查
epoll_ctl(_epfd, EPOLL_CTL_ADD, conn->SockFd(), &ev); // 加进监听
_conns[conn->SockFd()] = conn; // 维护映射:fd → 连接对象
}读事件 Recver:把字节流交给上层
数据连接的 EPOLLIN 触发时,调用 Recver 读数据。这里要记住 Reactor 的一个设计哲学(课件里的原话):读的时候我们关心数据是什么格式、协议是什么样吗?不关心!我们只负责把这轮能读的读干净,把读到的字节流交给上层,由上层去分析处理。 这叫"职责分离"——底层只搬字节,上层才懂协议。
// Recver —— 处理"普通连接可读":把读到的字节追加进 inbuffer,再交给上层 HandlerRequest
static void Recver(Connection *conn)
{
while (1) // ET 模式:一次把能读的都读光
{
// 注意:是非阻塞读。读到 EAGAIN 说明暂时没数据了,退出循环。
char buffer[1024];
ssize_t n = recv(conn->SockFd(), buffer, sizeof(buffer) - 1, 0);
if (n > 0) // 读到 n 个字节
{
buffer[n] = 0; // 转成字符串好追加(安全性看具体协议)
conn->AddInBuffer(buffer); // 追加进接收缓冲区,而不是覆盖!
}
else if (n == 0)
{
// 对端正常关闭(收到 FIN):这是一个"断开事件",走异常/清理回调
conn->_excepter(conn);
return;
}
else // 出错:根据 errno 区分
{
if (errno == EAGAIN) break; // 暂时没数据了,正常退出
else if (errno == EINTR) continue; // 被信号打断,重试
else { conn->_excepter(conn); return; } // 真正错误,清理连接
}
}
// 读到这儿,说明这轮读完了。把整个 inbuffer 交给上层去解析成"一条条完整报文"
HandlerRequest(conn);
}n == 0 判据是 TCP 的一个重要细节:recv 返回 0,意味着对端发来了 FIN(优雅关闭),而不是"没数据"。没数据时非阻塞读返回的是 -1 且 errno == EAGAIN。所以判断连接断开,看 recv 返回 0 就行——这是每个网络程序员都必须刻进 DNA 的知识点,后面"连接断开"一节还会细讲。
写事件 Sender:outbuffer 与 EPOLLOUT
Reactor 里最容易写崩的其实是写事件。这里有个初学者几乎必踩的坑:一上来就"收到请求马上 send 回包"。可问题是,send 不一定能把你要发的数据一次发完——发送缓冲区可能满了,尤其对方收得慢的时候。
Reactor 的写事件解法是"先攒缓冲,再靠事件驱动地慢慢发":
// Sender —— 处理"可写事件":尝试把 outbuffer 里的数据尽量发出去
static void Sender(Connection *conn)
{
std::string &outbuffer = conn->OutBuffer(); // 拿到发送缓冲区
while (1)
{
// 非阻塞 send:把 outbuffer 尽量发出去
ssize_t n = send(conn->SockFd(), outbuffer.c_str(), outbuffer.size(), 0);
if (n >= 0)
{
outbuffer.erase(0, (size_t)n); // 已经交给内核的,从我们缓冲区里删掉
if (outbuffer.empty()) break; // 全发完了,退出
}
else
{
if (errno == EAGAIN) break; // 内核发送缓冲满了,这次发不完,退出等下次
else if (errno == EINTR) continue; // 被信号打断,重试
else { conn->_excepter(conn); return; } // 真正错误,清理
}
}
// 走到这说明本轮"能发的都发了",但可能还没发完
if (!conn->OutBufferEmpty())
{
// 还有没发完的 → 开启对这个 fd EPOLLOUT 的关注,等下一次可写事件再来发
conn->_R->EnableReadWrite(conn->SockFd(), true, true);
}
else
{
// 发完了 → 关闭 EPOLLOUT 关注,否则会一直"可写"狂触发,空转烧 CPU
conn->_R->EnableReadWrite(conn->SockFd(), true, false);
}
}这段话点出了 EPOLLOUT 的"脾气":EPOLLOUT 几乎总是就绪的(只要内核发送缓冲没满,socket 就"可写")。如果你注册了它却常发不出数据,它就会高频触发,让事件循环空转。所以工程铁律是:平时不要注册 EPOLLOUT,只在"确实有数据没发完"(outbuffer 非空)时才临时注册,发完立刻注销。 这一条能救很多人的命,后面"坑"一节还会专门讲。
异常处理 Excepter:连接断开与清理
Excepter 负责"生命周期终结"这件事:连接出错、对端断开、发数据时对方关闭……所有这些"这条连接没用了"的时刻,都统一在这里收尾。清理是一个绕不开的顺序问题,顺序错了会直接崩溃:
// Excepter —— 异常/断开收尾:把连接从各个地方干净地移除
static void Excepter(Connection *conn)
{
// 清理三部曲的顺序很重要:
// 1. 先从 epoll 的关注里移除这个 fd,防止它继续触发事件导致悬空回调
// 2. 从服务器的映射表里移除"fd → conn"
// 3. 关闭 socket,释放内存
conn->_R->RemoveConnection(conn->SockFd()); // 内部会做 epoll_ctl(del) 和 map.erase
conn->Close(); // close(fd)
delete conn; // 释放这个 Connection 对象
// 注意:delete 之后绝不能再碰 conn 了,所有清理动作都要在 delete 之前完成
}一定要记住这个顺序:先摘(从 epoll 和 map 里移除),再关(close fd),最后删(delete 对象)。如果先 delete 再去做别的,后面的代码只要碰到已释放的 conn,就是悬垂指针(dangling pointer),直接崩溃。
等一下——这里有隐患:Sender 里调用 conn->_R->EnableReadWrite(...),Recver 里调用 conn->_excepter(conn),这些都是在事件循环的 Loop() 里,通过 conn 指针回调进来的。如果回调内部把 conn delete 了,回到 Loop() 里的 for 循环还可能继续用 conn——这就是经典的"回调过程中对象被销毁"的悬挂指针坑。更稳妥的做法是:清理时不立即 delete,而是标记为待清理(deferred),等这一轮分发循环全部结束、不再有人引用 conn 之后,再统一释放。这是所有严肃 Reactor(Muduo、libevent 等)都在用的手法,我在这里先埋个伏笔,后面"回调里别阻塞 / 生命周期"一节会展开。
HandlerRequest:上层业务处理 + 半包粘包
HandlerRequest 是"上层"的代表——它不管数据怎么读进来的,只负责把 inbuffer 里的字节解析成"一条条完整报文"来处理:
// HandlerRequest —— 上层的协议解析 + 业务处理
static void HandlerRequest(Connection *conn)
{
std::string &inbuffer = conn->InBuffer();
std::string message; // 一条"符合协议边界的完整报文"
// while (Decode(inbuffer, &message)):从 inbuffer 里抠一条完整报文
// 抠不出来(半包没凑齐)就返回,剩下的留在 inbuffer 里等下一次读事件
while (Decoder::Decode(inbuffer, &message))
{
// message 一定是完整的报文了
// 剩下的:反序列化 → 业务计算 → 序列化 → 加边界封装 → 追加进 outbuffer
// 这里 elided 业务细节,落点在"把应答追加到 outbuffer,交给 Sender 发"
conn->AddOutBuffer(build_response(message));
}
// 攒了一批要发的应答,就主动尝试"直接发一次"(而不是干等 EPOLLOUT)
// —— 因为通常这里 outbuffer 里有货,直接调 Sender 走一轮"尽量发"更高效
if (!conn->OutBufferEmpty())
{
conn->_sender(conn); // 注意:这只是"尽量发",不代表能全部发完!
// 发不完的部分,Sender 会自己注册 EPOLLOUT,等可写事件继续发
}
}这里闪现了两个网络必考概念,我提前预告一下(后面有专门小节展开):半包(一个报文还没凑齐,"半个")和粘包(两个报文粘在一起)。Decode 就是干这个的——从字节流里按协议边界切出完整报文,切不出完整报文就留在缓冲区继续等。
我上面这些类是"教学拆解版",为了让你看请每个角色的职责。现在把它们拼成一个真正能编译运行的完整程序。我写成单文件,方便你 g++ 一键跑起来:
// reactor_echo.cpp —— 单线程 Reactor 回声服务器(完整可编译运行)
// 用法:g++ -std=c++17 -O2 -pthread reactor_echo.cpp -o reactor_echo && ./reactor_echo
// 客户端:nc 127.0.0.1 8888 (或 telnet),随便发什么,原样回给你
#include <iostream>
#include <string>
#include <functional>
#include <unordered_map>
#include <cstring>
#include <cerrno>
#include <unistd.h>
#include <fcntl.h>
#include <sys/socket.h>
#include <sys/epoll.h>
#include <netinet/in.h>
#include <arpa/inet.h>
// ============ 工具函数 ============
void SetNonBlock(int fd) // 把 fd 设为非阻塞
{
int fl = fcntl(fd, F_GETFL); // 取当前标志
fcntl(fd, F_SETFL, fl | O_NONBLOCK); // 加上非阻塞位
}
// 一个简单的消息边界协议:每条消息前面用 4 字节大端表示长度(length-prefix)
// 这能彻底解决"半包/粘包",因为有了长度就能精确切出每个完整报文
struct ByteStream
{
std::string data; // 累积的未解析字节
// 尝试从 data 里取出一条完整报文(前面 4 字节是长度)
// 成功返回 true 并把报文装进 out;长度不够(半包)返回 false,data 里的不动
bool Decode(std::string &out)
{
if (data.size() < 4) return false; // 连长度头都没凑齐,半包,等下次
uint32_t len = 0;
for (int i = 0; i < 4; i++) // 组装大端长度
len = (len << 8) | (unsigned char)data[i];
if (data.size() < 4u + len) return false; // 长度够了,但数据还没到齐,仍是半包
out = data.substr(4, len); // 切出这条完整报文
data.erase(0, 4u + len); // 已消费的部分从缓冲区删除
return true;
}
static std::string Encode(const std::string &msg) // 给报文加上 4 字节长度头
{
uint32_t len = (uint32_t)msg.size();
std::string head(4, '0');
for (int i = 0; i < 4; i++) // 大端写长度
head[i] = (char)((len >> (8 * (3 - i))) & 0xFF);
return head + msg;
}
};
// ============ 前置声明 ============
class TcpServer;
class Connection;
using func_t = std::function<void(Connection *)>; // 回调类型:吃一个 Connection*
// ============ Connection:一个连接 ============
class Connection
{
public:
Connection(int fd, uint32_t ev, TcpServer *t)
: fd_(fd), events_(ev), server_(t) {}
void Reg(func_t r, func_t w, func_t e) { on_read_ = r; on_write_ = w; on_err_ = e; }
void AddIn(const std::string &d) { recv_buf_.data += d; } // 追加进接收侧
void AddOut(const std::string &d) { send_buf_ += d; } // 累积到发送缓冲
int fd() const { return fd_; }
uint32_t events() const { return events_; }
void SetEvents(uint32_t e) { events_ = e; }
std::string &OutBuf() { return send_buf_; }
ByteStream &InBuf() { return recv_buf_; }
void Close() { ::close(fd_); }
func_t on_read_; // 可读回调:监听连接→Accepter,数据连接→Recver
func_t on_write_; // 可写回调:Sender
func_t on_err_; // 异常回调:Excepter
TcpServer *server_; // 回指服务器,方便调用它的注册/移除接口
private:
int fd_; // socket fd
uint32_t events_; // 关注的事件位
ByteStream recv_buf_; // 接收字节流缓冲
std::string send_buf_;// 发送缓冲
};
// ============ 事件处理器(集中的"处理逻辑") ============
struct Handler
{
// Accepter:监听 socket 可读 → accept 新连接
static void Accepter(Connection *lconn)
{
while (1)
{
struct sockaddr_in peer; socklen_t len = sizeof(peer);
int cfd = ::accept(lconn->fd(), (struct sockaddr *)&peer, &len);
if (cfd > 0)
{
SetNonBlock(cfd);
char ip[INET_ADDRSTRLEN];
inet_ntop(AF_INET, &(peer.sin_addr), ip, sizeof(ip));
std::cout << "[+] new conn " << cfd << " from " << ip
<< ":" << ntohs(peer.sin_port) << std::endl;
// 造一个数据连接:只关心 可读(ET),回调都指向 Handler 的静态方法
Connection *c = new Connection(cfd, EPOLLIN | EPOLLET, lconn->server_);
c->Reg(Recver, Sender, Excepter);
lconn->server_->AddConnection(c); // 注册进 epoll 和映射表
}
else
{
if (errno == EAGAIN) break; // 新连接收完了
else if (errno == EINTR) continue;
else break;
}
}
std::cout << "[accepter] done for listen " << lconn->fd() << std::endl;
}
// Recver:数据连接可读 → 读入缓冲,再尝试解析成一条条报文
static void Recver(Connection *c)
{
char buf[4096];
while (1)
{
ssize_t n = ::recv(c->fd(), buf, sizeof(buf), 0);
if (n > 0) { c->AddIn(std::string(buf, (size_t)n)); } // 追加,不是覆盖
else if (n == 0) { DeferClose(c); return; } // 对端 FIN:断开
else {
if (errno == EAGAIN) break; // 读完一轮
else if (errno == EINTR) continue;
else { DeferClose(c); return; }
}
}
// 收完一轮,解析 + 业务 + 攒应答
std::string msg;
while (c->InBuf().Decode(msg)) // 一条条完整报文
{
std::cout << "[recv " << c->fd() << "] " << msg << std::endl;
std::string reply = ByteStream::Encode(msg); // 回声:原样回
c->AddOut(reply); // 先攒进发送缓冲
}
if (!c->OutBuf().empty()) c->on_write_(c); // 有货就立刻"尽量发"一轮
}
// Sender:尝试把发送缓冲发出去;发不完注册 EPOLLOUT
static void Sender(Connection *c)
{
std::string &ob = c->OutBuf();
while (1)
{
ssize_t n = ::send(c->fd(), ob.data(), ob.size(), 0);
if (n >= 0) {
ob.erase(0, (size_t)n);
if (ob.empty()) break; // 全发完
} else {
if (errno == EAGAIN) break; // 内核缓冲满,等下次 EPOLLOUT
else if (errno == EINTR) continue;
else { DeferClose(c); return; }
}
}
// 发完与否,决定是否还关心 EPOLLOUT
c->server_->EnableReadWrite(c->fd(), true, !ob.empty());
}
// Excepter:统一收尾(这里用"延迟关闭"避免悬垂)
static void Excepter(Connection *c) { DeferClose(c); }
// 延迟关闭:只打标记,真正的 delete 交给事件循环在安全时刻做
static void DeferClose(Connection *c) { c->server_->DeferClose(c); }
};
// ============ TcpServer:Reactor 主控 ============
class TcpServer
{
public:
TcpServer(uint16_t port) : port_(port) {}
void Start()
{
// 1. 建监听 socket
lfd_ = ::socket(AF_INET, SOCK_STREAM, 0);
int opt = 1;
setsockopt(lfd_, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt));
struct sockaddr_in a; std::memset(&a, 0, sizeof(a));
a.sin_family = AF_INET; a.sin_port = htons(port_); a.sin_addr.s_addr = INADDR_ANY;
::bind(lfd_, (struct sockaddr *)&a, sizeof(a));
::listen(lfd_, 512);
SetNonBlock(lfd_);
// 2. 建 epoll
epfd_ = ::epoll_create1(0);
// 3. 用工厂造"监听连接",回调指向 Accepter,只关心可读(ET)
Connection *lc = new Connection(lfd_, EPOLLIN | EPOLLET, this);
lc->Reg(Handler::Accepter, nullptr, nullptr);
AddConnection(lc);
std::cout << "[reactor] listening on " << port_ << std::endl;
Loop();
}
void AddConnection(Connection *c) // 把一个连接纳入 epoll + 映射
{
struct epoll_event ev{};
ev.events = c->events();
ev.data.fd = c->fd(); // 用 data.fd 携带回查线索
::epoll_ctl(epfd_, EPOLL_CTL_ADD, c->fd(), &ev);
conns_[c->fd()] = c;
}
void RemoveConnection(int fd) // 从 epoll 和映射里摘掉(不 delete)
{
::epoll_ctl(epfd_, EPOLL_CTL_DEL, fd, nullptr);
auto it = conns_.find(fd);
if (it != conns_.end()) { conns_.erase(it); }
}
void EnableReadWrite(int fd, bool rd, bool wr) // 动态调整某 fd 关注的事件位
{
auto it = conns_.find(fd);
if (it == conns_.end()) return;
Connection *c = it->second;
uint32_t ev = EPOLLET; // 永远保持边缘触发
if (rd) ev |= EPOLLIN;
if (wr) ev |= EPOLLOUT;
c->SetEvents(ev);
struct epoll_event ee{};
ee.events = ev; ee.data.fd = fd;
::epoll_ctl(epfd_, EPOLL_CTL_MOD, fd, &ee); // MOD 修改,不是 ADD
}
void DeferClose(Connection *c) // 记进待清理列表,稍后统一释放
{
dead_.push_back(c);
}
private:
void Loop() // 事件循环本体
{
const int MAXE = 64;
struct epoll_event evs[MAXE];
while (1)
{
int n = ::epoll_wait(epfd_, evs, MAXE, -1);
if (n < 0) { if (errno == EINTR) continue; break; }
for (int i = 0; i < n; i++)
{
int fd = evs[i].data.fd;
uint32_t ev = evs[i].events;
auto it = conns_.find(fd);
if (it == conns_.end()) continue; // 已被清理,跳过
Connection *c = it->second;
if (ev & EPOLLIN) { if (c->on_read_) c->on_read_(c); }
if (ev & EPOLLOUT) { if (c->on_write_) c->on_write_(c); }
if (ev & (EPOLLERR | EPOLLHUP))
{
if (c->on_err_) c->on_err_(c);
continue;
}
}
// ---- 本轮分发结束,此刻安全,统一释放"待关闭"连接 ----
for (Connection *c : dead_)
{
std::cout << "[-] close conn " << c->fd() << std::endl;
RemoveConnection(c->fd()); // 1 先从 epoll/map 摘掉
c->Close(); // 2 关闭 fd
delete c; // 3 释放内存
}
dead_.clear();
}
}
uint16_t port_;
int lfd_ = -1, epfd_ = -1;
std::unordered_map<int, Connection *> conns_; // fd → 连接对象
std::vector<Connection *> dead_; // 本轮回调节点产生的"待清理连接"
};
int main()
{
TcpServer s(8888); // 监听 8888 端口
s.Start(); // 启动,进入永不返回的事件循环
return 0;
}编译运行:
# 编译(用 C++17 及以上)
g++ -std=c++17 -O2 -pthread reactor_echo.cpp -o reactor_echo
# 启动服务器
./reactor_echo
# 另开一个终端,多开几个客户端连上来试试
nc 127.0.0.1 8888这个程序就是一个完整的、可编译可运行的单线程 Reactor 回声服务器。你可以在里面跑 stress 或同时开几十个 nc,观察它如何"一个线程管住所有连接"。这段代码我在"延迟关闭"(DeferClose)上做了工程化处理,正因为 Recver/Sender 里会通过回调调用 DeferClose,如果在回调里就 delete,回到 Loop 的 for 循环继续用 c 就是悬垂指针。所以一切销毁都推迟到本轮分发循环结束——这是你将来在真实服务器源码里会反复看到的模式。
事件注册与分发:Dispatcher 怎么知道该叫谁
现在我们把"事件注册"和"事件分发"这两个词嚼透,这是理解 Reactor 运转的钥匙。
事件注册(Event Registration):就是告诉 epoll"我关注这个 fd 的哪些事件"。用的 API 是 epoll_ctl,三种操作:
EPOLL_CTL_ADD:把一个 fd 和它关注的事件加进 epoll 的事件表;EPOLL_CTL_MOD:修改某个 fd 关注的事件位(比如从只关注可读 → 同时关注可写);EPOLL_CTL_DEL:把某个 fd 从事件表里移除。
对应到我们代码里就是 AddConnection(ADD)、EnableReadWrite(MOD)、RemoveConnection(DEL)。区分 ADD 和 MOD是新手最容易翻车的地方:对同一个 fd 第二次 epoll_ctl 用 ADD 会返回 EEXIST 错误;而用 MOD 去操作一个还没注册的 fd 会返回 ENOENT。所以一定要让"注册 / 修改"各安其位。
事件分发(Event Dispatch):就是 epoll_wait 返回一批就绪事件后,Reactor 把每个就绪 fd"分发"到它的回调。分发者(Loop)不关心回调具体干什么,它只拍板一个映射:哪个 fd 就绪了 → 调哪一个方法。这个映射的两端是:
- 输入端:
epoll_event.data.fd(epoll 就绪时带给我们的 fd); - 输出端:
conns_[fd]->on_read_ / on_write_ / on_err_(映射表反查出的回调)。
epoll_event 的 data 联合体是个很有用的设计:除了 data.fd,还可以用 data.ptr 直接指向某个自定义结构(比如你的 Connection 对象),省去一次 unordered_map 查找。但我们用 data.fd + map 的方式更直观,也好讲清楚"分发"这件事。
思考题 1:为什么说"监听 socket 的
EPOLLIN和普通连接的EPOLLIN是同一个事件位,却必须走不同的处理"?如果我只注册一个通用的EPOLLIN回调给所有 fd,会发生什么?详解:
epoll的事件位只描述"就绪状态",不描述"是谁"(监听还是数据)。同一个EPOLLIN位,对监听 socket 意味着"有连接请求可 accept",对数据连接意味着"有数据可 recv"。一个是"建新连接",一个是"收数据处理",语义完全不同。分发者不靠事件位区分这两类,靠的是"这个 fd 注册时塞进来的回调"——监听连接塞的是Accepter,数据连接塞的是Recver。所以"分发"的实质是"按 fd 找到当初注册的回调再调用",这正是conns_映射表存在的意义。如果你只注册一个通用
EPOLLIN回调给所有 fd,你就得在回调里分辨"这个 fd 是监听还是普通",要么靠判断 fd 是否等于lfd_,要么靠其它标志位。这会污染业务代码,也让"新增一种 fd 类型"变得很僵硬。Reactor 的价值恰恰就在于:用不同的回调封装不同的语义,分发者与回调完全解耦,加新事件类型不用改分发逻辑。
Accepter:把"新连接请求"变成"一条数据连接"
Accepter 是个很有意思的角色——它本身也是一个 Handler(处理"连接到来"这个事件),但它做的是"生产新连接"的活。它把监听到的"连接请求"变成"一条可被 Reactor 照管的数据连接",再交给 AddConnection 注册进事件循环。
这里有个 ET 模式必须掌握的细节(源材料里特别点了这个):ET 模式下事件只通知一次,你怎么保证只有一条连接到达时不漏接收? 答案不是去"猜"有几条,而是用 while 循环把 accept 收干为止。
ET(Edge Triggered,边缘触发)和 LT(Level Triggered,水平触发)的区别,是我们学 epoll 时反复强调过的:
- LT:只要缓冲区里还有数据(还有成功完成的条件),就会反复通知,哪怕你不处理它每次都报。优点是不容易漏,缺点是"明明读光了它还再叫一次",容易重复处理。
- ET:只在状态发生跳变(从无到有)的那一次通知,之后不再通知,直到你把它读光(读到 EAGAIN)。必须一次性处理干净,否则剩下的就永远没机会被通知。
正因为 ET "一次性",accept 和 recv 都必须用 while 读到 EAGAIN。这就是我们 Accepter 里 while(1) + if (errno==EAGAIN) break 的由来。在 ET 里,"读光为止"是纪律,不是可选项。
读事件 Recver:数据搬运与半包处理
读的职责再强调一次:只负责把"这一轮能读的字节"搬进 inbuffer,绝不越界去解析协议。 协议解析是上层 HandlerRequest 的事。这就是"职责分离"——底层不要关心业务,上层不要关心底层怎么把数据弄来的。
可是这里马上撞上一个 TCP 的硬事实:TCP 是字节流,没有消息边界。 你 recv 拿到的 1024 字节,可能只是"半个请求",也可能是"三个请求+半个",还可能是"一个半请求"。这就引出了两个网络编程的千古难题:
- 半包(partial packet):一条完整报文还没收齐,只收到了它的一部分。
- 粘包(sticky packet):一条
recv里塞进了多条完整报文的拼接。
解决它们的手段叫报文边界协议 / 粘包粘包处理,常见三种:
- 固定长度:每条报文定长 N 字节,不够就等,凑齐 N 就处理——最简单,但浪费(短报文也会占 N)。
- 长度前缀(length-prefix):每条报文开头先发 4 字节"正文长度",接收端先读长度,再按长度凑齐正文。这是最常用、最通用的方案(我们的
ByteStream就是它)。 - 分隔符 / 行协议:用
\r\n之类的字符分隔报文(HTTP 头、Redis 协议用的是它的变种)——直观,但要防止正文里恰好出现分隔符。
看到了吗,Recver 在处理粘包时,用的是 length-prefix + 缓冲累积的做法——这就是 Reactor 里"缓冲(buffer)"存在的意义。接收缓冲 inbuffer 就是为处理半包/粘包而存在的蓄水池:收不够就攒着,够一条就切一条,剩下的继续留在池里等下一轮。
思考题 2:如果
Recver每收到一段字节就立刻处理,而不攒进inbuffer,会出什么问题?详解:会出两个问题。第一是粘包:一次
recv可能拿到多条报文的拼接,如果你把整段当"一条"处理,就把多条请求混成一坨,业务逻辑直接错乱。第二是半包:一次recv也可能只拿到一条报文的"前一半",长度都没凑齐,你根本无从解析;如果你硬解析,会读到不完整的脏数据。把数据先"攒进缓冲"(inbuffer),再按边界协议(length-prefix)"等一条凑齐切一条",才能同时对付粘包和半包。缓冲的本质是"把 TCP 无序的字节流重排成有边界的消息流"。
写事件 Sender:发送缓冲与 EPOLLOUT 的取舍
写事件我前面已经打了预防针,这里把"为什么"讲透。
send 不是一个"你说发,它就全发"的接口。它只保证"尽力发那么多字节出去",返回值可能小于你要发的长度。什么时候会发不完?内核发送缓冲区满了——尤其对端应用层读得慢、或者 TCP 窗口为 0 的时候。数据就这么滞留在我们的 outbuffer 里。
Reactor 的正规处理是"缓冲 + 事件驱动的慢发":
- 业务方要发数据时,先
AddOutBuffer攒进outbuffer,而不是立刻send(能立刻发是"尽量",发不完才是常态)。 - 尝试调
Sender发一轮:能发多少发多少,发不完的留在outbuffer。 - 若发不完,
EnableReadWrite(fd, true, true)注册EPOLLOUT。内核发送缓冲一旦有空位,epoll 就报"可写",触发Sender再来发一轮。 - 发空了,
EnableReadWrite(fd, true, false)注销EPOLLOUT,防止无意义地狂触发。
这里最反直觉、也最值钱的经验是:EPOLLOUT 通常处于"几乎总是就绪"的状态。Linux 的 socket 只要内核发送缓冲没满,它就"可写"。如果你给一个空闲连接注册了 EPOLLOUT,事件循环会发现它"老是可写",于是每个循环都会回调它一次——哪怕你根本没有数据要发。这就是 "busy loop / 空转烧 CPU"。所以规矩是:
平时不要监听 EPOLLOUT;只有确实有数据滞留在 outbuffer 发不完时,才临时注册 EPOLLOUT,发空后立刻注销。
这一条比你想的重要得多,直接决定服务器是"接近零负载地等多方数据"还是"狂转空烧"。
异常处理 Excepter:连接断开与生命周期管理
Excepter 要处理的情况很多:对端 recv 返回 0(FIN,优雅关闭)、RST 造成 EPOLLERR 或 ECONNRESET、send 时收到 EPIPE(对端已关闭我们还写)……它们的共性:这条连接废了,要清场。
连接断开的两种典型信号,务必分清:
- 半关闭(FIN,优雅关闭):对端
close或shutdown,recv返回 0。表示"对端不再发送",但你还能收到它已发的数据、你可能还能回。 - 异常复位(RST,强制关闭):对端崩溃、或往一条已关闭的连接上写,内核发 RST,之后读会
recv返回 -1 且errno==ECONNRESET,写会触发SIGPIPE或EPIPE。
另外一个隐蔽坑:SIGPIPE。默认情况下,进程往一条对端已关闭的 socket 上 write/send,内核会给进程发 SIGPIPE 信号,默认行为是终止进程——你的服务器会"莫名其妙"地整进程挂掉。所以生产环境要么在启动时 signal(SIGPIPE, SIG_IGN)(忽略它,让 send 正常返回 EPIPE 错误自己处理),要么给 send 加 MSG_NOSIGNAL 标志。这是每个网络服务器都必须处理的第一件事,忘了它,掉一个客户端你的服务器就崩溃了。
清场的顺序前面强调过:先摘 → 再关 → 后删。而且更稳妥的是延迟销毁:不在回调内部立刻 delete,而是记入"待清理队列",等本轮分发结束再统一释放,从而彻底避开"回调过程中对象被销毁造成的悬垂指针"。
思考题 3:清理连接时为什么必须"先从 epoll 和映射表中移除,再 close,最后 delete"?顺序能不能换?
详解:三条理由决定了顺序。第一,先摘:如果不先把 fd 从 epoll 移除,而先
delete/close,那么下次epoll_wait可能带回这个"已关闭/已释放"的 fd,我们去conns_反查却找不到(或找到悬垂指针),直接崩溃;所以要先epoll_ctl(DEL)+map.erase,切断它被再次分发的可能。第二,再关:close只在 fd 还存在时合法,若已 delete 了对象、fd 也离开了我们手里,就没法安全 close 了,而且 fd 是有限资源(默认 1024/可调),拖越久越可能耗尽,所以要及时关。第三,后删:delete会真正释放Connection占用的内存,一旦 delete,之后再碰conn就是悬垂指针。所以所有"要用到 conn 的收尾动作"必须在 delete 之前完成。这也是为什么落库/打日志/递减计数这些动作都得放 delete 前做完。
回调解耦:服务器(Reactor)与业务解绑
这是 Reactor 模式最优雅的地方,值得单独撑起一节。
看我们 Loop 的分发代码,它从头到尾一次都没碰过"业务"——它不知道收到的数据是 HTTP 还是 DNS,不知道要不要计算,不知道应答格式。它只知道三件事:
- 这个 fd 就绪了;
- 我要不要调它的
on_read_(可读时); - 调完拉倒。
具体"往 inbuffer 外面做了什么",全在 Recver → HandlerRequest 这条业务链里,与事件循环无半毛钱关系。这,就是回调解耦(Callback Decoupling)。
好处立刻显现:
- 职责分离:IO(event demultiplex) 与 业务(business logic) 彻底分开,改协议不改事件循环,改事件循环不动业务。
- 可替换性:你想把这个 Reactor 从一个协议换到另一个协议,只需换
Connection注册进去的回调(换一组 Handler),Reactor 本身一行不用改。 - 可复用:同一份 Reactor 骨架,"填充"不同的回调,就变成 HTTP 服务器、RPC 服务器、聊天服务器……这就是为什么 Reactor 是绝大多数网络框架(libevent、libuv、Muduo、Nginx、Netty、Redis)的内核。
一句话:Reactor 是"演出的舞台和班车",Handler 是"上台的演员"。舞台不关心演员演什么,演员也不关心舞台怎么搭——换戏不换台,换台不换戏。
思考题 4:如果用一句话回答——"Reactor 模式到底值钱在哪?"你会怎么说?
详解:可以概括为三点,我更倾向于用这层递进:(1) 用单一(或极少数)事件循环承载海量并发连接,避免"一连接一线程"的线程爆炸;(2) 把"事件检测"(epoll 等)与"事件处理"(Handler 回调)解耦,让同样的 IO 骨架可以适配千变万化的业务;(3) 事件循环天然单线程化地消灭了大部分并发访问冲突(同一 fd 同一时刻只有一个回调在处理它),极大简化了加锁负担。这三者结合,构成了几乎一切高性能网络服务器可复用、可扩展、可维护的核心骨架。
单线程 vs 多线程 Reactor:三种标准变体
单线程 Reactor 就够了吗?不够——如果业务处理很耗时(比如要做一次复杂的计算、要访问数据库),那么回调会阻塞住事件循环,一个慢业务卡住,所有连接都被冻住。这就违背了 Reactor"高效并发"的初衷。
于是出现了多线程 Reactor 的三种标准变体,我把它们画清楚:
变体一:单 Reactor 单线程
一个线程
┌──────────────────────┐
│ Reactor (EventLoop) │
│ 监听 + accept + 读写 │ ← 所有连接的所有事件都在这一条线上串行处理
│ + 业务处理 │
└──────────────────────┘
优点:实现最简单、无任何锁。缺点:业务与 IO 全在一个线程,一个重业务就把整台服务器卡死。适合"业务极快、连接很多"的场景(比如回显、静态资源,Redis 的核心模型就很接近这一种)。
变体二:单 Reactor 多线程(业务线程池)
┌──────────────┐
│ Reactor │ ← 只负责:监听、accept、读写、分发
│ (单线程) │
└──────┬───────┘
│ 分发"耗时业务"给线程池
▼
┌─────────────────────┐
│ 业务线程池 Worker1 2 3...│ ← 真正吃 CPU 的业务放到线程池,不阻塞 Reactor
└─────────────────────┘
IO 仍在一个线程;耗时的业务计算交给线程池(Thread Pool)。Reactor 把请求扔进线程池后立刻回去监听下一个事件。优点:慢业务不再卡住 IO;缺点:Reactor 单线程承担所有 IO,吞吐瓶颈在这条线程;且线程池和 Reactor 之间共享数据要小心加锁。
变体三:多 Reactor 多线程(主从 Reactor,Main/Sub Reactor)
┌────────────┐ 一组 Sub Reactor(每个跑在一个线程里)
│ MainReactor │ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ 只做 │ 分 │ SubReactor│ │ SubReactor│ │ SubReactor│
│ accept │ ───▶ │ (EventLoop)│ │ ... │ │ ... │
│ (监听新连接) │ │ 处理IO+业务 │ └──────────┘ └──────────┘
└────────────┘ └──────────┘
这是 Netty、Nginx(worker 进程)、多种生产级框架的标配。Main Reactor 的线程只负责 accept 新连接;它 accept 出来的数据连接,按负载分摊给多个 Sub Reactor,每个 Sub Reactor 跑在一个独立线程的事件循环里,自己全权管理自己那批连接的完整生命周期。每个 Sub Reactor 管理自己的 fd 集合,彼此不串数据,因此各线程之间几乎没有需要共享的临界区——这就是它既高效又低锁冲突的秘密。
思考题 5:多 Reactor(主从)为什么能比"单 Reactor 多线程"更高效?它的前提假设是什么?
详解:核心在于**"减少共享,进而减少锁"。单 Reactor 多线程方案里,所有 IO 挤在一个线程,业务在另一堆线程,两条线之间必然要为一个连接的数据 "交接(handoff)" 而加锁,这个锁是全局的、高频的,成为吞吐的卡点。而多 Reactor 方案把"连接"按批次彻底分到各个 Sub Reactor,每条连接从出生到淘汰都只属于一个 Sub Reactor 线程,该线程自己处理自己的 IO 和业务,天然没有跨线程竞争——就像仓库分成若干独立小仓,每个仓自己管理自己的货,互不抢钥匙。它的前提是连接之间相对独立**(大多数服务端场景天然如此),只要业务不要求"一个请求访问到多个不同线程所持有的状态",这个模型就非常干净。如果业务必须跨连接共享全局状态(比如全局排行榜),那依然要引入锁或其它共享机制,就回到"有共享就要有同步"的普适规律了。
OTOL(One Thread One Loop):一线程一事件循环
聊完三种变体,引入一个更彻底的组织原则:OTOL(One Thread One Loop,一线程一事件循环)。
它说的是:每个线程跑一个独立的事件循环。这是一种"线程 = 事件循环 = 一组连接的自治单元"的组合。上面的多 Reactor 多线程,本质就是一个 OTOL 的实例化——每个 Sub Reactor 就是"一个线程一个 loop"。
OTOL 的核心理念,源材料里概括得很精辟:
- 只要把单 Reactor 写完,扩展多进程/多线程很容易——因为 Reactor 本身是"一份能独立自洽处理完整连接生命周期的代码",把它复制到多个执行流里,让它各自管一批连接即可。
- 每个进程/线程管理的连接和 fd 彼此不要重复——各行其道,互不相干,这就是"无共享"的前提。
- 每个 Reactor 自己处理自己的 sockfd 的完整生命周期,不涉及任何 IO 穿插与乱序——单线程内天然有序,没有并发数据竞争(同一 fd 不会同时被两个线程碰)。
那问题来了:多线程之间要"投递任务 / 互相唤醒",靠什么通信? 事件循环是"阻塞在 epoll_wait 上等事件"的,你没法直接"打断它、插一个任务进去"。你需要一个能被 epoll 关注到的通道来"从外部向事件循环注入事件"。有两条路:
- 管道(pipe):把一个管道读端注册进某线程的事件循环,别的线程往里
write,事件循环就会因"管道可读"被唤醒。 - eventfd:专门为此而生的事件通知文件描述符(下一节重点讲)。
源材料里还给出了一个用管道在多线程之间互传 fd 的设计:让 master 线程 accept 新的连接 fd,通过管道把整张 fd 表序列化传递给各个 worker,每个 worker 从管道里读到 fd 后直接加入自己的 Reactor。因为多线程是共享文件描述符表的,这个"传 fd"的玩法在 Linux 上是可行的(也是 UNIX 域 socket sendmsg 传描述符的经典技巧)。
用 eventfd 唤醒事件循环
为了跨线程/跨进程向事件循环"注入事件",Linux 给了我们一个比管道更轻、更精准的专用件:eventfd。
eventfd 是一个基于文件描述符的轻量级事件通知机制。它像一个小计数器文件,内核为它维护一个 64 位计数器,你 write 会累加计数器,read 会读取(并清零/递减)计数器。它最大的价值是:它能被 epoll 关注,于是"往 eventfd 写一笔"就变成"唤醒那个阻塞在 epoll_wait 上的事件循环"。
// eventfd_demo.c —— eventfd 跨进程通知最小示例(可编译运行:gcc eventfd_demo.c -o efd)
#include <stdio.h>
#include <stdlib.h>
#include <unistd.h>
#include <stdint.h>
#include <sys/eventfd.h>
int main(void)
{
// 创建 eventfd:初始计数 0,EFD_CLOEXEC 保证 exec 时自动关闭
int efd = eventfd(0, EFD_CLOEXEC);
if (efd == -1) { perror("eventfd"); return 1; }
pid_t pid = fork();
if (pid == -1) { perror("fork"); return 1; }
if (pid == 0) /* 子进程 */
{
uint64_t value;
// 阻塞读:等父进程写,读到事件值
if (read(efd, &value, sizeof(value)) < 0) perror("read");
printf("Child received event, counter was %llu\n",
(unsigned long long)value);
close(efd);
return 0;
}
else /* 父进程 */
{
uint64_t value = 1;
// 写 1 到计数器,相当于"发一次事件";子进程(或阻塞在 epoll_Wait 的事件循环)会被唤醒
if (write(efd, &value, sizeof(value)) < 0) perror("write");
printf("Parent sent event\n");
wait(NULL); /* 等子进程结束 */
close(efd);
return 0;
}
}抓重点:
- 低开销:内部就一个 64 位计数器,内核维护成本极低,读写基本是纯常数时间。
- 支持多路复用:能被 epoll / poll / select 关注,这正是它能"唤醒事件循环"的前提。
- 原子性:计数器的读写天然原子,适合高并发下"安全地攒通知"。
- 广播通知:可以多对多通知——多个写者往一个 eventfd 写,一个或多个读者被唤醒,用于"攒信号"很爽。
- 不能传内容:
eventfd只做"事件发生"的通知,不承载消息内容(管道可以承载数据流,eventfd 不能)。要传真正的数据/描述符,还是得走管道或 socket。 - 建议
EFD_NONBLOCK:非阻塞读,避免在一个实际没有新事件的空档卡死。
关于 eventfd 的两种工作模式,源材料给了非常直观的对照,这里我复述一遍这两张表:
模式一(默认,不设 EFD_SEMAPHORE):read 一次把计数器整个取走并清零。
模式二(设 EFD_SEMAPHORE):read 一次只让计数器 -1(信号量语义)。
模式一(默认): 模式二(EFD_SEMAPHORE):
写 1 → 计数器 1 写 1 → 计数器 1
写 2 → 计数器 3 写 2 → 计数器 3
read → 读出 3,计数器变 0 read → 读出 1,计数器变 2
read → 读出 0,计数器 0 read → 读出 1,计数器变 1
read → 读出 0,计数器 0 read → 读出 1,计数器变 0对我们"唤醒事件循环"的场景,用**模式一(默认)**就对了:因为"来了几个事件"我们往往不关心,只关心"要不要被叫醒";一次读出全部都节约了一次 hack。把 eventfd 注册进 epoll,别的线程向它 write 一笔,epoll_wait 就会因 EPOLLIN 返回,事件循环就被"唤醒"起来,正好完成跨线程信号注入。
思考题 6:eventfd 和管道(pipe)都能唤醒事件循环,为什么更推荐 eventfd?
详解:两点关键差异。第一,开销:管道是真正的内核 buffer(有读写缓冲、数据拷贝、两端的 wait queue),每次读写都要搬运数据;eventfd 内部只是一个 64 位计数器,写一笔就是原子累加,读一笔就是原子取走,没有数据拷贝、内核路径更短,在"高频唤醒"场景开销低得多。第二,语义:管道是"流"语义,得关心数据内容、可能读多读少、还要防阻塞;eventfd 是"事件计数"语义,天然服务于"通知一下"这件事,配合原子性在多生产者唤醒场景更省心。代价是 eventfd 不能携带内容——所以当你需要在线程间"投递任务参数 / 序列化 fd"时,还得用管道或 socket;eventfd 负责"喊你起床",管道/socket 负责"递东西给你"。
从 Reactor 出发的下一步
到这里,我们把 Reactor 反应堆模式的来龙去脉基本过完了:从"一连接一线程"的不可扩展,到事件驱动、事件循环、四角色的定义;从用 epoll 亲手实现事件槽到注册与分发;从 Accepter/Recver/Sender/Excepter 四个 Handler 到它们在读、写、异常三条战线上的分工;从回调解耦到单线程/多线程/多 Reactor 三种变体;最后落到 OTOL 和 eventfd 这两个支撑多进程多线程扩张的地基。
最后给你三个"以后写服务器一定用得上"的心法,当收尾的锚点:
- 读永远攒缓冲,写永远挂缓冲:
read一定要累积进inbuffer用协议切边界(否则半包粘包翻车);write一定要先攒进outbuffer,发不完靠EPOLLOUT续命,别在回调里阻塞硬发。 - EPOLLOUT 是"按需出现的贼",不是"常驻的老友":发不完才临时注册,发空立刻注销,否则忙等到烧 CPU。
- 回调里绝不要做耗时的事:业务重,就搬去线程池;一切
delete都等本轮分发结束再统一做。记住这句话——事件循环的命门是"别让任何一次回调把整条循环拖住"。
Reactor 是绝大多数网络框架的内心。掌握了它,你再看 libevent、libuv、Muduo、Netty 的源码,会发现全是熟面孔:不过是把我们现在手写的东西,做得更健壮、更通用、更高效。愿你有一天打开 Nginx 或者 Netty 的 EventLoop 源码时,能会心一笑:哦,这就是那个"反应堆"啊。
还没有评论 — 第一条由你来留。