欢迎回到 Aloe。笔者近期有兴趣来细致地梳理有关于 epoll 的原理、 应用的相关的知识点。笔者会结合 UKVEngine 的应用,尽可能贴合实际来为各位读者介绍。

前言:为什么要理解 epoll?

我们先建立一个最简单的 PING-PONG 服务端程序:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
if (listen(server_fd, 16) < 0) {
std::cerr << "监听失败!" << std::endl;
return -1;
}
std::cout << "服务端启动成功,正在监听 8888 端口..." << std::endl;

for (;;)
{
int addrlen = sizeof(address);

int client_fd = accept(server_fd, (struct sockaddr*)&address, (socklen_t*)&addrlen);
if (client_fd < 0) {
std::cerr << "接收连线失败!" << std::endl;
return -1;
}
std::cout << "有客户端连进来了!" << std::endl;

char buffer[1024] = {};
read(client_fd, buffer, 1024);
std::cout << "收到客户端消息: " << buffer << std::endl;

std::string response = "PONG\n";
write(client_fd, response.c_str(), response.length());

close(client_fd);
}
close(server_fd);

省略了创建 Socket 的步骤,编译并运行,我们使用 nc localhost 8888 命令, 随便敲一些东西进去:

1
2
3
4
5
❯ nc localhost 8888
text
PONG
text

果然出现了 PONG 的回应,但是似乎又有什么地方不合预期

  • 为什么第二次输入 text 后,或者说再敲一个回车,nc自动退出了?

  • 开两个 nc 的时候,第一个人不发消息,为什么第二个人会被卡住

思考: 这两个问题需要怎么解决?会带来哪些新的问题?有没有两全其美的办法?

单线程阻塞式服务端的困境

我们来分别分析这两个问题背后的原因:

  • 循环的结尾,我们有一行 close(client_fd);,服务端发起四次挥手关闭连线, 因此 nc 会自动退出。 如果我们不去 close,是否就能达到预期?

危险: 再次观察 accept 一行代码,思考 accept 默认的行为,此操作的危险在哪里?

很显然如果注释掉 close,行为上 nc 的确不会直接退出了, 但是每次 acceptclient_fd 都被覆写,这会造成资源泄漏的严重问题! 此时哪怕客户端按下 ^C 关闭 nc,客户端发起四次挥手关闭连线, 但是服务端已经永久泄漏了上一个连线,服务端状态将永远卡在 CLOSE_WAIT 状态,无法调用 close 发送 FIN 进入 LAST_ACK 状态, 并且当连线数接近 1000 时,会触发 Too many open files 错误, 服务彻底瘫痪,很显然这是错误的做法

补充:TCP 四次挥手步骤

这里以客户端主动断开为例。

  1. 客户端发送 FIN 报文,进入 FIN_WAIT1 状态,表明没有数据要发送了, 但是仍然可以接收并处理数据。

    • 对应着客户端的 close(client_fd) 调用。
  2. 服务端收到 FIN,发送 ACK 报文,进入 CLOSE_WAIT 状态, 表明接收到对方的关闭请求,但是可能还会有数据要发送。

    客户端收到 ACK,进入 FIN_WAIT2 状态。

    • 这由内核自动处理,不需要写代码处理这次挥手过程。
  3. 服务端发送 FIN 报文,进入 LAST_ACK 状态,表明所有数据发送完毕, 等待客户端的最后确认。

    • 对应着服务端的 close(client_fd) 调用。
  4. 客户端收到 FIN,发送 ACK 报文,进入 TIME_WAIT 状态, 等待 2MSL 后进入 CLOSED 状态,

    服务端收到 ACK 后直接进入 CLOSED 状态。

    • 这由内核自动处理,不需要写代码处理这次挥手过程。

accept 用于从全连接队列取出一个连线,read 用于拷贝数据到缓冲区, 然而这两个调用都是阻塞的。

这意味着代码运行到这一行时,如果队列里没有连线, 或者数据内部缓冲区没有数据可供拷贝,操作系统内核会挂起该线程, 变为等待或睡眠的状态,剥夺 CPU 使用权。只有当全连接队列再次出现了一个连线, 内核才会重新唤醒这个线程。

因此这就能解答第二个问题了:

  • 第一个 nc 使服务端阻塞在了 read 调用,

  • 第二个 nc 在服务端甚至没能accept 从全连接队列取出。

  • 在第一个 nc 按下回车,服务端返回 PONG 后,我们看到第二个 nc 瞬间 就输出了 PONG,本质上是内核唤醒了被阻塞的线程。程序进入下一轮循环后立刻完成了 accept - read - write 三个步骤。

进化尝试一:多线程模型

既然单线程会被 read 阻塞,那么第一直觉就是:

  • 给每个连线建立并分配一个独立的线程,

  • 子线程去面对阻塞的 read 调用

  • 主线程去面对 accept

听起来似乎是一个完美的解决方案,很容易写出这样的代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
for (;;)
{
int addrlen = sizeof(address);

int client_fd = accept(server_fd, (struct sockaddr*)&address, (socklen_t*)&addrlen);
if (client_fd < 0) {
std::cerr << "接收连线失败!" << std::endl;
return -1;
}
std::cout << "有客户端连进来了!" << std::endl;

std::thread([client_fd] {
char buffer[1024] = {0};
read(client_fd, buffer, 1024);
std::cout << "收到客户端消息: " << buffer << std::endl;

std::string response = "PONG\n";
write(client_fd, response.c_str(), response.length());
close(client_fd);
}).detach();
}

close(server_fd);

看起来一切都是那么美好,甚至看起来是一个“两全其美”的解决方案。 事实上早期的互联网时代,大多数网页服务器真的是靠这种方法支撑着运作。 但是把目光放到现代,像 UKVEngine 这样主打高并发的服务端, 面对数以万计的并发,这种多线程模型会暴露出巨大的弱点

  • 在很多 Linux 发行版上,默认线程栈大小常见为 8MiB,即使可以调小栈大小, “一个连线一个线程”的模型在高并发下仍会带来显著内存与调度开销。

  • 线程间的切换,CPU 调度开销绝对不算小。因而在 10000 个线程中互相切换, 大多数时间都将耗费在等待和调度中,导致真正处理用户数据的效率极低

    • 我们知道线程间的切换需要保存寄存器、计数器等上下文状态, 并且经历用户态到内核态的切换
  • 线程建立的其他开销同样不能忽视,除了刚刚提到的内存分配,还有系统调用的开销, 建立线程需要调用系统 API,这同样需要经历用户态 - 内核态 - 用户态的切换。

  • 新线程刚刚运行,指令还未被加载到 CPU 缓存,导致运行初期缓存命中率下降

  • 因此哪怕服务器性能足够,也无法支撑起如此巨大的并发数, 这迫使我们要找到一种新的模式

进化尝试二:非阻塞 I/O 与忙轮询

既然多线程模型在面对海量并发的时候行不通, 那看来我们还是要回到单线程的路线。再次回想我们初版服务端面临的问题: 程序被 acceptread 调用阻塞,导致程序死等连线或数据。 那么我们不禁思考,是否能让程序在原本被阻塞的时候,去做一些别的事情, 比如去管理其他的连线

非阻塞 I/O 的配置

在 Linux 中,我们可以通过 fcntl 调用来改变文件描述符 (fd) 的行为。 现在介绍通过 fcntl 修改 fd 为非阻塞的写法:

1
2
int flags = fcntl(fd, F_GETFL, 0);
fcntl(fd, F_SETFL, flags | O_NONBLOCK);

设置为非阻塞 I/O 后,readaccept 调用的行为会发生翻天覆地的变化, 原本会导致阻塞的情形(队列里没有连线、Socket 里没有数据), 现在会直接返回 -1,并且带有错误码 EAGAIN。 表明现在没有数据,稍后再来

忙轮询服务端的写法

借助非阻塞 I/O,我们拥有了轮询所有连线的能力,我们用一个循环来管理连线:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
int flags = fcntl(server_fd, F_GETFL, 0);
fcntl(server_fd, F_SETFL, flags | O_NONBLOCK);
std::vector<int> client_fds;

for (;;)
{
int addrlen = sizeof(address);

int client_fd = accept4(server_fd, (struct sockaddr*)&address, (socklen_t*)&addrlen, SOCK_NONBLOCK);
if (client_fd > 0) {
client_fds.push_back(client_fd);
}

for (auto it = client_fds.begin(); it != client_fds.end(); ) {
char buffer[1024] = {0};
ssize_t bytes_read = read(*it, buffer, 1024);

if (bytes_read > 0) {
write(*it, "PONG\n", 5);
}
else if (bytes_read < 0 && errno == EAGAIN) {}
else {
close(*it);
it = client_fds.erase(it);

continue;
}

++it;
}
}

close(server_fd);

补充:accept4 调用介绍

我们在代码里使用了 accept4。它大体上与 accept 调用类似, 但允许我们额外传入 Socket 标志, 常用的有 SOCK_NONBLOCKSOCK_CLOEXEC

1
int client_fd = accept4(server_fd, ..., SOCK_NONBLOCK | SOCK_CLOEXEC);

可以原子性为新建立的 fd 建立标志,大致等效于这几行代码:

1
2
3
4
5
6
7
int client_fd = accept(server_fd, ...);

int file_status_flags = fcntl(client_fd, F_GETFL, 0);
fcntl(client_fd, F_SETFL, file_status_flags | O_NONBLOCK);

int fd_flags = fcntl(client_fd, F_GETFD, 0);
fcntl(client_fd, F_SETFD, fd_flags | FD_CLOEXEC);

在保证功能相同的情况下,accept4 仅需一次调用系统调用开销更小。 非阻塞 I/O 通常都更倾向于使用 accept4

忙轮询的缺陷

我们朝着问题的解决又推进了一步,并且看似非常完美: 通过非阻塞 I/O 搭配轮询,我们免去了高并发下线程建立、切换的巨大开销, 用一个循环就照顾到了所有连线。

但是这仍然存在问题。

致命代价:空转与忙轮询

假设此时有 10000 个客户端连进来了,但它们都是“安静的听众”, 谁也不发消息。 我们的代码就会在这个 for (;;) 循环里疯狂空转,不断地调用 read, 不断地拿到 EAGAIN

这种毫无意义的轮询会在一瞬间让系统 CPU 飙满到 100%。 整个服务器都在做无用功,白白浪费电量和算力。

寻找更理想的解决方案

我们梳理一下当下的矛盾:

  • 阻塞 I/O:线程会被卡死无法处理其他连线。

  • 非阻塞 I/O + 用户态轮询:线程不会卡死,但是 CPU 空转

类比 Linux 内核唤醒阻塞 I/O 调用的流程:借助等待队列 - 运行队列的移动, 以及线程状态的变化,实现了线程从阻塞到继续运行的过程。 并且令人兴奋的是,内核知道要唤醒哪个线程,这个过程不是 CPU 轮询的结果。

那么我们有理由相信:对于一个 fd,内核一定也有某种办法,让我们不使用轮询, 而是让内核去了解哪个 fd 发生了什么事件,避免用户态空转。

进一步,我们希望存在这样的流程:我们把 10000 个 fd 交给系统内核监听, 当所有的 fd 都没有数据,我们的线程就睡眠;一旦有若干个 fd 来了数据, 内核会唤醒线程,并且精准告诉我们哪些 fd 来的数据, 我们根据这份“名单”去处理对应 fd 的事件。

简单来说,就是在单个线程(或进程)内同时管理多个 Socket 连线, 有数据时才会通知程序处理。这种设计思想叫做 I/O 多路复用

在漫长的 Linux 演进中,这个“看门大爷”的角色发生了三次更迭:

特性/指标 select poll epoll
底层数据结构 Bitmap 链表 红黑树 + 双向就绪链表
时间复杂度 O(N) O(N) epoll_ctl O(logN),epoll_wait O(k)
就绪 fd 得知 遍历 fd 列表 遍历 fd 列表 遍历就绪列表
触发模式 LT LT LT & ET

select, poll 的无效遍历

select 和 poll 需要主动轮询整个 fd 列表,具体来说:

  • 用户态拷贝:每次调用 select / poll,都需要把整个 fd 列表传入内核;

  • 内核盲目轮询:内核不知道哪个 fd 有数据,需要遍历列表并检查状态;

  • 程序二次轮询:有数据后,内核修改 fd 状态并返回, 但是不告诉哪个 fd 变了,需要我们在代码手动遍历,找出究竟是哪个 fd 就绪;

  • O(N) 复杂度:如果有 10000 个 fd,至少需要遍历 20000 次。

相比之下,epoll 的优势:

  • 不重复拷贝:程序只需要调用 epoll_ctl 一次, 把 fd 加入内核的红黑树中,后续不需要再次传入;

  • 内核事件回调:当某个 fd 来了数据,内核的网络栈会自动触发对应的回调, 这个回调函数会自动把这个 fd 放到就绪链表中;

  • 精准通知:当程序调用 epoll_wait 的时候,不扫描全部监听 fd, 只处理 ready list 中的就绪项,如果不为空, 则直接把链表中的 fd 拷贝返回给调用处;

  • 极致时间复杂度:只要有 1 个连线活跃, epoll_wait 就能在 O(1) 的时间内返回,而不会随 fd 数量增长而增长。

epoll 的用法与原理

目前 Linux 系统性能最好的 I/O 多路复用实现就是 epoll, 现在我们介绍 epoll三大核心函数,以及配套使用的状态/事件宏

使用 epoll 的函数与宏,需要包含 sys/epoll.h 头文件。

三大核心函数

epoll 的建立、控制、等待分别由三个函数负责:

建立 epoll 实例

1
2
int epoll_create(int size);
int epoll_create1(int flags);

create 系列函数用于在内核开辟空间,建立底层的红黑树和就绪链表。

  • size:原本是用于提示这颗红黑树的大小。 这个参数在现代 Linux 已被忽略,使用时传入一个大于 0 的值即可。

  • flags:现代 Linux 更推荐使用 epoll_create1,需要传入一个标志, 通常可以选择传入 0EPOLL_CLOEXEC

    • EPOLL_CLOEXEC:当前进程调用 fork 并执行 exec 时, 自动关闭这个 epoll_fd,防止 fd 泄漏。fork 本身会继承 fd。
  • 返回值:调用成功后会返回一个 epoll_fd,代表整个 epoll 实例, 使用完毕后需要 close 释放内核资源

控制红黑树的增、删、改

1
int epoll_ctl(int epfd, int op, int fd, epoll_event* event);

epoll_ctl 用于控制内核红黑树上的 fd 节点。

  • epfd:使用 create 系列函数拿到的 epoll_fd

  • op:操作类型,通常传入这三个宏:

    • EPOLL_CTL_ADD:注册一个 fd,新增至红黑树;

    • EPOLL_CTL_MOD修改一个已存在 fd 的监听事件;

    • EPOLL_CTL_DEL:取消监听一个 fd,从红黑树中删除

  • fd:要操作的 fd。

  • event:指向 epoll_event 结构体的指针,用于告诉内核要监听什么事件。

等待事件

1
int epoll_wait(int epfd, epoll_event* events, int maxevents, int timeout);

epoll_wait 通常会阻塞当前线程,直到监听的 fd 有事件发生,或者超时。

  • epfd:使用 create 系列函数拿到的 epoll_fd

  • events:用户分配的 epoll_event 数组, 内核会从就绪链表中拷贝最多 maxevents 个事件到传入的 events 数组中, 并返回实际拷贝的事件数量。

  • maxevents:告诉内核 events 数组最多能装多少个事件

  • timeout:设置超时时间(毫秒),其中:

    • > 0,设置对应的超时时间;

    • 0,不阻塞,立刻返回

    • -1持续阻塞至有事件;

  • 返回值:

    • > 0:实际就绪的 fd 数量,同时也对应着 events 数组中的事件个数;

    • 0:超时返回;

    • -1:发生错误。

结构体定义

在调用 epoll_ctlepoll_wait 的时候,都会接触到 epoll_event, 这是一个结构体,声明如下:

1
2
3
4
5
6
7
8
9
10
11
struct epoll_event {
uint32_t events; /* epoll 事件宏 */
epoll_data_t data; /* 用户自定义数据 */
};

typedef union epoll_data {
void* ptr; /* 通常绑定一个新结构体 */
int fd; /* 直接存 fd */
uint32_t u32;
uint64_t u64;
} epoll_data_t;

其中需要注意的是:epoll_data_t 是一个联合,这也就意味着绝对不能 同时给 ptrfd 赋值。另外 ptr 指向的内存通常是自行分配的, 需要在移出红黑树或 close手动释放内存。

取舍:fd 还是 ptr

既然如此,那么不仅要问:我们在 data 中是要传 ptr 还是 fd? 事实上二者的差异无外乎这几点:

特性 data.fd data.ptr
信息量 4 字节整数 无限制(可以指向任意大小对象)
内存开销 一般需要手动管理指向的内存
配套资源查找 使用 fd 查表,O(log n) 或哈希表 O(1) O(1)
适合场景 短连线、LT、简单处理 长连线、ET、复杂或高性能框架

UKVEngine 的场景事实上更适合后者,这也是 UKVEngine 后续迭代的重点。

典型的实现可能是这样:

1
2
3
4
struct ClientContext {
int fd;
RespParser parser;
};

这样就不用拿着 client_fd 去查哈希表找 parser 了, 复杂度可以从哈希表层面的“算法 O(1)”变为物理层面的 O(1), 同时也可以省去哈希表的锁

常用的宏

这些宏通常用在 epoll_ctlop 参数中,或者 event.events 中。

操作类型宏(epoll_ctlop 参数)

  • EPOLL_CTL_ADD:将 fd 加入内核的监听,新增红黑树节点。

  • EPOLL_CTL_DEL:将 fd 移出内核的监听,删除红黑树节点。

  • EPOLL_CTL_MOD:修改 fd 的监听事件,例如从监听“可读”切换到“可写”, 或者在配置为 EPOLLONESHOT 的事件被触发后,再次注册。

事件类型宏(event.events

  • EPOLLIN:对应的 fd 可读,例如有数据流入、连线关闭或 EOF, 新连线进入等。

  • EPOLLOUT:对应的 fd 可写,例如内核的发送缓冲区未满。

  • EPOLLRDHUP:对端中断了连线或关闭了写端, 实现了可以在 read() 返回 0 之前就处理断开事件。

  • EPOLLPRI:有紧急数据可读。

  • EPOLLERR:对应的 fd 发生错误 (无需手动注册,epoll_wait 自动触发)。

  • EPOLLHUP:对应的 fd 被挂断/断开连线(同样无需手动注册)。

其中 EPOLLRDHUPEPOLLHUP 都表示连线断开, 但是其中有一些细节不同:

  • EPOLLRDHUP 通常代表对端发送了 FIN 包,代表第一次挥手, 或者关闭了写半边。对端不会再发送数据, 但本机仍然可以继续读取发送剩余数据, 在处理完毕后从 epoll 删除这个 fd 并 close。 这可以作为半关闭的检测。

    • 这需要手动注册
  • EPOLLHUP 通常代表挂断事件,包括收到 RST,以及其他未知原因被挂断。 只要触发这个事件,就代表这个 Socket 已经不可用了。 原则上直接从 epoll 删除并 close 即可;但对于 stream socket, 也应结合 EPOLLIN/read 返回值以及 EPOLLERR/SO_ERROR 判断, 避免丢弃尚未读完的数据。

    • 默认监听,不需要手动注册这个事件。

行为控制宏(event.events

  • EPOLLET:将 epoll 设置为边缘触发, 只有在 fd 状态发生变化时才产生通知;它常用于高性能网络程序。 注意一点:LT 模式也可以写出高性能服务端,只是 ET 对事件处理的要求更严格。

  • EPOLLONESHOT:单次监听。确保一个 Socket 在同一时间内, 只能由一个线程独占处理。这有效解决了多个线程同时处理一个 Socket 的 数据错乱问题。

  • EPOLLEXCLUSIVE:主要用于多个 epoll 实例监听同一个 fd 时减少惊群, 典型场景是多 worker 监听同一个 listen socket。

epoll 的底层原理

讲清楚了 epoll 的用法,现在我们深入系统内核,来看看系统为我们准备了什么。

我们要知道一点:在内核中,传进来的 fd 都会被包装成一个 epitem, 它大致包含了这样几个字段:rbn 红黑树节点、rdllink 双向链表节点、 ffd fd 数据、event 注册的事件类型。

红黑树:高效的 fd 登记处

在内核态充当一个高效的常驻查找表,我们使用 EPOLL_CTL_ADD 添加的每个 fd, 就会变成 epitem 挂在这棵树上。

它解决的问题以及优点:

  • 避免全量拷贝selectpoll 每次调用,都要拷贝一万个 fd 到内核, 而红黑树确保了 fd 只在 epoll_ctl 调用的时候拷贝一次,后续不再传入。

  • 高效的增、删、改:相比于其他数据结构, 红黑树可以做到稳定的 O(log n) 时间复杂度,而其他的数据结构:

    • 数组、链表:随机存取(未知索引)的时间复杂度是 O(N),这是重走 selectpoll 的老路,不可取。

    • 哈希表:虽然哈希表有算法上的 O(1) 时间复杂度,但是这并不稳定。 以拉链法实现为例,需要预先维护较大的桶数组, 这对于珍贵的内核空间无法接受; 同时会有 Rehash 的问题,发生瞬时延迟升高; 当攻击者恶意构造 Socket 制造大量哈希碰撞,哈希表会退化成链表, 时间复杂度退化为 O(N)。

    • B+ 树:大材小用。B+ 树的强项在于范围查找和减少磁盘 I/O。 而 epoll 的环境下纯粹是内存访问,不涉及磁盘, 并且也没有类似“查找 10-100 fd 的所有连线”这样的需求。 另外,B+ 树会进行节点的分裂与合并,这相比于红黑树的变色和旋转开销更大。

  • 事件订阅中心epitem 结构体不仅记录 fd,也会向网络栈注册一个回调。 当数据来临,内核会触发 Socket 等待队列上的回调函数,直接找到 epitem, 把节点放到就绪链表。

双向就绪链表:精准的交作业名单

这个双向链表只存已经发生事件的 fd,等待 epoll_wait 高效取出。

回调函数会把 epitem 放在就绪链表的结尾, 此时如果发现等待队列仍有线程,此时会将其唤醒。 epoll_wait 同时把就绪链表中数据拷贝到 events 数组里。

一个 Socket 的流转

我们继续深入系统内核,以内核的视角来观察数据包从网卡到程序的流转过程。

展示 epoll 事件触发流程
  1. 当来自客户端的数据包被服务器的网卡接收,会率先写入内存的环形缓冲区。 随后触发中断解析封包内容。

  2. 内核的 TCP/IP 协议栈会解析封包,确认这个封包的 Socket 归属后, 把数据移到该 Socket 的接收缓冲区中。随后触发 Socket 等待队列的事件。

  3. 等待队列上挂着我们当初使用 epoll_ctl自动生成的一个回调函数, 此时就会调用这个回调。

  4. 由于精妙的结构体设计, 回调可以直接拿到该 Socket 对应的 epitem 对象。 随后按顺序做这两件事情:

    • 加入就绪链表:检查 epitem 内的 rdllink 是否已经在链表中, 如果不在则挂在就绪链表中。

    • 唤醒线程:检查 epoll 实例自己的等待队列,这里面装着由于 epoll_wait 而处于阻塞状态的线程。回调会唤醒对应的线程, 状态从阻塞 (TASK_INTERRUPTIBLE) 变为可运行 (TASK_RUNNING)。

  5. 返回用户态epoll_wait 拷贝就绪队列内的事件到 events 数组中, 此时最多拷贝 maxevents 个事件,这在调用 epoll_wait 时传入。 随后 epoll_wait 返回,等待后续的 read 等操作。

epoll 的应用与 Reactor 模型初步

了解 epoll 的使用方法和原理后,我们终于可以利用它来写真实的程序了。 在这里我们将利用 epoll 的三大函数来优化我们的 PING-PONG 服务端、 了解 LT 与 ET 的区别,以及简要介绍 Reactor 模型。

基于 epoll 的服务端改造

我们在“非阻塞 I/O + 轮询”服务端的基础上继续迭代,使用 I/O 多路复用, 实现真正的“高性能”服务端。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
// 设置为非阻塞 I/O
int flags = fcntl(server_fd, F_GETFL, 0);
fcntl(server_fd, F_SETFL, flags | O_NONBLOCK);

// 建立 epoll 实例,为我们准备红黑树和就绪链表
int epoll_fd = epoll_create1(0);

// 填写监听事件信息 (accept) 并注册
epoll_event event;
event.events = EPOLLIN;
event.data.fd = server_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, server_fd, &event);

// epoll_event 回传数组
epoll_event events[16] = {};

for (;;)
{
int addrlen = sizeof(address);

// 在此等待事件触发:新连线 (accept) 与数据接收 (client)
// 若等待到事件则继续向下进行,否则在此阻塞
// 返回就绪的 fd 数量,用这个数量遍历回传数组
int ready_count = epoll_wait(epoll_fd, events, 16, -1);

for (int i = 0; i < ready_count; i++) {
int active_fd = events[i].data.fd;

// accept4 触发了 epoll
if (active_fd == server_fd) {
int client_fd = accept4(server_fd, (struct sockaddr*)&address, (socklen_t*)&addrlen, SOCK_NONBLOCK);
if (client_fd < 0) continue;

// 连线建立,把客户端发消息事件注册进 epoll
epoll_event client_event;
client_event.events = EPOLLIN;
client_event.data.fd = client_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, client_fd, &client_event);
}
// 客户端发消息触发了 epoll
else {
char buffer[1024] = {0};
ssize_t bytes_read = read(active_fd, buffer, 1024);

if (bytes_read > 0) {
write(active_fd, "PONG\n", 5);
}
else if (bytes_read < 0 && errno == EAGAIN) {

}
else {
epoll_ctl(epoll_fd, EPOLL_CTL_DEL, active_fd, nullptr);
close(active_fd);
}
}
}
}

close(epoll_fd);
close(server_fd);

尝试运行并使用 nc 与服务端连线,体验上与之前的“非阻塞 I/O + 轮询”服务端 没有差别,都做到了 readaccept 不阻塞第二个 nc 的连线。 但是观察前后两次的 top 命令结果,前者的 CPU 是接近满载的, 后者最多 0.1% 的占用,差距悬殊。

介绍一下代码逻辑:

  • 使用 socket 建立一个 server_fd,用于 accept4 的监听。

    • 同时利用 fcntl 设置为非阻塞 I/O。
  • 使用 epoll_create1 建立一个 epoll 实例,准备红黑树和就绪链表。

  • server_fd 注册为事件,事件类型为 EPOLLIN。当可读时触发响应。

    • 即有新连线连入时触发。
  • 进入主循环,在此等待事件发生。包括连入事件和客户端发消息的事件。

    • 客户端发消息事件是后面注册的。
  • 当事件触发,epoll_wait 把就绪的所有 event 拷贝到 events 数组中。

  • 利用 active_fd 判断事件类型,如果是:

    • server_fd,说明是由新连线连入触发了事件,接受这个连线; 利用拿到的 client_fd,把这个连线的后续发消息事件注册进 epoll, 事件类型为 EPOLLIN。当可读时触发响应。

    • client_fd,说明是已经建立连线的客户端发来了消息, 服务端回复一个 PONG。

      • 如果 read 返回值不大于 0,说明连线异常。我们取消注册这个连线, 并且关闭这个 fd。

epoll 边缘触发

在介绍 epoll 的行为控制宏的时候,我们介绍了 EPOLLET 宏, 讲到这个宏用于将 epoll 的触发模式切换为边缘触发 (ET),如果不指定这个宏, 则默认为水平触发 (LT)。

从行为上,或者直观上来讲, LT 会在就绪链表非空的时候持续触发 epoll_wait。 如果不能一次读完缓冲区的内容,LT 允许我们一次一读,慢慢地读干净; ET 模式下,事件主要在 fd 状态发生变化时产生; 如果本次没有把数据读到 EAGAIN,剩余数据可能不会再次触发新的通知。 因此工程上通常要求 fd 非阻塞,并在一次事件处理中循环读/写直到 EAGAIN

以 Socket 为例,LT 与 ET 在内核处理上的差别:

  • LTepoll_wait 拷贝完某个 epitem 后, 会额外检查这个 Socket 的接收缓冲区是否还有残留数据,如果仍有, 则会把这个 event 再次塞回就绪链表中;

  • ETepoll_wait 拷贝后,直接把 fd 从就绪链表拔除,不关心读取情况。

因此在使用 ET 的时候,就有几个必须遵守的规则:

  • 循环读取:必须使用一个循环持续读取, 直到返回 -1errno == EAGAIN,表明彻底读空了,才能退出循环。 否则由于 fd 已经从就绪链表拿掉,旧数据将永远留在内核。

  • 非阻塞 I/O:由于需要循环读取,如果是阻塞的 read, 读空后的调用不会返回 -1,而是会直接阻塞线程,导致整个业务流程瘫痪。

因此使用 ET 的 PING-PONG 服务端要改成这样:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
// 非阻塞 I/O,使处理 server_fd 的 accept4 不阻塞线程
int flags = fcntl(server_fd, F_GETFL, 0);
fcntl(server_fd, F_SETFL, flags | O_NONBLOCK);

int epoll_fd = epoll_create1(0);

epoll_event event;

// 加上行为控制宏 EPOLLET
event.events = EPOLLIN | EPOLLET;
event.data.fd = server_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, server_fd, &event);

epoll_event events[16] = {};

for (;;)
{
int addrlen = sizeof(address);
int ready_count = epoll_wait(epoll_fd, events, 16, -1);

for (int i = 0; i < ready_count; i++) {
int active_fd = events[i].data.fd;

if (active_fd == server_fd) {
// 循环 accept
for (;;) {
// 非阻塞 I/O,使处理 client_fd 的 read 不阻塞线程
int client_fd = accept4(server_fd, (struct sockaddr*)&address, (socklen_t*)&addrlen, SOCK_NONBLOCK);

// 直到读空,退出死循环
if (client_fd < 0 && errno == EAGAIN) {
break;
}
else if (client_fd < 0) {
std::cerr << errno << ": " << strerror(errno) << std::endl;
break;
}

epoll_event client_event;

// 加上行为控制宏 EPOLLET
client_event.events = EPOLLIN | EPOLLET;
client_event.data.fd = client_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, client_fd, &client_event);
}
}
else {
// 循环 read
for (;;) {
char buffer[1024] = {0};
ssize_t bytes_read = read(active_fd, buffer, 1024);

if (bytes_read > 0) {
write(active_fd, "PONG\n", 5);
}
// 直到读空,退出死循环
else if (bytes_read < 0 && errno == EAGAIN) {
break;
}
else {
epoll_ctl(epoll_fd, EPOLL_CTL_DEL, active_fd, nullptr);
close(active_fd);
break;
}
}
}
}
}

close(epoll_fd);
close(server_fd);

Reactor 模型

Reactor 模型是一种事件驱动的网络程序设计模型,在高并发 I/O 中有广泛的应用。 核心做法就是利用 I/O 多路复用来监听多个事件,并将事件分发给对应的处理器, 实现高效、无阻塞的处理。

稍微思考一下 epoll 三大函数的使用方法,就不难发现: epoll 把整个流程封装成三个函数后, 其实是在推动着我们在写 Reactor 模式, 或者说用上 epoll 后,就很难写出不遵循 Reactor 模型的程序。 所以,我们使用 epoll 改造过的 PING-PONG 服务端程序, 毫无疑问遵循 Reactor 模式。现在我们深入介绍这个设计模式。

Reactor 模型的角色划分

Reactor 模型分为三个核心角色,结合我们的服务端程序:

  • Reactor:负责监听与分配事件,当事件发生时,会把事件分发给 Handler。

    • epoll_wait 与外层 for(;;) 担任监听事件的任务, for (int i = 0; i < ready_count; i++) 担任事件分发的任务。
  • Acceptor:专门负责处理客户端的连线请求。

    • if (active_fd == server_fd) 的分支担任处理连线的任务, 使用 accept4 接受连线,并通过 epoll_ctl 注册回 epoll
  • Handler:负责执行实际 I/O 操作与业务逻辑

    • else 的分支担任实际操作的任务,read 调用读取客户端输入、 write 调用给出回应、close 调用关闭连线。

Reactor 模型的常见架构

根据 Reactor 和线程的数量,大致可以分为三类:

  • 单 Reactor 单线程(PING-PONG 服务端):Reactor, Acceptor, Handler 都在同一个线程内,编码简单但无法发挥多核 CPU 的性能; 如果 Handler 卡住则会影响整个业务流程

  • 单 Reactor 多线程(UKVEngine):Handler 的业务被拿到线程池处理, 能够充分利用多核资源。但单一 Reactor 在高并发下可能会成为瓶颈

  • 主从 Reactor 多线程(业内主流):设置有主、从 Reactor。 主 Reactor 只负责接受连线 (Acceptor),从 Reactor 有自己的 epoll_fd, 主 Reactor 接受连线后挑选一个从 Reactor 线程注册进去。

多线程编程与线程安全

前面在介绍 Reactor 架构的时候提到了线程池,一个简易的线程池的实现如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
class ThreadPool {
public:
explicit ThreadPool(size_t thread_num) : running_(true) {
for (size_t i = 0; i < thread_num; i++) {
workers_.emplace_back([this] {
for (;;) {
std::function<void()> task;

{
std::unique_lock<std::mutex> lock(mutex_);
cv_.wait(lock,
[this] { return !running_ || !tasks_.empty(); });

if (!running_ && tasks_.empty()) {
return;
}

task = std::move(this->tasks_.front());
tasks_.pop();
}

task();
}
});
}
}

ThreadPool(const ThreadPool&) = delete;
ThreadPool& operator=(const ThreadPool&) = delete;

template<typename T>
void EnQueue(T&& f) {
{
std::unique_lock<std::mutex> lock(mutex_);
tasks_.emplace(std::forward<T>(f));
}

cv_.notify_one();
}

~ThreadPool() {
{
std::unique_lock<std::mutex> lock(mutex_);
running_ = false;
}

cv_.notify_all();
for (auto& worker : workers_) {
worker.join();
}
}

private:
std::vector<std::thread> workers_;
std::queue<std::function<void()> > tasks_;

std::mutex mutex_;
std::condition_variable cv_;

bool running_;
};

我们在这里不过多介绍线程池本身的实现。这里简单介绍一下线程池的优点:

  • 降低开销:线程的建立、销毁的开销较大。如果随时建立、销毁线程, 会带来不小的系统开销。线程池维护若干个现有线程,不必频繁建立与销毁。

  • 加快响应:业务代码直接使用已有的线程,响应速度较快,无需等待线程建立。

  • 统一管理:线程池可以对线程进行统一的监控与管理,避免无限制创建。

引入多线程后,充分发挥 CPU 的多核性能,提高并发。 但如果利用不当,会引起较为严重的问题:惊群效应线程安全问题

惊群效应

在 epoll/Reactor 的语境下,惊群效应大致可以分为两个场景:

accept 惊群

  • 发生场景早期或某些实现中, 多个线程/进程阻塞在同一个 listen socket 的 accept 上时, 一个新连线可能唤醒多个等待者,最终只有一个 accept 成功,造成惊群。

    现代 Linux 的阻塞 accept 路径, 已经通过 exclusive wait / wake-one 机制基本缓解了传统 accept 惊群。

    但这并不意味着所有与 listenfd 相关的惊群都不存在: 如果多个 epoll 实例同时监听同一个 listenfd, 仍然需要考虑 epoll 层面的事件通知惊群, 此时可以考虑 EPOLLEXCLUSIVESO_REUSEPORT、accept mutex, 或单 Acceptor。

epoll_wait 惊群

这又分为两个场景和一个变种:

  • epoll 实例,多个线程在 epoll_wait 阻塞

    • 发生场景:多个线程都有这样的代码:

      1
      epoll_wait(epoll_fd, events, maxevents, -1);

      在 LT 或事件未被及时消费的场景下,多个线程共享同一个 epoll 实例, 可能导致多个线程先后拿到同一个 fd 的就绪事件, 形成类似惊群或重复处理问题;

      而对于同一个 epoll_fd + EPOLLET, Linux 有只唤醒一个等待者的优化。

    • 解决方案

      • 单 Acceptor 或 accept 锁:本质相同。控制只有一个线程才能去 accept。

        • 通常用于解决 listen_fd 的 accept 分发问题。
      • SO_REUSEPORT:建立多个 Socket,因此 fd 不同。 但是这些 Socket 绑定的地址和端口都一致,因此开启端口复用, 借助负载均衡来避免惊群。

        • 通常用于解决 listen_fd 的 accept 分发问题。
      • EPOLLET:Linux 内核特性,如果多个线程阻塞在同一 epoll_fd, 并且使用 ET 触发,那么只会唤醒一个等待者。

  • epoll 实例,多个线程 epoll_ctl 注册同一个 fd

    • 发生场景:每个工作线程都有自己独立的 epoll 实例, 但是都使用 EPOLL_CTL_ADD 监听了同一个 fd。默认情况下, 多个线程都能收到事件。

    • 解决方案

      • EPOLLEXCLUSIVE:如果没有这个标志, 一个新连线到来会唤醒所有等待的进程,造成资源浪费; 设置了 EPOLLEXCLUSIVE,多个 epoll 实例监听同一个目标 fd 时, 内核会以排他唤醒方式减少惊群,通常唤醒一个或少数等待者, 但语义上是 one or more,不是严格保证永远只唤醒一个。
  • 变种:一个事件准备塞入线程池执行业务,多个线程同时去处理事件业务

    • 发生场景:这不属于典型的惊群效应,但是的确是一个并发设计问题。 对于一个已经拿到的 fd,如果多个工作线程都去处理这个事件, 会造成数据错乱的问题。比如在 LT 下,同一个 fd 持续通知, 多个线程都参与了处理。

    • 解决方案

      • EPOLLONESHOT:表示某个 fd 的事件被 epoll_wait 返回后, 该 fd 会在红黑树中被临时禁用;后续事件不会再通过 epoll 返回, 直到调用 EPOLL_CTL_MOD 重新 armed。

线程池条件变量惊群

这个不是很常见,属于线程池设计失误。

一般出现这个问题的根源,在于混淆notify_onenotify_all。 导致本应唤醒一个线程的操作,被误操作为唤醒所有线程。

线程安全

通常认为,如果无论操作系统如何调度线程,程序都能正确处理线程间共享的变量, 都能表现出符合逻辑、符合预期的行为,那么我们说这段代码是线程安全的。

线程安全问题的罪魁祸首就是共享的写操作,比如:

  • 线程 A、B 都要对一个变量 cnt +1 操作。分为两个步骤:读旧值、写新值。 按照逻辑和预期,cnt 最终要被 +2,但是可能发生这样的问题:

    • 线程 A 读取到了 cnt 的旧值为 100,准备修改为 101;

    • 极短时间后,线程 B 此时也要修改,由于 A 还没有做出修改操作, 因此 B 读到的旧值居然也是 100。两个线程都要写成 101。不合我们的预期。

  • 线程 A 在对 buffer 写入字符串,线程 B 准备读取字符串, 期间可能发生这样的问题:

    • 线程 A 准备写入 "Hello World",没来得及写完, 只写了 "Hello W"

    • 此时线程 B 准备读取 buffer 的内容,读到的不是 "Hello World", 而是 "Hello W",直接导致业务出错。

以上的例子我们可以说是线程不安全的代码。保证线程安全的措施一般有:

  • 不共享、不可变对象

  • 原子变量

  • 条件变量

  • 信号量

这里我们不去详细介绍每一种措施的原理。 而是选择最常用的锁、原子变量等进行介绍,直接以实际应用入手, 告诉读者应该如何选择与使用。

不共享、不可变对象

这种想法非常简单粗暴:既然线程不安全的罪魁祸首就是共享的写, 那干脆不去共享对象,将对象对每个线程私有,自然解决了安全问题。

或者这个对象在运行时完全不变,比如配置文件、元数据等, 那么这种对象多线程共享也是绝对安全的,因为不涉及写操作。

通常能接触到的锁有互斥锁共享锁,以及自旋锁

在 C++ 中,前两种锁在标准库中有实现: std::mutexstd::shared_mutex,而自旋锁也可以简单地实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
#include <atomic>

class SpinLock {
private:
std::atomic_flag flag = ATOMIC_FLAG_INIT;

public:
void lock() {
// 当 test_and_set 返回 true 时,说明锁已经被其他线程占用,继续循环(自旋)
while (flag.test_and_set(std::memory_order_acquire)) {}
}

void unlock() {
// 释放锁
flag.clear(std::memory_order_release);
}
};

同时 C++ 实现了 RAII 风格的锁管理器:

  • std::lock_guard:最为轻量, 仅支持初始化时 lock 与析构时unlock,适合自动管理 std::mutex。 不适合多锁管理与读写锁管理。

    • 我们给出的自旋锁实现也可以使用 std::lock_guard 管理。
  • std::unique_lock:功能最全, 不仅支持 RAII 风格,同时也支持延迟上锁、提前解锁。 并且条件变量的锁和读写锁排他写也需要使用 std::unique_lock

  • std::shared_lock:为读写锁共享读设计。

  • std::scoped_lock:多锁友好,用于解决锁配置不当导致的死锁问题, 因此它支持传入任意数量的互斥锁,并且从根本上杜绝死锁问题。

互斥锁

这属于最基础、最简单的线程安全手段,含义是同一时间只允许一个线程进入临界区。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
#include <mutex>

std::mutex mtx;

void update() {
// Do something...

{
std::lock_guard<std::mutex> lock(mtx);
var_name++;
}

// Do something...
}

互斥锁适合保护:

  • 共享 map

  • 共享队列

  • 日志缓冲区

优点是简单可靠。

缺点也很明显:

  • 可能造成阻塞

  • 可能造成锁竞争

  • 可能造成死锁

  • 粗粒度锁会降低并发性能

读写锁

支持共享读排他写的锁,适合读多写少的场景。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
#include <shared_mutex>

std::shared_mutex mtx;

void get_var() {
// Do something...

{
std::shared_lock<std::shared_mutex> lock(mtx);
parser = client_parser[client_fd];
}

// Do something...
}

void set_var() {
// Do something...

{
std::unique_lock<std::shared_mutex> lock(mtx);
client_parser.erase(client_fd);
}

// Do something...
}

读写锁适合保护:

  • 配置表

  • 路由表

  • 用户缓存

  • 其他读多写少的数据

如果写操作很多,性能可能不如普通互斥锁好。

自旋锁

纯用户态实现。不会让线程睡眠,而是循环等待锁释放。

自旋锁适合:

  • 临界区极短(改少量状态变量)

因为它不会进行线程切换,不会陷入内核态,能够减少一部分开销。 或者内核、底层编程下很常见。

使用时需再三确定临界区等待时间,否则会造成 CPU 空转,这个时候更适合互斥锁。

原子变量

对于简单的变量,可以考虑使用原子操作来避免加锁

例如:

1
2
3
4
5
6
7
8
9
#include <atomic>

std::atomic<bool> g_running{false};

void Run() {
g_running = true;

// Do something...
}

原子变量通过硬件原子指令、CAS 循环或平台相关机制, 保证对单个变量的操作不会发生数据竞争。

原子变量适合:

  • 计数器

  • 状态标志

  • 引用计数

  • 开关

原子变量只适合保护单个变量简单状态,不适合保护复杂对象的不变量。

PING-PONG 服务端的多线程改造

现在就把我们的服务端改造成“单 epoll + 多线程”的模式。

我们还记得,使用多线程解决的问题主要就是防止一个长 Handler 操作阻塞线程, 因此改造思路很明确,我们把 readwrite 操作移入线程池中。 为了体现线程安全性,我们全局维护一个计数器 cnt, 记录总共发出 PONG 的数量。

注意: 下面代码用于演示 epoll + 线程池 + EPOLLONESHOT 的核心流程, 省略了生产环境必须处理的错误检查、非阻塞写、半关闭、协议分包等细节。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
// 有 8 个工作线程的线程池
ThreadPool pool(8);

int cnt = 0;
std::mutex mtx;

int flags = fcntl(server_fd, F_GETFL, 0);
fcntl(server_fd, F_SETFL, flags | O_NONBLOCK);

int epoll_fd = epoll_create1(0);

epoll_event event;
event.events = EPOLLIN | EPOLLET;
event.data.fd = server_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, server_fd, &event);

epoll_event events[16] = {};

for (;;)
{
int addrlen = sizeof(address);
int ready_count = epoll_wait(epoll_fd, events, 16, -1);

for (int i = 0; i < ready_count; i++) {
int active_fd = events[i].data.fd;

if (active_fd == server_fd) {
for (;;) {
int client_fd = accept4(server_fd, (struct sockaddr*)&address,
(socklen_t*)&addrlen, SOCK_NONBLOCK);

if (client_fd < 0 && errno == EAGAIN) {
break;
}
else if (client_fd < 0) {
std::cerr << errno << ": " << strerror(errno) << std::endl;
break;
}

// EPOLLONESHOT 防止多线程处理同一 fd
epoll_event client_event;
client_event.events = EPOLLIN | EPOLLET | EPOLLONESHOT;
client_event.data.fd = client_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, client_fd, &client_event);
}
}
else {
// 把任务扔到线程池运行,不影响主线程
pool.EnQueue([active_fd, epoll_fd, &cnt, &mtx] {
for (;;) {
char buffer[1024] = {0};
ssize_t bytes_read = read(active_fd, buffer, 1024);

if (bytes_read > 0) {
write(active_fd, "PONG\n", 5);

// 临界区,使用互斥锁保护共享的变量
// 由于这里保护的是单变量,可以考虑用原子变量。
std::lock_guard<std::mutex> lock(mtx);
cnt++;
printf("%d\n", cnt);
}
else if (bytes_read < 0 && errno == EAGAIN) {
// EPOLLONESHOT 的重新注册
epoll_event re_event;
re_event.events = EPOLLIN | EPOLLET | EPOLLONESHOT;
re_event.data.fd = active_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_MOD, active_fd, &re_event);
break;
}
else {
epoll_ctl(epoll_fd, EPOLL_CTL_DEL, active_fd, nullptr);
close(active_fd);
break;
}
}
});
}
}
}

close(epoll_fd);
close(server_fd);

注意临界区附近,我们这里使用互斥锁保护 cnt,这里也可以使用原子变量, 因为符合简单单一变量的原则。

写在最后

这篇文章不仅是笔者对 epoll 的理解,同时也是学习多线程的来时路。 希望笔者的理解能给各位读者带来帮助,我们下一篇文章见。