文章目录
- 从 select 到 epoll:C++17 TCP 服务器的六层渐进式重构
-
- 一、为什么要采用“分层升级”
- 二、先建立正确的 epoll 心智模型
-
- 2.1 epoll 不是异步 I/O
- 2.2 interest list 与 ready list
- 2.3 epoll 并不是“完全零拷贝”
- 2.4 不要把 epoll 的复杂度概括成“所有操作都是 O(1)”
- 第一部分:第 0 层——原始 select 版本
-
- 三、第 0 层的目的
- 3.1 完整服务端:`00_select_original/server.cpp`
- 3.2 完整客户端:`00_select_original/client.cpp`
- 3.3 原始版本的具体问题
-
- 问题一:每轮复制并修改 fd_set
- 问题二:返回后仍然遍历全部客户端
- 问题三:受到 fd_set 表示能力限制
- 问题四:把 EOF 和错误混在一起
- 问题五:默认阻塞 I/O
- 问题六:没有处理部分发送
- 问题七:TCP 没有消息边界
- 第二部分:第 1 层——最小 epoll LT 迁移
-
- 四、本层目标
- 4.1 从 select 到 epoll 的最小映射
- 4.2 完整服务端:`01_epoll_lt_minimal/server.cpp`
- 4.3 完整客户端:`01_epoll_lt_minimal/client.cpp`
- 4.4 本层修复了什么
-
- 修复一:不再每轮复制完整监控集合
- 修复二:不再遍历所有空闲客户端
- 修复三:不再依赖 `max_fd + 1`
- 4.5 本层仍然存在的问题
- 第三部分:第 2 层——非阻塞 epoll LT
-
- 五、为什么 LT 也建议使用非阻塞套接字
- 5.1 本层新增的规则
-
- 规则一:监听 fd 非阻塞
- 规则二:通信 fd 也设为非阻塞
- 规则三:循环读取到 EAGAIN
- 规则四:本层遇到发送缓冲区满时直接关闭连接
- 5.2 完整服务端:`02_epoll_lt_nonblocking/server.cpp`
- 5.3 完整客户端:`02_epoll_lt_nonblocking/client.cpp`
- 5.4 本层修复了什么
- 5.5 为什么不能把 EAGAIN 当成错误
- 5.6 本层仍然存在的问题
- 第四部分:LT 与 ET 的真正区别
-
- 六、水平触发 LT
- 七、边缘触发 ET
- 第五部分:ET 的两个典型错误
-
- 八、错误一:EPOLLET + 单次 read
-
- 8.1 错误逻辑
- 8.2 完整错误复现代码
- 九、错误二:EPOLLET + 阻塞套接字 + 循环 read
- 9.1 完整错误复现代码
- 十、ET 正确读法
- 第六部分:第 3 层——正确的 epoll ET 读取
-
- 十一、本层相对第 2 层改了什么
- 11.1 完整服务端:`03_epoll_et_correct_read/server.cpp`
- 11.2 完整客户端:`03_epoll_et_correct_read/client.cpp`
- 11.3 本层修复了什么
- 11.4 本层仍然没有解决写路径
- 第七部分:为什么不能永久监听 EPOLLOUT
-
- 十二、套接字“可写”意味着什么
- 十三、为什么每个连接都需要独立发送状态
- 第八部分:第 4 层——ET + 发送队列 + EPOLLOUT
-
- 十四、本层新增的状态机
- 14.1 完整服务端:`04_epoll_et_write_queue/server.cpp`
- 14.2 完整客户端:`04_epoll_et_write_queue/client.cpp`
- 14.3 `BuildClientEvents()` 的意义
- 14.4 为什么收到数据后先主动发送
- 14.5 `EPOLLRDHUP` 不能直接等同于“所有数据已经读完”
- 14.6 本层修复了什么
- 14.7 本层的新风险:内存无限增长
- 第九部分:背压是什么
-
- 十五、背压的核心含义
- 十六、高低水位
- 十七、为什么还需要硬上限
- 第十部分:ET、公平预算和 EPOLLONESHOT
-
- 十八、为什么不能无限排空一个非常活跃的连接
- 十九、ET 下提前停止读取为什么危险
- 二十、EPOLLONESHOT 的语义
- 第十一部分:第 5 层——背压、预算与资源边界
-
- 二十一、最终层的状态字段
- 21.1 完整服务端:`05_epoll_et_backpressure/server.cpp`
- 21.2 完整客户端:`05_epoll_et_backpressure/client.cpp`
- 21.3 最终层的保护机制
-
- 连接数上限
- 单连接发送队列硬上限
- 高低水位背压
- 每事件读写预算
- EPOLLONESHOT 重装
- 21.4 最终层完整流程图
- 第十二部分:epoll 三大 API 逐项解析
-
- 二十二、`epoll_create1()`
- 二十三、`epoll_ctl()`
-
- ADD
- MOD
- DEL
- 二十四、`epoll_wait()`
- 第十三部分:常用事件标志
-
- 二十五、事件位速查表
- 二十六、为什么使用 SO_ERROR
- 第十四部分:accept、recv、send 的完整状态表
-
- 二十七、非阻塞 accept
- 二十八、非阻塞 recv
- 二十九、非阻塞 send
- 第十五部分:TCP 半关闭
-
- 三十、`shutdown(SHUT_WR)` 的意义
- 第十六部分:LT 与 ET 应该如何选
-
- 三十一、对比表
- 三十二、ET 一定比 LT 快吗
- 第十七部分:压力测试客户端
-
- 三十三、完整大数据校验客户端
- 33.1 测试 1 MiB 回显
- 33.2 为什么不能用一次 write 作为压力测试
- 第十八部分:编译与运行
-
- 三十四、编译全部版本
- 三十五、运行指定层次
-
- 第 0 层
- 第 1 层
- 第 2 层
- 第 3 层
- 第 4 层
- 第 5 层
- 第十九部分:调试方法
-
- 三十六、使用 strace 观察事件循环
- 三十七、使用 ss 查看连接
- 三十八、模拟慢客户端
- 第二十部分:常见错误清单
-
- 三十九、错误:ET 只加 EPOLLET,不改读取方式
- 四十、错误:ET 使用阻塞 fd 循环读
- 四十一、错误:把 EAGAIN 当成连接失败
- 四十二、错误:永久监听 EPOLLOUT
- 四十三、错误:忽略部分发送
- 四十四、错误:发送队列没有上限
- 四十五、错误:ET 设置处理预算但不重新激活
- 四十六、错误:EPOLLRDHUP 立即 close
- 四十七、错误:在 `epoll_event.data.ptr` 中保存裸指针却不管理生命周期
- 第二十一部分:进一步工程化方向
-
- 四十八、应用层协议拆包
- 四十九、定时器
- 五十、信号
- 五十一、多线程
- 五十二、跨线程唤醒
- 第二十二部分:FAQ
-
- 五十三、close 前必须 EPOLL_CTL_DEL 吗
- 五十四、accepted socket 会继承 O_NONBLOCK 吗
- 五十五、EAGAIN 和 EWOULDBLOCK 是否相同
- 五十六、LT 可以使用阻塞套接字吗
- 五十七、ET 是否每次都必须读到 EAGAIN
- 五十八、epoll 可以监控普通文件吗
- 五十九、epoll 能替代线程池吗
- 第二十三部分:六层演进总结
-
- 六十、每一层到底学到了什么
-
- 第 0 层:select
- 第 1 层:epoll LT 最小迁移
- 第 2 层:非阻塞 LT
- 第 3 层:ET 正确读取
- 第 4 层:写事件状态机
- 第 5 层:工程边界
- 六十一、最终核心原则
- 附录 A:源码目录
- 附录 B:官方资料名称
从 select 到 epoll:C++17 TCP 服务器的六层渐进式重构
平台:Linux 语言标准:C++17 示例类型:单线程 TCP Echo 服务器 核心目标:不是一次性给出一份复杂的“最终代码”,而是从原始 select 程序开始,一层一层升级到具备非阻塞 I/O、ET 正确读写、发送队列、背压和公平性控制的 epoll 事件循环。
一、为什么要采用“分层升级”
很多 epoll 教程存在同一个问题:
本文不采用这种跳跃式讲法,而是固定为六个可以独立编译、独立运行的版本:
| 第 0 层 | 原始 select | 建立对照基线 |
| 第 1 层 | 最小 epoll LT | 把“扫描全部 fd”改成“只返回就绪事件” |
| 第 2 层 | 非阻塞 epoll LT | 排空连接队列、区分 EAGAIN、避免单个 I/O 阻塞事件循环 |
| 第 3 层 | 正确的 epoll ET 读取 | 非阻塞套接字 + 循环处理到 EAGAIN |
| 第 4 层 | ET 发送队列 | 处理部分发送、发送缓冲区满和 EPOLLOUT |
| 第 5 层 | 背压与公平性 | 限制内存、暂停读取、处理预算、EPOLLONESHOT 重装 |
每一层都包含:
- 完整服务端代码;
- 完整客户端代码;
- 相对上一层改动了什么;
- 修复了什么;
- 本层仍然没有解决什么;
- 编译与运行方法。
二、先建立正确的 epoll 心智模型
2.1 epoll 不是异步 I/O
epoll 是一种 I/O 就绪事件通知机制。
它告诉应用程序:
某个文件描述符现在可能可以执行读取、写入或其他操作。
它不会替应用程序执行 recv() 或 send(),也不会自动创建线程。
典型执行过程仍然是同步的:
应用线程
│
│ epoll_wait()
▼
等待一个或多个 fd 就绪
│
│ 返回就绪事件
▼
应用自己调用 recv()/send()/accept()
│
▼
处理业务状态
│
└──────────────→ 再次 epoll_wait()
所以更准确的描述是:
epoll 负责“告诉你谁可能能做 I/O”,应用程序仍负责真正的 I/O 和状态管理。
2.2 interest list 与 ready list
从用户态视角,可以把一个 epoll 实例理解成两部分:
epoll 实例
┌─────────────────────────────┐
│ │
│ interest list │
│ 我关心哪些 fd、哪些事件 │
│ │
│ fd=3: EPOLLIN │
│ fd=7: EPOLLIN|EPOLLRDHUP │
│ fd=9: EPOLLIN|EPOLLOUT │
│ │
├─────────────────────────────┤
│ │
│ ready list │
│ 当前已经就绪的监控项 │
│ │
│ fd=7: EPOLLIN │
│ fd=9: EPOLLOUT │
│ │
└─────────────────────────────┘
三个核心 API 的角色分别是:
epoll_create1()
创建 epoll 实例
epoll_ctl()
修改 interest list
ADD / MOD / DEL
epoll_wait()
从 ready list 取出当前就绪事件
2.3 epoll 并不是“完全零拷贝”
select() 每轮都要把待监控的 fd_set 传给内核,并在返回时把修改后的集合带回用户态。
epoll 的优势是:
- 监控集合通过 epoll_ctl() 持久保存在内核中;
- epoll_wait() 不需要每轮重新提交整个监控集合;
- 返回时只复制本轮就绪的 epoll_event。
但 epoll_wait() 返回事件时仍然需要把就绪事件复制到用户缓冲区,所以不应该把 epoll 简化成“零拷贝”。
2.4 不要把 epoll 的复杂度概括成“所有操作都是 O(1)”
更稳妥的理解是:
- epoll_ctl() 要查询并修改内核中的监控项;
- epoll_wait() 的主要用户态处理量与返回的就绪事件数量相关;
- 应用不再像 select() 那样线性检查每一个可能的 fd;
- 内核实现细节可能随版本演化,不应把某一种内部数据结构当作稳定 ABI。
第一部分:第 0 层——原始 select 版本
三、第 0 层的目的
这一层不追求健壮性,只保留最初代码的结构,用于回答:
从 select 迁移到 epoll 时,最先发生变化的到底是什么?
原始版本的主要结构是:
master_fds 保存全部监控对象
│
▼
read_fds = master_fds
│
▼
select() 修改 read_fds
│
▼
FD_ISSET(listen_fd)
│
├─ accept()
│
└─ 遍历 clients
│
└─ FD_ISSET(client_fd)
3.1 完整服务端:00_select_original/server.cpp
#include <arpa/inet.h>
#include <sys/select.h>
#include <sys/socket.h>
#include <unistd.h>
#include <algorithm>
#include <cstring>
#include <iostream>
#include <string>
#include <vector>
using namespace std;
int main() {
// 1. 创建 TCP 监听套接字
int fd = socket(AF_INET, SOCK_STREAM, 0);
if (fd == –1) {
cerr << "socket() failed" << endl;
return –1;
}
// 2. 设置端口复用
int opt = 1;
setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt));
// 3. 绑定地址并监听
sockaddr_in addr{};
addr.sin_family = AF_INET;
addr.sin_addr.s_addr = htonl(INADDR_ANY);
addr.sin_port = htons(8888);
if (bind(fd, reinterpret_cast<sockaddr*>(&addr), sizeof(addr)) == –1 ||
listen(fd, 5) == –1) {
cerr << "bind/listen failed" << endl;
close(fd);
return –1;
}
cout << "Server listening on port 8888…" << endl;
// 4. 初始化 select 监控集合
fd_set master_fds;
fd_set read_fds;
FD_ZERO(&master_fds);
FD_SET(fd, &master_fds);
int max_fd = fd;
vector<int> clients;
while (true) {
// select 会修改 fd_set,因此每轮必须复制。
read_fds = master_fds;
int activity =
select(max_fd + 1, &read_fds, nullptr, nullptr, nullptr);
if (activity < 0) {
perror("select");
break;
}
// 5. 监听套接字可读:接收一个新连接
if (FD_ISSET(fd, &read_fds)) {
sockaddr_in client_addr{};
socklen_t len = sizeof(client_addr);
int cfd = accept(
fd, reinterpret_cast<sockaddr*>(&client_addr), &len);
if (cfd != –1) {
char ip_str[INET_ADDRSTRLEN]{};
inet_ntop(AF_INET,
&client_addr.sin_addr,
ip_str,
sizeof(ip_str));
cout << "New connection from "
<< ip_str << ':'
<< ntohs(client_addr.sin_port)
<< endl;
FD_SET(cfd, &master_fds);
max_fd = max(max_fd, cfd);
clients.push_back(cfd);
}
}
// 6. 线性遍历所有客户端,检查哪些 fd 可读
for (auto it = clients.begin(); it != clients.end();) {
int cfd = *it;
if (FD_ISSET(cfd, &read_fds)) {
string buf;
buf.resize(64);
ssize_t n = read(cfd, buf.data(), buf.size());
if (n <= 0) {
cout << "Client disconnected (fd="
<< cfd << ')' << endl;
close(cfd);
FD_CLR(cfd, &master_fds);
it = clients.erase(it);
} else {
buf.resize(static_cast<size_t>(n));
cout << "Received: " << buf << endl;
const string reply = "Hello from server";
write(cfd, reply.c_str(), reply.size());
++it;
}
} else {
++it;
}
}
}
close(fd);
return 0;
}
3.2 完整客户端:00_select_original/client.cpp
#include <arpa/inet.h>
#include <cstring>
#include <iostream>
#include <sys/socket.h>
#include <unistd.h>
using namespace std;
int main() {
// 1. 创建 TCP 套接字
int fd = socket(AF_INET, SOCK_STREAM, 0);
// 2. 设置服务器地址
sockaddr_in addr{};
addr.sin_family = AF_INET;
addr.sin_addr.s_addr = inet_addr("127.0.0.1");
addr.sin_port = htons(8888);
// 3. 连接服务器
connect(fd, reinterpret_cast<sockaddr*>(&addr), sizeof(addr));
cout << "Connected to server" << endl;
// 4. 发送一条消息
const char* msg = "Hello from client";
write(fd, msg, strlen(msg));
cout << "Sent: " << msg << endl;
// 5. 阻塞等待服务器回复
char buf[64] = {};
read(fd, buf, sizeof(buf));
cout << "Received: " << buf << endl;
close(fd);
return 0;
}
3.3 原始版本的具体问题
问题一:每轮复制并修改 fd_set
read_fds = master_fds;
select(max_fd + 1, &read_fds, nullptr, nullptr, nullptr);
select() 会原地修改 read_fds,所以必须保留 master_fds,并在每轮重新复制。
问题二:返回后仍然遍历全部客户端
即使只有一个客户端有数据,应用仍然会检查 clients 中的每一个 fd:
for (auto it = clients.begin(); it != clients.end(); ++it) {
if (FD_ISSET(*it, &read_fds)) {
// …
}
}
问题三:受到 fd_set 表示能力限制
在常见 Linux/glibc 环境中,fd_set 的容量由 FD_SETSIZE 约束。文件描述符数值一旦超出可表示范围,继续使用 FD_SET() 会导致未定义行为。
问题四:把 EOF 和错误混在一起
if (n <= 0) {
// …
}
这里没有区分:
n == 0 对端有序关闭写方向,接收端到达 EOF
n == -1 系统调用失败,需要检查 errno
问题五:默认阻塞 I/O
即使 select() 报告套接字可读,后续操作仍可能因为竞态、错误的循环方式或不完整的状态判断而阻塞。
问题六:没有处理部分发送
write(cfd, reply.data(), reply.size());
TCP write()/send() 成功返回时,返回值可能小于请求发送长度。不能把“一次调用成功”理解为“全部数据已经发送”。
问题七:TCP 没有消息边界
一次 write() 不保证对应对端的一次 read():
发送端:
write("ABC")
write("DEF")
接收端可能看到:
read() -> "ABCDEF"
也可能看到:
read() -> "AB"
read() -> "CDE"
read() -> "F"
第二部分:第 1 层——最小 epoll LT 迁移
四、本层目标
这一层只完成一件事:
用 epoll_create1()、epoll_ctl()、epoll_wait() 替换 fd_set 和 select()。
为了让差异足够清楚,本层刻意保留:
- 阻塞套接字;
- 每次只 accept() 一个连接;
- 每次只 read() 一次;
- 单次 write();
- 没有发送队列。
4.1 从 select 到 epoll 的最小映射
| FD_ZERO() | epoll_create1() |
| FD_SET(fd) | epoll_ctl(ADD, fd) |
| 修改监控事件 | epoll_ctl(MOD, fd) |
| FD_CLR(fd) | epoll_ctl(DEL, fd) |
| select() | epoll_wait() |
| 遍历全部客户端 | 只遍历返回的就绪事件 |
核心循环从:
read_fds = master_fds;
select(...);
for (int fd : clients) {
if (FD_ISSET(fd, &read_fds)) {
// …
}
}
变成:
int count = epoll_wait(...);
for (int i = 0; i < count; ++i) {
int fd = events[i].data.fd;
// fd 已经是本轮就绪对象
}
4.2 完整服务端:01_epoll_lt_minimal/server.cpp
#include <arpa/inet.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <cstring>
#include <iostream>
#include <string>
#include <unordered_set>
#include <vector>
namespace {
constexpr uint16_t kPort = 8888;
constexpr int kBacklog = 5;
constexpr int kMaxEvents = 64;
constexpr std::size_t kBufferSize = 64;
bool AddReadEvent(int epoll_fd, int fd) {
epoll_event event{};
event.events = EPOLLIN; // 未指定 EPOLLET,默认是 LT
event.data.fd = fd;
if (epoll_ctl(epoll_fd, EPOLL_CTL_ADD, fd, &event) == –1) {
std::perror("epoll_ctl(ADD)");
return false;
}
return true;
}
void CloseClient(int epoll_fd,
int client_fd,
std::unordered_set<int>& clients) {
// 对于单线程简单程序,close() 最终也会让内核移除监控项;
// 显式 DEL 可以让用户态生命周期更加清晰。
if (epoll_ctl(epoll_fd,
EPOLL_CTL_DEL,
client_fd,
nullptr) == –1 &&
errno != ENOENT &&
errno != EBADF) {
std::perror("epoll_ctl(DEL)");
}
close(client_fd);
clients.erase(client_fd);
}
}
int main() {
int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
if (listen_fd == –1) {
std::perror("socket");
return 1;
}
int reuse = 1;
if (setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse)) == –1) {
std::perror("setsockopt");
close(listen_fd);
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
server_addr.sin_port = htons(kPort);
if (bind(listen_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("bind");
close(listen_fd);
return 1;
}
if (listen(listen_fd, kBacklog) == –1) {
std::perror("listen");
close(listen_fd);
return 1;
}
// epoll 实例本身也是一个文件描述符。
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
if (epoll_fd == –1) {
std::perror("epoll_create1");
close(listen_fd);
return 1;
}
if (!AddReadEvent(epoll_fd, listen_fd)) {
close(epoll_fd);
close(listen_fd);
return 1;
}
std::vector<epoll_event> ready_events(kMaxEvents);
std::unordered_set<int> clients;
std::cout << "epoll LT minimal server listening on 0.0.0.0:"
<< kPort << '\\n';
while (true) {
int ready_count =
epoll_wait(epoll_fd,
ready_events.data(),
static_cast<int>(ready_events.size()),
–1);
if (ready_count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
// epoll_wait 只返回本轮已经就绪的事件,不再遍历全部连接。
for (int i = 0; i < ready_count; ++i) {
const int current_fd = ready_events[i].data.fd;
const uint32_t events = ready_events[i].events;
if (current_fd == listen_fd) {
// 为了保持“最小迁移”,本层仍然只 accept 一次。
sockaddr_in client_addr{};
socklen_t client_len = sizeof(client_addr);
int client_fd =
accept(listen_fd,
reinterpret_cast<sockaddr*>(&client_addr),
&client_len);
if (client_fd == –1) {
std::perror("accept");
continue;
}
char ip[INET_ADDRSTRLEN]{};
inet_ntop(AF_INET,
&client_addr.sin_addr,
ip,
sizeof(ip));
std::cout << "connected: "
<< ip << ':'
<< ntohs(client_addr.sin_port)
<< ", fd=" << client_fd << '\\n';
if (!AddReadEvent(epoll_fd, client_fd)) {
close(client_fd);
continue;
}
clients.insert(client_fd);
continue;
}
bool close_now = false;
if ((events & EPOLLIN) != 0U) {
char buffer[kBufferSize];
ssize_t n = read(current_fd, buffer, sizeof(buffer));
if (n > 0) {
std::cout << "received from fd="
<< current_fd << ": ";
std::cout.write(buffer, n);
std::cout << '\\n';
// 本层刻意保留“单次 write”的旧结构。
ssize_t sent =
write(current_fd,
buffer,
static_cast<std::size_t>(n));
if (sent == –1) {
std::perror("write");
close_now = true;
}
} else if (n == 0) {
close_now = true;
} else {
std::perror("read");
close_now = true;
}
}
if ((events & (EPOLLERR | EPOLLHUP)) != 0U) {
close_now = true;
}
if (close_now) {
std::cout << "closed fd=" << current_fd << '\\n';
CloseClient(epoll_fd, current_fd, clients);
}
}
}
for (int client_fd : clients) {
close(client_fd);
}
close(epoll_fd);
close(listen_fd);
return 0;
}
4.3 完整客户端:01_epoll_lt_minimal/client.cpp
#include <arpa/inet.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <cstring>
#include <iostream>
#include <string>
namespace {
bool SendAll(int fd, const char* data, std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
if (socket_fd == –1) {
std::perror("socket");
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(8888);
if (inet_pton(AF_INET,
"127.0.0.1",
&server_addr.sin_addr) != 1) {
std::cerr << "invalid server address\\n";
close(socket_fd);
return 1;
}
if (connect(socket_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("connect");
close(socket_fd);
return 1;
}
const std::string message = "Hello from epoll LT client";
if (!SendAll(socket_fd, message.data(), message.size())) {
close(socket_fd);
return 1;
}
std::cout << "sent: " << message << '\\n';
char buffer[256];
ssize_t n;
do {
n = recv(socket_fd, buffer, sizeof(buffer), 0);
} while (n == –1 && errno == EINTR);
if (n > 0) {
std::cout << "received: ";
std::cout.write(buffer, n);
std::cout << '\\n';
} else if (n == 0) {
std::cout << "server closed connection\\n";
} else {
std::perror("recv");
}
close(socket_fd);
return n < 0 ? 1 : 0;
}
4.4 本层修复了什么
修复一:不再每轮复制完整监控集合
监控项通过 epoll_ctl() 加入 epoll 实例,之后由内核维护。
修复二:不再遍历所有空闲客户端
epoll_wait() 直接返回本轮就绪事件数组:
for (int i = 0; i < ready_count; ++i) {
int current_fd = ready_events[i].data.fd;
}
修复三:不再依赖 max_fd + 1
epoll_wait() 不需要最大 fd 参数。
4.5 本层仍然存在的问题
- 监听套接字仍然阻塞;
- 每次监听事件只 accept() 一个连接;
- 通信套接字仍然阻塞;
- 没有正确区分 EINTR、EAGAIN;
- 单次 read() 可能无法读完当前数据;
- 单次 write() 可能部分发送;
- 没有处理慢客户端;
- 没有背压。
因此本层只是“API 迁移版”,不是最终版本。
第三部分:第 2 层——非阻塞 epoll LT
五、为什么 LT 也建议使用非阻塞套接字
LT 模式的语义和 poll() 接近:只要条件仍然满足,就可能继续报告就绪。
这并不意味着阻塞套接字一定安全。
例如:
epoll_wait 报告 fd 可读
│
▼
另一个执行单元先读走数据
│
▼
当前线程调用阻塞 recv()
│
▼
可能睡眠等待新数据
在本文的单线程示例里这种竞态较少,但统一采用非阻塞 I/O 可以建立更稳定的事件循环不变量:
任何可能等待的系统调用,都必须以 EAGAIN/EWOULDBLOCK 返回控制权,而不是阻塞整个事件循环。
5.1 本层新增的规则
规则一:监听 fd 非阻塞
这样才能安全地循环 accept():
while (true) {
int cfd = accept(...);
if (cfd >= 0) {
// 继续 accept
} else if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
}
规则二:通信 fd 也设为非阻塞
Linux 上,accept() 得到的新套接字不会因为监听套接字设置了 O_NONBLOCK 就自动继承该状态,因此必须单独设置,或者使用带 SOCK_NONBLOCK 的 accept4()。
规则三:循环读取到 EAGAIN
LT 并不强制必须这样做,但排空当前接收缓冲区可以减少重复进入 epoll_wait()。
规则四:本层遇到发送缓冲区满时直接关闭连接
这是一个有意保留的中间策略:
好处:不会阻塞事件循环,也不会悄悄丢弃未发送数据
代价:慢客户端会被断开
真正保留慢客户端连接,需要在第 4 层引入发送队列。
5.2 完整服务端:02_epoll_lt_nonblocking/server.cpp
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
#include <unordered_set>
#include <vector>
namespace {
constexpr uint16_t kPort = 8888;
constexpr int kBacklog = 128;
constexpr int kMaxEvents = 64;
constexpr std::size_t kBufferSize = 4096;
bool SetNonBlocking(int fd) {
int old_flags = fcntl(fd, F_GETFL, 0);
if (old_flags == –1) {
return false;
}
return fcntl(fd,
F_SETFL,
old_flags | O_NONBLOCK) != –1;
}
bool AddClient(int epoll_fd, int client_fd) {
epoll_event event{};
event.events = EPOLLIN | EPOLLRDHUP; // 默认 LT
event.data.fd = client_fd;
return epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&event) != –1;
}
void CloseClient(int epoll_fd,
int client_fd,
std::unordered_set<int>& clients) {
if (epoll_ctl(epoll_fd,
EPOLL_CTL_DEL,
client_fd,
nullptr) == –1 &&
errno != ENOENT &&
errno != EBADF) {
std::perror("epoll_ctl(DEL)");
}
close(client_fd);
clients.erase(client_fd);
std::cout << "closed fd=" << client_fd << '\\n';
}
// 本层尚未引入发送队列。
// 为避免非阻塞 send() 阻塞事件循环,如果不能一次发送完,直接关闭连接。
bool SendImmediatelyOrFail(int fd,
const char* data,
std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
if (n == –1 &&
(errno == EAGAIN || errno == EWOULDBLOCK)) {
std::cerr
<< "send buffer full; layer 02 has no output queue\\n";
return false;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int listen_fd = socket(AF_INET,
SOCK_STREAM | SOCK_CLOEXEC,
0);
if (listen_fd == –1) {
std::perror("socket");
return 1;
}
if (!SetNonBlocking(listen_fd)) {
std::perror("fcntl(listen_fd)");
close(listen_fd);
return 1;
}
int reuse = 1;
if (setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse)) == –1) {
std::perror("setsockopt");
close(listen_fd);
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
server_addr.sin_port = htons(kPort);
if (bind(listen_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("bind");
close(listen_fd);
return 1;
}
if (listen(listen_fd, kBacklog) == –1) {
std::perror("listen");
close(listen_fd);
return 1;
}
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
if (epoll_fd == –1) {
std::perror("epoll_create1");
close(listen_fd);
return 1;
}
epoll_event listen_event{};
listen_event.events = EPOLLIN; // 默认 LT
listen_event.data.fd = listen_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
listen_fd,
&listen_event) == –1) {
std::perror("epoll_ctl(ADD listen_fd)");
close(epoll_fd);
close(listen_fd);
return 1;
}
std::vector<epoll_event> ready_events(kMaxEvents);
std::unordered_set<int> clients;
std::cout << "epoll LT nonblocking server listening on 0.0.0.0:"
<< kPort << '\\n';
while (true) {
int ready_count =
epoll_wait(epoll_fd,
ready_events.data(),
static_cast<int>(ready_events.size()),
–1);
if (ready_count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
for (int i = 0; i < ready_count; ++i) {
int current_fd = ready_events[i].data.fd;
uint32_t events = ready_events[i].events;
if (current_fd == listen_fd) {
// 监听 fd 已经非阻塞,因此可以一直 accept 到 EAGAIN。
while (true) {
sockaddr_in client_addr{};
socklen_t client_len =
sizeof(client_addr);
int client_fd =
accept(listen_fd,
reinterpret_cast<sockaddr*>(
&client_addr),
&client_len);
if (client_fd == –1) {
if (errno == EINTR ||
errno == ECONNABORTED) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("accept");
break;
}
if (!SetNonBlocking(client_fd)) {
std::perror("fcntl(client_fd)");
close(client_fd);
continue;
}
if (!AddClient(epoll_fd, client_fd)) {
std::perror("epoll_ctl(ADD client)");
close(client_fd);
continue;
}
char ip[INET_ADDRSTRLEN]{};
inet_ntop(AF_INET,
&client_addr.sin_addr,
ip,
sizeof(ip));
std::cout << "connected: "
<< ip << ':'
<< ntohs(client_addr.sin_port)
<< ", fd=" << client_fd << '\\n';
clients.insert(client_fd);
}
continue;
}
bool close_now = false;
if ((events & EPOLLERR) != 0U) {
int socket_error = 0;
socklen_t error_len = sizeof(socket_error);
getsockopt(current_fd,
SOL_SOCKET,
SO_ERROR,
&socket_error,
&error_len);
std::cerr << "socket error on fd="
<< current_fd
<< ": errno=" << socket_error
<< '\\n';
close_now = true;
}
if (!close_now &&
(events & EPOLLIN) != 0U) {
// 即使 LT 不强制排空,非阻塞循环读取可以减少重复唤醒。
while (true) {
char buffer[kBufferSize];
ssize_t n = recv(current_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
std::cout << "received "
<< n << " bytes from fd="
<< current_fd << '\\n';
if (!SendImmediatelyOrFail(
current_fd,
buffer,
static_cast<std::size_t>(n))) {
close_now = true;
break;
}
continue;
}
if (n == 0) {
close_now = true;
break;
}
if (errno == EINTR) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("recv");
close_now = true;
break;
}
}
if ((events &
(EPOLLRDHUP | EPOLLHUP)) != 0U) {
close_now = true;
}
if (close_now) {
CloseClient(epoll_fd,
current_fd,
clients);
}
}
}
for (int client_fd : clients) {
close(client_fd);
}
close(epoll_fd);
close(listen_fd);
return 0;
}
5.3 完整客户端:02_epoll_lt_nonblocking/client.cpp
客户端仍然使用 select() 同时监控:
- 标准输入;
- 服务器套接字。
这样做是为了把本章的变量限制在服务端 epoll 重构上,而不是同时改变客户端模型。
#include <arpa/inet.h>
#include <sys/select.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
namespace {
bool SendAll(int fd, const char* data, std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
if (socket_fd == –1) {
std::perror("socket");
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(8888);
if (inet_pton(AF_INET,
"127.0.0.1",
&server_addr.sin_addr) != 1) {
std::cerr << "invalid IPv4 address\\n";
close(socket_fd);
return 1;
}
if (connect(socket_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("connect");
close(socket_fd);
return 1;
}
std::cout << "connected to 127.0.0.1:8888\\n"
<< "type text; Ctrl+D stops sending but keeps receiving\\n";
bool stdin_open = true;
while (true) {
fd_set read_set;
FD_ZERO(&read_set);
if (stdin_open) {
FD_SET(STDIN_FILENO, &read_set);
}
FD_SET(socket_fd, &read_set);
int max_fd = socket_fd > STDIN_FILENO
? socket_fd
: STDIN_FILENO;
int ready =
select(max_fd + 1,
&read_set,
nullptr,
nullptr,
nullptr);
if (ready == –1) {
if (errno == EINTR) {
continue;
}
std::perror("select");
break;
}
if (stdin_open &&
FD_ISSET(STDIN_FILENO, &read_set)) {
char buffer[4096];
ssize_t n = read(STDIN_FILENO,
buffer,
sizeof(buffer));
if (n > 0) {
if (!SendAll(socket_fd,
buffer,
static_cast<std::size_t>(n))) {
close(socket_fd);
return 1;
}
} else if (n == 0) {
stdin_open = false;
if (shutdown(socket_fd, SHUT_WR) == –1 &&
errno != ENOTCONN) {
std::perror("shutdown");
}
} else if (errno != EINTR) {
std::perror("read(stdin)");
break;
}
}
if (FD_ISSET(socket_fd, &read_set)) {
char buffer[4096];
ssize_t n = recv(socket_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
std::cout.write(buffer, n);
std::cout.flush();
} else if (n == 0) {
std::cout << "\\nserver closed connection\\n";
close(socket_fd);
return 0;
} else if (errno != EINTR) {
std::perror("recv");
break;
}
}
}
close(socket_fd);
return 1;
}
5.4 本层修复了什么
- 监听套接字不再因为多调用一次 accept() 而阻塞;
- 一次监听事件可以排空当前已完成连接队列;
- 通信套接字不再让事件循环长期阻塞;
- 正确处理 EINTR;
- 正确处理 EAGAIN/EWOULDBLOCK;
- 使用 EPOLLRDHUP 感知对端关闭写方向;
- 使用 MSG_NOSIGNAL 避免向已关闭连接发送数据时触发 SIGPIPE 终止进程;
- 读取路径可以一次处理当前已经到达的多批数据。
5.5 为什么不能把 EAGAIN 当成错误
非阻塞套接字中:
recv() == -1
errno == EAGAIN / EWOULDBLOCK
表示:
当前没有更多数据可以立即读取,本轮处理正常结束。
同理,非阻塞 send() 返回 EAGAIN 表示:
当前发送缓冲区暂时没有足够空间,需要稍后等待可写事件。
它们是事件驱动程序中的正常控制流,而不是异常崩溃条件。
5.6 本层仍然存在的问题
本层的读路径已经较稳健,但写路径仍然采用:
立即写完
│
├─ 成功:继续
└─ 写不完/EAGAIN:关闭连接
下一步切换 ET 时,先保持这个限制不变,只专注解决 ET 读取规则。
第四部分:LT 与 ET 的真正区别
六、水平触发 LT
LT 可以理解为:
只要 fd 当前仍然满足就绪条件,就可能继续通知。
假设接收缓冲区有 5000 字节,而应用每次只读 1024 字节:
第 1 次 epoll_wait -> EPOLLIN
read 1024,剩余 3976
第 2 次 epoll_wait -> 仍可能 EPOLLIN
read 1024,剩余 2952
……
七、边缘触发 ET
ET 更接近:
当监控对象的就绪状态发生变化时通知;收到通知后,应用应把该 fd 当作可操作对象处理,直到非阻塞 I/O 返回 EAGAIN。
注意,ET 不能被机械地简化为严格的“缓冲区从 0 字节变成非 0 字节只通知一次”。多个数据到达过程仍可能产生多个事件。真正可靠的编程契约是:
收到 ET 事件
│
▼
认为 fd 当前可以继续处理
│
▼
非阻塞 read/write 循环
│
└─ 直到 EAGAIN 才把控制权交回 epoll_wait
第五部分:ET 的两个典型错误
八、错误一:EPOLLET + 单次 read
8.1 错误逻辑
event.events = EPOLLIN | EPOLLET;
// 收到事件后只读取一次
read(fd, buffer, 1024);
假设客户端发送 5000 字节:
接收缓冲区
┌────────────────────────────────────────────┐
│ 前 1024 字节 │ 剩余 3976 字节 │
└────────────────────────────────────────────┘
▲
└─ 单次 read 只取走这里
剩余数据并没有被内核“丢掉”,而是仍留在接收缓冲区中。
问题在于:
应用已经消费了这一轮边缘通知,却没有把当前可读状态处理完;之后可能长时间收不到新的通知。
所以更准确的术语是:
- 未读数据滞留;
- 读取停滞;
- 事件通知遗漏后的应用层饥饿。
不应简单称为“TCP 数据丢失”。
8.2 完整错误复现代码
该代码只用于实验,不应作为正确实现使用。
// 错误演示:EPOLLET + 非阻塞套接字 + 单次 read。
// 结果不是“内核丢包”,而是未读数据可能长期留在接收缓冲区,
// 事件循环却不一定再次收到通知。
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
bool SetNonBlocking(int fd) {
int flags = fcntl(fd, F_GETFL, 0);
return flags != –1 &&
fcntl(fd, F_SETFL, flags | O_NONBLOCK) != –1;
}
int main() {
int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
int reuse = 1;
setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse));
SetNonBlocking(listen_fd);
sockaddr_in addr{};
addr.sin_family = AF_INET;
addr.sin_addr.s_addr = htonl(INADDR_ANY);
addr.sin_port = htons(8888);
bind(listen_fd,
reinterpret_cast<sockaddr*>(&addr),
sizeof(addr));
listen(listen_fd, 128);
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
epoll_event listen_event{};
listen_event.events = EPOLLIN | EPOLLET;
listen_event.data.fd = listen_fd;
epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
listen_fd,
&listen_event);
epoll_event events[16]{};
while (true) {
int count = epoll_wait(epoll_fd,
events,
16,
–1);
if (count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
for (int i = 0; i < count; ++i) {
int fd = events[i].data.fd;
if (fd == listen_fd) {
while (true) {
int client_fd =
accept(listen_fd,
nullptr,
nullptr);
if (client_fd == –1) {
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
break;
}
SetNonBlocking(client_fd);
epoll_event client_event{};
client_event.events =
EPOLLIN | EPOLLET;
client_event.data.fd = client_fd;
epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&client_event);
}
} else {
char buffer[1024];
// 错误点:只读取一次,没有排空到 EAGAIN。
ssize_t n = read(fd,
buffer,
sizeof(buffer));
std::cout << "single read returned "
<< n << " bytes\\n";
}
}
}
close(epoll_fd);
close(listen_fd);
return 0;
}
九、错误二:EPOLLET + 阻塞套接字 + 循环 read
有人发现单次读取不够,于是改成:
while (true) {
read(fd, buffer, sizeof(buffer));
}
但如果 fd 仍是阻塞模式:
第 1 次 read -> 1024 字节
第 2 次 read -> 1024 字节
……
最后一批 read -> 成功
下一次 read -> 接收缓冲区为空
│
▼
当前线程睡眠
│
▼
整个事件循环停止处理其他 fd
这通常不是严格意义上的“死锁”,因为新数据到来后 read() 仍可能返回。
更准确的描述是:
单个阻塞套接字把整个单线程事件循环长期卡住,其他连接发生饥饿。
9.1 完整错误复现代码
该代码只用于实验,不应作为正确实现使用。
// 错误演示:EPOLLET + 阻塞通信套接字 + while(read)。
// 当接收缓冲区被读空后,下一次 read 会阻塞整个事件循环。
// 这不是严格意义上的“死锁”,而是事件循环被单个连接长期阻塞。
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
bool SetNonBlocking(int fd) {
int flags = fcntl(fd, F_GETFL, 0);
return flags != –1 &&
fcntl(fd, F_SETFL, flags | O_NONBLOCK) != –1;
}
int main() {
int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
int reuse = 1;
setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse));
// 监听套接字非阻塞,方便循环 accept。
SetNonBlocking(listen_fd);
sockaddr_in addr{};
addr.sin_family = AF_INET;
addr.sin_addr.s_addr = htonl(INADDR_ANY);
addr.sin_port = htons(8888);
bind(listen_fd,
reinterpret_cast<sockaddr*>(&addr),
sizeof(addr));
listen(listen_fd, 128);
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
epoll_event listen_event{};
listen_event.events = EPOLLIN | EPOLLET;
listen_event.data.fd = listen_fd;
epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
listen_fd,
&listen_event);
epoll_event events[16]{};
while (true) {
int count = epoll_wait(epoll_fd,
events,
16,
–1);
if (count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
for (int i = 0; i < count; ++i) {
int fd = events[i].data.fd;
if (fd == listen_fd) {
while (true) {
int client_fd =
accept(listen_fd,
nullptr,
nullptr);
if (client_fd == –1) {
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
break;
}
// 错误点:没有把 client_fd 设置为非阻塞。
epoll_event client_event{};
client_event.events =
EPOLLIN | EPOLLET;
client_event.data.fd = client_fd;
epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&client_event);
}
} else {
char buffer[1024];
while (true) {
// 数据读空后,这里会睡眠等待下一批数据,
// 事件循环无法处理其他 fd。
ssize_t n = read(fd,
buffer,
sizeof(buffer));
if (n <= 0) {
break;
}
std::cout << "read "
<< n << " bytes\\n";
}
}
}
}
close(epoll_fd);
close(listen_fd);
return 0;
}
十、ET 正确读法
ET 读取路径应保持三个不变量:
1. fd 必须是非阻塞模式
2. 收到可读事件后循环 recv()
3. 只有遇到 EAGAIN/EWOULDBLOCK 才结束本轮读取
状态表:
| n > 0 | 读到 n 字节 | 处理数据,继续读 |
| n == 0 | 接收方向到达 EOF | 标记对端写方向关闭 |
| n == -1, EINTR | 被信号中断 | 重试 |
| n == -1, EAGAIN/EWOULDBLOCK | 当前已经读空 | 正常结束本轮 |
| 其他错误 | 连接异常 | 关闭连接 |
第六部分:第 3 层——正确的 epoll ET 读取
十一、本层相对第 2 层改了什么
代码结构几乎不变,关键差异只有事件掩码:
// LT
EPOLLIN | EPOLLRDHUP
// ET
EPOLLIN | EPOLLRDHUP | EPOLLET
因为第 2 层已经完成:
- 非阻塞监听 fd;
- 非阻塞通信 fd;
- 循环 accept() 到 EAGAIN;
- 循环 recv() 到 EAGAIN;
所以切换 ET 不需要一次性重写整个事件循环。
这正是分层升级的价值。
11.1 完整服务端:03_epoll_et_correct_read/server.cpp
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
#include <unordered_set>
#include <vector>
namespace {
constexpr uint16_t kPort = 8888;
constexpr int kBacklog = 128;
constexpr int kMaxEvents = 64;
constexpr std::size_t kBufferSize = 4096;
bool SetNonBlocking(int fd) {
int old_flags = fcntl(fd, F_GETFL, 0);
if (old_flags == –1) {
return false;
}
return fcntl(fd,
F_SETFL,
old_flags | O_NONBLOCK) != –1;
}
bool AddClient(int epoll_fd, int client_fd) {
epoll_event event{};
event.events =
EPOLLIN | EPOLLRDHUP | EPOLLET;
event.data.fd = client_fd;
return epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&event) != –1;
}
void CloseClient(int epoll_fd,
int client_fd,
std::unordered_set<int>& clients) {
if (epoll_ctl(epoll_fd,
EPOLL_CTL_DEL,
client_fd,
nullptr) == –1 &&
errno != ENOENT &&
errno != EBADF) {
std::perror("epoll_ctl(DEL)");
}
close(client_fd);
clients.erase(client_fd);
std::cout << "closed fd=" << client_fd << '\\n';
}
// 本层重点是把 ET 的“读路径”做正确。
// 写路径仍采用“必须立即全部写完,否则关闭连接”的简化策略。
bool SendImmediatelyOrFail(int fd,
const char* data,
std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
if (n == –1 &&
(errno == EAGAIN || errno == EWOULDBLOCK)) {
std::cerr
<< "send buffer full; layer 03 has no EPOLLOUT queue\\n";
return false;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int listen_fd = socket(AF_INET,
SOCK_STREAM | SOCK_CLOEXEC,
0);
if (listen_fd == –1) {
std::perror("socket");
return 1;
}
if (!SetNonBlocking(listen_fd)) {
std::perror("fcntl(listen_fd)");
close(listen_fd);
return 1;
}
int reuse = 1;
if (setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse)) == –1) {
std::perror("setsockopt");
close(listen_fd);
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
server_addr.sin_port = htons(kPort);
if (bind(listen_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("bind");
close(listen_fd);
return 1;
}
if (listen(listen_fd, kBacklog) == –1) {
std::perror("listen");
close(listen_fd);
return 1;
}
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
if (epoll_fd == –1) {
std::perror("epoll_create1");
close(listen_fd);
return 1;
}
epoll_event listen_event{};
listen_event.events = EPOLLIN | EPOLLET;
listen_event.data.fd = listen_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
listen_fd,
&listen_event) == –1) {
std::perror("epoll_ctl(ADD listen_fd)");
close(epoll_fd);
close(listen_fd);
return 1;
}
std::vector<epoll_event> ready_events(kMaxEvents);
std::unordered_set<int> clients;
std::cout << "epoll ET read-correct server listening on 0.0.0.0:"
<< kPort << '\\n';
while (true) {
int ready_count =
epoll_wait(epoll_fd,
ready_events.data(),
static_cast<int>(ready_events.size()),
–1);
if (ready_count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
for (int i = 0; i < ready_count; ++i) {
int current_fd = ready_events[i].data.fd;
uint32_t events = ready_events[i].events;
if (current_fd == listen_fd) {
// ET:必须一直 accept 到 EAGAIN。
while (true) {
sockaddr_in client_addr{};
socklen_t client_len =
sizeof(client_addr);
int client_fd =
accept(listen_fd,
reinterpret_cast<sockaddr*>(
&client_addr),
&client_len);
if (client_fd == –1) {
if (errno == EINTR ||
errno == ECONNABORTED) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("accept");
break;
}
if (!SetNonBlocking(client_fd)) {
std::perror("fcntl(client_fd)");
close(client_fd);
continue;
}
if (!AddClient(epoll_fd, client_fd)) {
std::perror("epoll_ctl(ADD client)");
close(client_fd);
continue;
}
char ip[INET_ADDRSTRLEN]{};
inet_ntop(AF_INET,
&client_addr.sin_addr,
ip,
sizeof(ip));
std::cout << "connected: "
<< ip << ':'
<< ntohs(client_addr.sin_port)
<< ", fd=" << client_fd << '\\n';
clients.insert(client_fd);
}
continue;
}
bool close_now = false;
if ((events & EPOLLERR) != 0U) {
int socket_error = 0;
socklen_t error_len = sizeof(socket_error);
getsockopt(current_fd,
SOL_SOCKET,
SO_ERROR,
&socket_error,
&error_len);
std::cerr << "socket error on fd="
<< current_fd
<< ": errno=" << socket_error
<< '\\n';
close_now = true;
}
if (!close_now &&
(events & EPOLLIN) != 0U) {
// ET 核心:非阻塞 + 循环 recv,直到 EAGAIN。
while (true) {
char buffer[kBufferSize];
ssize_t n = recv(current_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
std::cout << "received "
<< n << " bytes from fd="
<< current_fd << '\\n';
if (!SendImmediatelyOrFail(
current_fd,
buffer,
static_cast<std::size_t>(n))) {
close_now = true;
break;
}
continue;
}
if (n == 0) {
close_now = true;
break;
}
if (errno == EINTR) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
// 当前接收缓冲区已经排空。
break;
}
std::perror("recv");
close_now = true;
break;
}
}
if ((events &
(EPOLLRDHUP | EPOLLHUP)) != 0U) {
close_now = true;
}
if (close_now) {
CloseClient(epoll_fd,
current_fd,
clients);
}
}
}
for (int client_fd : clients) {
close(client_fd);
}
close(epoll_fd);
close(listen_fd);
return 0;
}
11.2 完整客户端:03_epoll_et_correct_read/client.cpp
#include <arpa/inet.h>
#include <sys/select.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
namespace {
bool SendAll(int fd, const char* data, std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
if (socket_fd == –1) {
std::perror("socket");
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(8888);
if (inet_pton(AF_INET,
"127.0.0.1",
&server_addr.sin_addr) != 1) {
std::cerr << "invalid IPv4 address\\n";
close(socket_fd);
return 1;
}
if (connect(socket_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("connect");
close(socket_fd);
return 1;
}
std::cout << "connected to 127.0.0.1:8888\\n"
<< "type text; Ctrl+D stops sending but keeps receiving\\n";
bool stdin_open = true;
while (true) {
fd_set read_set;
FD_ZERO(&read_set);
if (stdin_open) {
FD_SET(STDIN_FILENO, &read_set);
}
FD_SET(socket_fd, &read_set);
int max_fd = socket_fd > STDIN_FILENO
? socket_fd
: STDIN_FILENO;
int ready =
select(max_fd + 1,
&read_set,
nullptr,
nullptr,
nullptr);
if (ready == –1) {
if (errno == EINTR) {
continue;
}
std::perror("select");
break;
}
if (stdin_open &&
FD_ISSET(STDIN_FILENO, &read_set)) {
char buffer[4096];
ssize_t n = read(STDIN_FILENO,
buffer,
sizeof(buffer));
if (n > 0) {
if (!SendAll(socket_fd,
buffer,
static_cast<std::size_t>(n))) {
close(socket_fd);
return 1;
}
} else if (n == 0) {
stdin_open = false;
if (shutdown(socket_fd, SHUT_WR) == –1 &&
errno != ENOTCONN) {
std::perror("shutdown");
}
} else if (errno != EINTR) {
std::perror("read(stdin)");
break;
}
}
if (FD_ISSET(socket_fd, &read_set)) {
char buffer[4096];
ssize_t n = recv(socket_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
std::cout.write(buffer, n);
std::cout.flush();
} else if (n == 0) {
std::cout << "\\nserver closed connection\\n";
close(socket_fd);
return 0;
} else if (errno != EINTR) {
std::perror("recv");
break;
}
}
}
close(socket_fd);
return 1;
}
11.3 本层修复了什么
- 监听 fd 在 ET 模式下排空连接队列;
- 通信 fd 在 ET 模式下排空接收缓冲区;
- 数据不会因为只读一次而长期滞留;
- 事件循环不会因为阻塞 recv() 卡住;
- EAGAIN 成为“本轮处理完成”的边界。
11.4 本层仍然没有解决写路径
当前逻辑仍然要求:
收到一批数据
│
▼
立即 send
│
├─ 全部发送:正常
└─ 部分发送/EAGAIN:关闭连接
这不会悄悄丢数据,但会错误地拒绝慢客户端。
要正确处理它,必须理解 EPOLLOUT。
第七部分:为什么不能永久监听 EPOLLOUT
十二、套接字“可写”意味着什么
对 TCP 套接字而言,EPOLLOUT 通常表示:
发送缓冲区当前有空间,执行非阻塞发送可能取得进展。
新建立的 TCP 连接大多数时间都是可写的。
如果无条件永久注册:
EPOLLIN | EPOLLOUT
事件循环可能不断得到 EPOLLOUT:
epoll_wait -> fd 可写
没有数据需要发送
再次 epoll_wait -> fd 仍可写
没有数据需要发送
……
这会形成无意义的高频唤醒。
因此常见策略是:
没有待发送数据:
只监听 EPOLLIN
出现待发送数据:
先主动 send
全部发送完成:
不监听 EPOLLOUT
send 返回 EAGAIN 或只完成一部分:
保存剩余数据
注册 EPOLLOUT
后续 EPOLLOUT 到来:
继续发送
十三、为什么每个连接都需要独立发送状态
TCP 发送是每连接独立的,因此需要:
struct ClientState {
std::string output;
std::size_t sent_offset;
};
假设:
output = "ABCDEFGHIJ"
sent_offset = 4
表示:
已经发送:ABCD
等待发送:EFGHIJ
下一次发送地址为:
output.data() + sent_offset
剩余长度为:
output.size() – sent_offset
第八部分:第 4 层——ET + 发送队列 + EPOLLOUT
十四、本层新增的状态机
收到客户端数据
│
▼
追加 output
│
▼
主动尝试 send
│ │
│ └─ EAGAIN / 部分发送
│ │
▼ ▼
全部发完 保留 sent_offset
│ │
移除 EPOLLOUT 注册 EPOLLOUT
│
▼
下次可写时继续
14.1 完整服务端:04_epoll_et_write_queue/server.cpp
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
#include <string>
#include <unordered_map>
#include <vector>
namespace {
constexpr uint16_t kPort = 8888;
constexpr int kBacklog = 128;
constexpr int kMaxEvents = 64;
constexpr std::size_t kBufferSize = 4096;
struct ClientState {
std::string peer;
std::string output;
std::size_t sent_offset = 0;
bool peer_closed_write = false;
};
bool SetNonBlocking(int fd) {
int old_flags = fcntl(fd, F_GETFL, 0);
if (old_flags == –1) {
return false;
}
return fcntl(fd,
F_SETFL,
old_flags | O_NONBLOCK) != –1;
}
std::size_t PendingBytes(const ClientState& client) {
return client.output.size() – client.sent_offset;
}
void CompactOutput(ClientState& client) {
if (client.sent_offset == 0) {
return;
}
if (client.sent_offset == client.output.size()) {
client.output.clear();
client.sent_offset = 0;
return;
}
if (client.sent_offset >= 4096 &&
client.sent_offset * 2 >= client.output.size()) {
client.output.erase(0, client.sent_offset);
client.sent_offset = 0;
}
}
uint32_t BuildClientEvents(const ClientState& client) {
uint32_t events = EPOLLET | EPOLLRDHUP;
if (!client.peer_closed_write) {
events |= EPOLLIN;
}
if (PendingBytes(client) > 0) {
events |= EPOLLOUT;
}
return events;
}
bool ModifyClientEvents(int epoll_fd,
int client_fd,
const ClientState& client) {
epoll_event event{};
event.events = BuildClientEvents(client);
event.data.fd = client_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_MOD,
client_fd,
&event) == –1) {
std::perror("epoll_ctl(MOD client)");
return false;
}
return true;
}
void CloseClient(
int epoll_fd,
int client_fd,
std::unordered_map<int, ClientState>& clients) {
if (epoll_ctl(epoll_fd,
EPOLL_CTL_DEL,
client_fd,
nullptr) == –1 &&
errno != ENOENT &&
errno != EBADF) {
std::perror("epoll_ctl(DEL)");
}
close(client_fd);
clients.erase(client_fd);
std::cout << "closed fd=" << client_fd << '\\n';
}
bool FlushOutput(int client_fd,
ClientState& client) {
while (client.sent_offset < client.output.size()) {
const char* data =
client.output.data() + client.sent_offset;
std::size_t remain =
client.output.size() – client.sent_offset;
ssize_t n =
send(client_fd,
data,
remain,
MSG_NOSIGNAL);
if (n > 0) {
client.sent_offset +=
static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
if (n == –1 &&
(errno == EAGAIN || errno == EWOULDBLOCK)) {
// 发送缓冲区暂时已满,保留未发送数据,
// 等待下一次 EPOLLOUT。
CompactOutput(client);
return true;
}
std::perror("send");
return false;
}
CompactOutput(client);
return true;
}
}
int main() {
int listen_fd = socket(AF_INET,
SOCK_STREAM | SOCK_CLOEXEC,
0);
if (listen_fd == –1) {
std::perror("socket");
return 1;
}
if (!SetNonBlocking(listen_fd)) {
std::perror("fcntl(listen_fd)");
close(listen_fd);
return 1;
}
int reuse = 1;
if (setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse)) == –1) {
std::perror("setsockopt");
close(listen_fd);
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
server_addr.sin_port = htons(kPort);
if (bind(listen_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("bind");
close(listen_fd);
return 1;
}
if (listen(listen_fd, kBacklog) == –1) {
std::perror("listen");
close(listen_fd);
return 1;
}
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
if (epoll_fd == –1) {
std::perror("epoll_create1");
close(listen_fd);
return 1;
}
epoll_event listen_event{};
listen_event.events = EPOLLIN | EPOLLET;
listen_event.data.fd = listen_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
listen_fd,
&listen_event) == –1) {
std::perror("epoll_ctl(ADD listen_fd)");
close(epoll_fd);
close(listen_fd);
return 1;
}
std::vector<epoll_event> ready_events(kMaxEvents);
std::unordered_map<int, ClientState> clients;
std::cout << "epoll ET output-queue server listening on 0.0.0.0:"
<< kPort << '\\n';
while (true) {
int ready_count =
epoll_wait(epoll_fd,
ready_events.data(),
static_cast<int>(ready_events.size()),
–1);
if (ready_count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
for (int i = 0; i < ready_count; ++i) {
int current_fd = ready_events[i].data.fd;
uint32_t events = ready_events[i].events;
if (current_fd == listen_fd) {
while (true) {
sockaddr_in client_addr{};
socklen_t client_len =
sizeof(client_addr);
int client_fd =
accept(listen_fd,
reinterpret_cast<sockaddr*>(
&client_addr),
&client_len);
if (client_fd == –1) {
if (errno == EINTR ||
errno == ECONNABORTED) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("accept");
break;
}
if (!SetNonBlocking(client_fd)) {
std::perror("fcntl(client_fd)");
close(client_fd);
continue;
}
char ip[INET_ADDRSTRLEN]{};
inet_ntop(AF_INET,
&client_addr.sin_addr,
ip,
sizeof(ip));
ClientState client;
client.peer =
std::string(ip) + ":" +
std::to_string(
ntohs(client_addr.sin_port));
epoll_event client_event{};
client_event.events =
EPOLLIN | EPOLLRDHUP | EPOLLET;
client_event.data.fd = client_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&client_event) == –1) {
std::perror(
"epoll_ctl(ADD client)");
close(client_fd);
continue;
}
std::cout << "connected: "
<< client.peer
<< ", fd=" << client_fd << '\\n';
clients.emplace(client_fd,
std::move(client));
}
continue;
}
auto it = clients.find(current_fd);
if (it == clients.end()) {
continue;
}
ClientState& client = it->second;
bool close_now = false;
if ((events & EPOLLERR) != 0U) {
int socket_error = 0;
socklen_t error_len = sizeof(socket_error);
getsockopt(current_fd,
SOL_SOCKET,
SO_ERROR,
&socket_error,
&error_len);
std::cerr << "socket error on fd="
<< current_fd
<< ": errno=" << socket_error
<< '\\n';
close_now = true;
}
if (!close_now &&
(events &
(EPOLLIN | EPOLLRDHUP | EPOLLHUP)) != 0U) {
while (true) {
char buffer[kBufferSize];
ssize_t n = recv(current_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
CompactOutput(client);
client.output.append(
buffer,
static_cast<std::size_t>(n));
std::cout << "queued "
<< n << " echo bytes for fd="
<< current_fd << '\\n';
continue;
}
if (n == 0) {
client.peer_closed_write = true;
break;
}
if (errno == EINTR) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("recv");
close_now = true;
break;
}
}
// EPOLLRDHUP/EPOLLHUP 只是“对端关闭方向发生变化”的提示。
// 仍要通过 recv() 返回 0 确认接收方向已经到达 EOF,
// 因为 FIN 之前可能还有尚未读取的数据。
// 无论本轮是否带 EPOLLOUT,只要存在待发送数据,
// 都先主动尝试一次发送;写不动再订阅 EPOLLOUT。
if (!close_now &&
PendingBytes(client) > 0) {
if (!FlushOutput(current_fd, client)) {
close_now = true;
}
}
if (!close_now &&
client.peer_closed_write &&
PendingBytes(client) == 0) {
close_now = true;
}
if (close_now) {
CloseClient(epoll_fd,
current_fd,
clients);
continue;
}
if (!ModifyClientEvents(epoll_fd,
current_fd,
client)) {
CloseClient(epoll_fd,
current_fd,
clients);
}
}
}
for (const auto& [client_fd, client] : clients) {
(void)client;
close(client_fd);
}
close(epoll_fd);
close(listen_fd);
return 0;
}
14.2 完整客户端:04_epoll_et_write_queue/client.cpp
#include <arpa/inet.h>
#include <sys/select.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
namespace {
bool SendAll(int fd, const char* data, std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
if (socket_fd == –1) {
std::perror("socket");
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(8888);
if (inet_pton(AF_INET,
"127.0.0.1",
&server_addr.sin_addr) != 1) {
std::cerr << "invalid IPv4 address\\n";
close(socket_fd);
return 1;
}
if (connect(socket_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("connect");
close(socket_fd);
return 1;
}
std::cout << "connected to 127.0.0.1:8888\\n"
<< "type text; Ctrl+D stops sending but keeps receiving\\n";
bool stdin_open = true;
while (true) {
fd_set read_set;
FD_ZERO(&read_set);
if (stdin_open) {
FD_SET(STDIN_FILENO, &read_set);
}
FD_SET(socket_fd, &read_set);
int max_fd = socket_fd > STDIN_FILENO
? socket_fd
: STDIN_FILENO;
int ready =
select(max_fd + 1,
&read_set,
nullptr,
nullptr,
nullptr);
if (ready == –1) {
if (errno == EINTR) {
continue;
}
std::perror("select");
break;
}
if (stdin_open &&
FD_ISSET(STDIN_FILENO, &read_set)) {
char buffer[4096];
ssize_t n = read(STDIN_FILENO,
buffer,
sizeof(buffer));
if (n > 0) {
if (!SendAll(socket_fd,
buffer,
static_cast<std::size_t>(n))) {
close(socket_fd);
return 1;
}
} else if (n == 0) {
stdin_open = false;
if (shutdown(socket_fd, SHUT_WR) == –1 &&
errno != ENOTCONN) {
std::perror("shutdown");
}
} else if (errno != EINTR) {
std::perror("read(stdin)");
break;
}
}
if (FD_ISSET(socket_fd, &read_set)) {
char buffer[4096];
ssize_t n = recv(socket_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
std::cout.write(buffer, n);
std::cout.flush();
} else if (n == 0) {
std::cout << "\\nserver closed connection\\n";
close(socket_fd);
return 0;
} else if (errno != EINTR) {
std::perror("recv");
break;
}
}
}
close(socket_fd);
return 1;
}
14.3 BuildClientEvents() 的意义
事件掩码不能再写死:
uint32_t BuildClientEvents(const ClientState& client) {
uint32_t events = EPOLLET | EPOLLRDHUP;
if (!client.peer_closed_write) {
events |= EPOLLIN;
}
if (PendingBytes(client) > 0) {
events |= EPOLLOUT;
}
return events;
}
它把连接状态映射成内核兴趣集合:
连接仍可接收数据 -> EPOLLIN
存在待发送数据 -> EPOLLOUT
使用边缘触发 -> EPOLLET
关注对端半关闭 -> EPOLLRDHUP
14.4 为什么收到数据后先主动发送
收到数据并加入发送队列后,不必等下一次 EPOLLOUT 才第一次尝试发送。
套接字此时很可能已经可写:
client.output.append(buffer, n);
FlushOutput(fd, client);
只有确实写不动时,才把 EPOLLOUT 加入事件掩码。
14.5 EPOLLRDHUP 不能直接等同于“所有数据已经读完”
这是一个非常容易写错的细节。
TCP 数据和 FIN 的顺序可能是:
业务数据 1
业务数据 2
FIN
当 epoll 报告 EPOLLRDHUP 时,接收缓冲区中仍可能存在 FIN 之前到达但尚未读取的业务数据。
因此正确做法是:
EPOLLRDHUP / EPOLLHUP
│
▼
继续执行非阻塞 recv()
│
├─ n > 0:继续处理遗留数据
└─ n == 0:才确认接收方向 EOF
EPOLLRDHUP 是关闭方向变化的提示,recv() == 0 才是接收方向到达 EOF 的确认。
14.6 本层修复了什么
- 正确处理部分发送;
- 正确处理 send() 返回 EAGAIN;
- 不再因慢客户端暂时写不动而立即断开;
- 仅在确实存在待发送数据时监听 EPOLLOUT;
- 支持客户端执行 shutdown(SHUT_WR) 后继续接收回显;
- 发送完最后的数据后再关闭连接;
- 使用 MSG_NOSIGNAL 防止 SIGPIPE 结束进程。
14.7 本层的新风险:内存无限增长
如果客户端持续发送,但从不读取服务端回显:
客户端不断发送
│
▼
服务端不断 append 到 output
│
▼
客户端不接收
│
▼
内核发送缓冲区满
│
▼
output 越积越多
│
▼
进程内存耗尽
因此“有发送队列”还不够,必须增加背压。
第九部分:背压是什么
十五、背压的核心含义
背压不是一个特定系统调用,而是一种流量控制策略:
当下游处理速度跟不上上游生产速度时,限制或暂停上游继续产生数据。
在 Echo 服务器中:
客户端发送速度 = 上游生产速度
客户端读取回显速度 = 下游消费速度
慢客户端造成的问题是:
recv 很快
send 很慢
output 持续增长
十六、高低水位
使用两个阈值而不是一个阈值,可以避免事件掩码频繁抖动:
output < 128 KiB
正常读取
output >= 256 KiB
暂停 EPOLLIN
暂停后 output 逐渐发送
output <= 128 KiB
恢复 EPOLLIN
ASCII 图:
待发送字节
▲
1M │ 硬上限:超过则关闭
│
256K├────────────── 高水位:暂停读取
│ ╲
│ ╲ 发送队列逐步下降
128K├────────────────╲── 低水位:恢复读取
│
0 └────────────────────────────────→ 时间
十七、为什么还需要硬上限
高水位只是暂停读取,但以下情况仍可能让发送队列突破预期:
- 单次读取的数据本身很大;
- 业务层一次生成大量响应;
- 多个处理步骤在暂停生效前已经追加数据;
- 代码缺陷导致继续入队。
因此还应设置硬上限:
if (pending + incoming > kOutputHardLimit) {
close_connection();
}
这不是为了“惩罚客户端”,而是为了保护整个进程不被单个连接耗尽内存。
第十部分:ET、公平预算和 EPOLLONESHOT
十八、为什么不能无限排空一个非常活跃的连接
ET 的经典规则是“一直处理到 EAGAIN”。
但如果某个客户端持续高速发送,单线程可能长期停留在该连接的读取循环中:
fd=7 持续有数据
│
▼
不断 recv()
│
▼
fd=8、fd=9 长时间得不到处理
这叫 事件循环饥饿。
一种工程化做法是给每个事件设置预算:
每次最多读取 64 KiB
每次最多发送 64 KiB
预算耗尽后,把控制权交回事件循环,让其他就绪连接获得处理机会。
十九、ET 下提前停止读取为什么危险
如果使用普通 ET:
读了 64 KiB
但还没有读到 EAGAIN
直接返回 epoll_wait
接收缓冲区仍然可读,却不一定出现新的边缘通知,连接可能停滞。
解决方法之一是:
- 使用用户态就绪队列;
- 或使用 EPOLLONESHOT,处理后通过 EPOLL_CTL_MOD 重新激活。
本文最终层采用第二种方案。
二十、EPOLLONESHOT 的语义
注册:
EPOLLIN | EPOLLET | EPOLLONESHOT
事件被 epoll_wait() 交付后,这个监控项暂时禁用。
应用处理完成后必须:
epoll_ctl(epoll_fd,
EPOLL_CTL_MOD,
fd,
&new_event);
重新激活。
这带来两个作用:
第十一部分:第 5 层——背压、预算与资源边界
二十一、最终层的状态字段
struct ClientState {
std::string output;
std::size_t sent_offset;
bool peer_closed_write;
bool read_paused;
};
其中:
| output | 已生成但尚未完全发送的数据 |
| sent_offset | output 中已经发送的前缀长度 |
| peer_closed_write | recv() 已返回 0 |
| read_paused | 因待发送数据过多而暂停监听 EPOLLIN |
21.1 完整服务端:05_epoll_et_backpressure/server.cpp
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <algorithm>
#include <cerrno>
#include <cstdio>
#include <iostream>
#include <string>
#include <unordered_map>
#include <vector>
namespace {
constexpr uint16_t kPort = 8888;
constexpr int kBacklog = 256;
constexpr int kMaxEvents = 128;
constexpr std::size_t kBufferSize = 16 * 1024;
constexpr std::size_t kMaxClients = 10000;
constexpr std::size_t kReadBudgetPerEvent = 64 * 1024;
constexpr std::size_t kWriteBudgetPerEvent = 64 * 1024;
constexpr std::size_t kOutputLowWatermark = 128 * 1024;
constexpr std::size_t kOutputHighWatermark = 256 * 1024;
constexpr std::size_t kOutputHardLimit = 1024 * 1024;
struct ClientState {
std::string peer;
std::string output;
std::size_t sent_offset = 0;
bool peer_closed_write = false;
bool read_paused = false;
};
bool SetNonBlocking(int fd) {
int old_flags = fcntl(fd, F_GETFL, 0);
if (old_flags == –1) {
return false;
}
return fcntl(fd,
F_SETFL,
old_flags | O_NONBLOCK) != –1;
}
std::size_t PendingBytes(const ClientState& client) {
return client.output.size() – client.sent_offset;
}
void CompactOutput(ClientState& client) {
if (client.sent_offset == 0) {
return;
}
if (client.sent_offset == client.output.size()) {
client.output.clear();
client.sent_offset = 0;
return;
}
if (client.sent_offset >= 4096 &&
client.sent_offset * 2 >= client.output.size()) {
client.output.erase(0, client.sent_offset);
client.sent_offset = 0;
}
}
uint32_t BuildClientEvents(const ClientState& client) {
uint32_t events =
EPOLLET | EPOLLONESHOT | EPOLLRDHUP;
if (!client.peer_closed_write &&
!client.read_paused) {
events |= EPOLLIN;
}
if (PendingBytes(client) > 0) {
events |= EPOLLOUT;
}
return events;
}
bool RearmClient(int epoll_fd,
int client_fd,
const ClientState& client) {
epoll_event event{};
event.events = BuildClientEvents(client);
event.data.fd = client_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_MOD,
client_fd,
&event) == –1) {
std::perror("epoll_ctl(MOD/rearm)");
return false;
}
return true;
}
void CloseClient(
int epoll_fd,
int client_fd,
std::unordered_map<int, ClientState>& clients) {
if (epoll_ctl(epoll_fd,
EPOLL_CTL_DEL,
client_fd,
nullptr) == –1 &&
errno != ENOENT &&
errno != EBADF) {
std::perror("epoll_ctl(DEL)");
}
close(client_fd);
clients.erase(client_fd);
std::cout << "closed fd=" << client_fd << '\\n';
}
bool QueueEcho(ClientState& client,
const char* data,
std::size_t size) {
CompactOutput(client);
if (PendingBytes(client) + size >
kOutputHardLimit) {
return false;
}
client.output.append(data, size);
if (PendingBytes(client) >=
kOutputHighWatermark) {
client.read_paused = true;
}
return true;
}
bool FlushOutput(int client_fd,
ClientState& client,
std::size_t budget) {
std::size_t written_this_round = 0;
while (client.sent_offset < client.output.size() &&
written_this_round < budget) {
std::size_t remain =
client.output.size() – client.sent_offset;
std::size_t allowed =
std::min(remain,
budget – written_this_round);
ssize_t n =
send(client_fd,
client.output.data() +
client.sent_offset,
allowed,
MSG_NOSIGNAL);
if (n > 0) {
std::size_t written =
static_cast<std::size_t>(n);
client.sent_offset += written;
written_this_round += written;
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
if (n == –1 &&
(errno == EAGAIN || errno == EWOULDBLOCK)) {
break;
}
std::perror("send");
return false;
}
CompactOutput(client);
if (client.read_paused &&
PendingBytes(client) <=
kOutputLowWatermark) {
client.read_paused = false;
}
return true;
}
}
int main() {
int listen_fd = socket(AF_INET,
SOCK_STREAM | SOCK_CLOEXEC,
0);
if (listen_fd == –1) {
std::perror("socket");
return 1;
}
if (!SetNonBlocking(listen_fd)) {
std::perror("fcntl(listen_fd)");
close(listen_fd);
return 1;
}
int reuse = 1;
if (setsockopt(listen_fd,
SOL_SOCKET,
SO_REUSEADDR,
&reuse,
sizeof(reuse)) == –1) {
std::perror("setsockopt");
close(listen_fd);
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_addr.s_addr = htonl(INADDR_ANY);
server_addr.sin_port = htons(kPort);
if (bind(listen_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("bind");
close(listen_fd);
return 1;
}
if (listen(listen_fd, kBacklog) == –1) {
std::perror("listen");
close(listen_fd);
return 1;
}
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
if (epoll_fd == –1) {
std::perror("epoll_create1");
close(listen_fd);
return 1;
}
epoll_event listen_event{};
listen_event.events = EPOLLIN | EPOLLET;
listen_event.data.fd = listen_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
listen_fd,
&listen_event) == –1) {
std::perror("epoll_ctl(ADD listen_fd)");
close(epoll_fd);
close(listen_fd);
return 1;
}
std::vector<epoll_event> ready_events(kMaxEvents);
std::unordered_map<int, ClientState> clients;
std::cout
<< "epoll ET backpressure server listening on 0.0.0.0:"
<< kPort << '\\n';
while (true) {
int ready_count =
epoll_wait(epoll_fd,
ready_events.data(),
static_cast<int>(ready_events.size()),
–1);
if (ready_count == –1) {
if (errno == EINTR) {
continue;
}
std::perror("epoll_wait");
break;
}
for (int i = 0; i < ready_count; ++i) {
int current_fd = ready_events[i].data.fd;
uint32_t events = ready_events[i].events;
if (current_fd == listen_fd) {
while (true) {
sockaddr_in client_addr{};
socklen_t client_len =
sizeof(client_addr);
int client_fd =
accept(listen_fd,
reinterpret_cast<sockaddr*>(
&client_addr),
&client_len);
if (client_fd == –1) {
if (errno == EINTR ||
errno == ECONNABORTED) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("accept");
break;
}
if (clients.size() >= kMaxClients) {
std::cerr
<< "reject fd=" << client_fd
<< ": client limit reached\\n";
close(client_fd);
continue;
}
if (!SetNonBlocking(client_fd)) {
std::perror("fcntl(client_fd)");
close(client_fd);
continue;
}
char ip[INET_ADDRSTRLEN]{};
inet_ntop(AF_INET,
&client_addr.sin_addr,
ip,
sizeof(ip));
ClientState client;
client.peer =
std::string(ip) + ":" +
std::to_string(
ntohs(client_addr.sin_port));
epoll_event client_event{};
client_event.events =
BuildClientEvents(client);
client_event.data.fd = client_fd;
if (epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&client_event) == –1) {
std::perror(
"epoll_ctl(ADD client)");
close(client_fd);
continue;
}
std::cout << "connected: "
<< client.peer
<< ", fd=" << client_fd << '\\n';
clients.emplace(client_fd,
std::move(client));
}
continue;
}
auto it = clients.find(current_fd);
if (it == clients.end()) {
continue;
}
ClientState& client = it->second;
bool close_now = false;
if ((events & EPOLLERR) != 0U) {
int socket_error = 0;
socklen_t error_len = sizeof(socket_error);
getsockopt(current_fd,
SOL_SOCKET,
SO_ERROR,
&socket_error,
&error_len);
std::cerr << "socket error on fd="
<< current_fd
<< ": errno=" << socket_error
<< '\\n';
close_now = true;
}
if (!close_now &&
!client.read_paused &&
(events &
(EPOLLIN | EPOLLRDHUP | EPOLLHUP)) != 0U) {
std::size_t read_this_round = 0;
while (read_this_round <
kReadBudgetPerEvent) {
char buffer[kBufferSize];
std::size_t allowed =
std::min(
sizeof(buffer),
kReadBudgetPerEvent –
read_this_round);
ssize_t n = recv(current_fd,
buffer,
allowed,
0);
if (n > 0) {
std::size_t received =
static_cast<std::size_t>(n);
if (!QueueEcho(client,
buffer,
received)) {
std::cerr
<< "fd=" << current_fd
<< " output hard limit exceeded\\n";
close_now = true;
break;
}
read_this_round += received;
if (client.read_paused) {
break;
}
continue;
}
if (n == 0) {
client.peer_closed_write = true;
break;
}
if (errno == EINTR) {
continue;
}
if (errno == EAGAIN ||
errno == EWOULDBLOCK) {
break;
}
std::perror("recv");
close_now = true;
break;
}
}
// 不直接把 EPOLLRDHUP/EPOLLHUP 等同于“所有数据已读完”。
// FIN 前面仍可能排着业务数据;只有 recv() 返回 0 才确认 EOF。
if (!close_now &&
PendingBytes(client) > 0) {
if (!FlushOutput(current_fd,
client,
kWriteBudgetPerEvent)) {
close_now = true;
}
}
if (!close_now &&
client.peer_closed_write &&
PendingBytes(client) == 0) {
close_now = true;
}
if (close_now) {
CloseClient(epoll_fd,
current_fd,
clients);
continue;
}
// EPOLLONESHOT 在事件交付后会禁用该 fd。
// 这里通过 MOD 重新激活,同时写入最新的读写兴趣集合。
if (!RearmClient(epoll_fd,
current_fd,
client)) {
CloseClient(epoll_fd,
current_fd,
clients);
}
}
}
for (const auto& [client_fd, client] : clients) {
(void)client;
close(client_fd);
}
close(epoll_fd);
close(listen_fd);
return 0;
}
21.2 完整客户端:05_epoll_et_backpressure/client.cpp
#include <arpa/inet.h>
#include <sys/select.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <iostream>
namespace {
bool SendAll(int fd, const char* data, std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
std::perror("send");
return false;
}
return true;
}
}
int main() {
int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
if (socket_fd == –1) {
std::perror("socket");
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(8888);
if (inet_pton(AF_INET,
"127.0.0.1",
&server_addr.sin_addr) != 1) {
std::cerr << "invalid IPv4 address\\n";
close(socket_fd);
return 1;
}
if (connect(socket_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("connect");
close(socket_fd);
return 1;
}
std::cout << "connected to 127.0.0.1:8888\\n"
<< "type text; Ctrl+D stops sending but keeps receiving\\n";
bool stdin_open = true;
while (true) {
fd_set read_set;
FD_ZERO(&read_set);
if (stdin_open) {
FD_SET(STDIN_FILENO, &read_set);
}
FD_SET(socket_fd, &read_set);
int max_fd = socket_fd > STDIN_FILENO
? socket_fd
: STDIN_FILENO;
int ready =
select(max_fd + 1,
&read_set,
nullptr,
nullptr,
nullptr);
if (ready == –1) {
if (errno == EINTR) {
continue;
}
std::perror("select");
break;
}
if (stdin_open &&
FD_ISSET(STDIN_FILENO, &read_set)) {
char buffer[4096];
ssize_t n = read(STDIN_FILENO,
buffer,
sizeof(buffer));
if (n > 0) {
if (!SendAll(socket_fd,
buffer,
static_cast<std::size_t>(n))) {
close(socket_fd);
return 1;
}
} else if (n == 0) {
stdin_open = false;
if (shutdown(socket_fd, SHUT_WR) == –1 &&
errno != ENOTCONN) {
std::perror("shutdown");
}
} else if (errno != EINTR) {
std::perror("read(stdin)");
break;
}
}
if (FD_ISSET(socket_fd, &read_set)) {
char buffer[4096];
ssize_t n = recv(socket_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
std::cout.write(buffer, n);
std::cout.flush();
} else if (n == 0) {
std::cout << "\\nserver closed connection\\n";
close(socket_fd);
return 0;
} else if (errno != EINTR) {
std::perror("recv");
break;
}
}
}
close(socket_fd);
return 1;
}
21.3 最终层的保护机制
连接数上限
constexpr std::size_t kMaxClients = 10000;
真实可支持数量还受以下因素约束:
- RLIMIT_NOFILE;
- 系统级文件描述符限制;
- epoll watch 限制;
- 每连接用户态状态;
- socket 缓冲区;
- 进程内存;
- CPU 和业务处理成本。
因此“epoll 没有 1024 限制”不等于“可以无限建立连接”。
单连接发送队列硬上限
constexpr std::size_t kOutputHardLimit = 1024 * 1024;
单个客户端不能无限占用内存。
高低水位背压
constexpr std::size_t kOutputLowWatermark = 128 * 1024;
constexpr std::size_t kOutputHighWatermark = 256 * 1024;
高水位暂停读取,低水位恢复读取。
每事件读写预算
constexpr std::size_t kReadBudgetPerEvent = 64 * 1024;
constexpr std::size_t kWriteBudgetPerEvent = 64 * 1024;
防止一个活跃连接长时间独占事件循环。
EPOLLONESHOT 重装
每轮处理结束后:
RearmClient(epoll_fd, fd, client);
把最新的读取、写入、暂停状态重新提交给 epoll。
21.4 最终层完整流程图
┌──────────────────────────────────────────────────────────────┐
│ epoll_wait │
└──────────────────────────────┬───────────────────────────────┘
│
▼
取出本轮就绪 epoll_event
│
┌─────────────────┴──────────────────┐
│ │
▼ ▼
fd == listen_fd client fd
│ │
▼ ▼
循环 accept 到 EAGAIN 检查 EPOLLERR
│ │
▼ ▼
新 fd 设置非阻塞 按预算循环 recv
│ │
▼ ▼
注册 EPOLLIN/ET 数据加入 output
│
▼
达到高水位?
│ │
是 否
│ │
▼ ▼
暂停读 继续读
│
▼
按预算循环 send
│
┌───────────────────────┴─────────────┐
│ │
▼ ▼
全部发完 EAGAIN/预算耗尽
│ │
▼ ▼
移除 EPOLLOUT 保留 EPOLLOUT
│
▼
低于低水位时恢复 EPOLLIN
│
▼
EPOLL_CTL_MOD 重新激活 ONESHOT
第十二部分:epoll 三大 API 逐项解析
二十二、epoll_create1()
int epoll_create1(int flags);
推荐:
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
EPOLL_CLOEXEC 可以避免程序以后调用 exec() 时意外把 epoll fd 泄漏到新程序。
返回值:
>= 0 epoll 实例文件描述符
-1 失败,检查 errno
二十三、epoll_ctl()
int epoll_ctl(int epoll_fd,
int operation,
int target_fd,
epoll_event* event);
三种常用操作:
| EPOLL_CTL_ADD | 增加新的监控项 |
| EPOLL_CTL_MOD | 修改已有监控项的事件掩码或用户数据 |
| EPOLL_CTL_DEL | 删除监控项 |
ADD
epoll_event event{};
event.events = EPOLLIN;
event.data.fd = client_fd;
epoll_ctl(epoll_fd,
EPOLL_CTL_ADD,
client_fd,
&event);
MOD
event.events = EPOLLIN | EPOLLOUT | EPOLLET;
epoll_ctl(epoll_fd,
EPOLL_CTL_MOD,
client_fd,
&event);
DEL
epoll_ctl(epoll_fd,
EPOLL_CTL_DEL,
client_fd,
nullptr);
close() 会在相关 open file description 的引用最终消失时清理 epoll 关联,但显式 DEL 可以让单线程示例中的用户态生命周期更加清楚。复杂程序还必须考虑 dup()、fork() 和描述符复用产生的对象所有权问题。
二十四、epoll_wait()
int epoll_wait(int epoll_fd,
epoll_event* events,
int max_events,
int timeout_ms);
超时参数:
| -1 | 一直等待 |
| 0 | 不等待,立即返回 |
| > 0 | 最多等待对应毫秒数 |
返回值:
| > 0 | 返回的就绪事件数量 |
| 0 | 超时 |
| -1 | 出错 |
信号中断:
if (count == –1 && errno == EINTR) {
continue;
}
通常不是致命错误。
第十三部分:常用事件标志
二十五、事件位速查表
| EPOLLIN | 可能可读 | accept()、recv()、读取管道等 |
| EPOLLOUT | 可能可写 | 继续发送待发送队列 |
| EPOLLRDHUP | 流式连接对端关闭写方向 | 继续读完剩余数据,等待 recv()==0 |
| EPOLLERR | 异步错误 | 用 getsockopt(SO_ERROR) 获取错误 |
| EPOLLHUP | 挂断 | 仍应考虑读取残留数据 |
| EPOLLET | 边缘触发 | 非阻塞并处理到 EAGAIN |
| EPOLLONESHOT | 一次交付后禁用 | EPOLL_CTL_MOD 重新激活 |
| EPOLLEXCLUSIVE | 某些多 epoll 等待场景的独占唤醒 | 主要用于缓解特定惊群场景 |
EPOLLERR 和 EPOLLHUP 即使没有显式加入兴趣掩码,也可能被返回,因此事件循环必须处理它们。
二十六、为什么使用 SO_ERROR
int socket_error = 0;
socklen_t error_len = sizeof(socket_error);
getsockopt(fd,
SOL_SOCKET,
SO_ERROR,
&socket_error,
&error_len);
EPOLLERR 只是告诉你套接字存在错误状态,SO_ERROR 才用于读取挂起的具体套接字错误码。
第十四部分:accept、recv、send 的完整状态表
二十七、非阻塞 accept
| fd >= 0 | 成功取得一个连接 |
| -1, EINTR | 被信号中断,重试 |
| -1, ECONNABORTED | 队列中的连接已经中止,可继续 |
| -1, EAGAIN/EWOULDBLOCK | 当前连接队列已排空 |
| 其他错误 | 记录并根据策略处理 |
二十八、非阻塞 recv
| n > 0 | 收到 n 字节 |
| n == 0 | 接收方向 EOF |
| -1, EINTR | 重试 |
| -1, EAGAIN/EWOULDBLOCK | 当前已无数据可读 |
| 其他错误 | 连接异常 |
二十九、非阻塞 send
| n > 0 | 发送了 n 字节,可能小于请求长度 |
| -1, EINTR | 重试 |
| -1, EAGAIN/EWOULDBLOCK | 当前发送缓冲区满,等待 EPOLLOUT |
| -1, EPIPE | 对端已经关闭相应方向 |
| 其他错误 | 连接异常 |
第十五部分:TCP 半关闭
三十、shutdown(SHUT_WR) 的意义
客户端按下 Ctrl+D 后,本文客户端不会立即 close(),而是:
shutdown(socket_fd, SHUT_WR);
这表示:
客户端:
不再发送新数据
但仍然可以接收服务端数据
对应的 TCP 状态是发送 FIN,但连接的另一个方向仍可继续传输。
流程:
客户端发送业务数据
│
▼
shutdown(SHUT_WR)
│
▼
服务端读到业务数据
│
▼
服务端 recv() 最终返回 0
│
▼
服务端发送完剩余回显
│
▼
服务端 close()
│
▼
客户端读到全部回显和 EOF
这也是为什么服务端不能一看到 EPOLLRDHUP 就立即丢弃发送队列并关闭。
第十六部分:LT 与 ET 应该如何选
三十一、对比表
| 默认模式 | 是 | 否,需要 EPOLLET |
| 编程难度 | 较低 | 较高 |
| 非阻塞要求 | 强烈建议 | 实际使用中应当采用 |
| 读取策略 | 可以分批处理 | 必须确保不会遗留未处理就绪状态 |
| 写入策略 | 仍要处理部分发送和 EAGAIN | 同样需要,且事件状态更敏感 |
| 重复通知 | 条件持续满足时可能重复 | 通知更依赖状态变化 |
| 出错风险 | 相对低 | 忘记排空、错误暂停更容易停滞 |
| 适用场景 | 大多数服务端的安全起点 | 明确理解状态机并需要精细控制时 |
三十二、ET 一定比 LT 快吗
不一定。
性能取决于:
- 活跃连接比例;
- 每次事件的数据量;
- 系统调用频率;
- 业务处理成本;
- 缓冲区大小;
- 事件循环公平性;
- 内存分配;
- 锁竞争;
- 网络延迟;
- CPU 缓存行为。
在许多服务中,业务逻辑、序列化、数据库和日志的成本远高于 LT/ET 的差异。
选择原则应是:
先实现正确、可测量的 LT
│
▼
通过性能分析确认事件调度成为瓶颈
│
▼
再评估 ET 是否能带来实际收益
不要仅凭“ET 更高级”选择 ET。
第十七部分:压力测试客户端
三十三、完整大数据校验客户端
该客户端:
#include <arpa/inet.h>
#include <sys/socket.h>
#include <unistd.h>
#include <cerrno>
#include <cstdio>
#include <cstdlib>
#include <iostream>
#include <string>
namespace {
bool SendAll(int fd,
const char* data,
std::size_t size) {
std::size_t offset = 0;
while (offset < size) {
ssize_t n = send(fd,
data + offset,
size – offset,
MSG_NOSIGNAL);
if (n > 0) {
offset += static_cast<std::size_t>(n);
continue;
}
if (n == –1 && errno == EINTR) {
continue;
}
std::perror("send");
return false;
}
return true;
}
}
int main(int argc, char* argv[]) {
std::size_t message_size = 1024 * 1024;
if (argc >= 2) {
message_size =
static_cast<std::size_t>(
std::strtoull(argv[1], nullptr, 10));
}
int socket_fd = socket(AF_INET, SOCK_STREAM, 0);
if (socket_fd == –1) {
std::perror("socket");
return 1;
}
sockaddr_in server_addr{};
server_addr.sin_family = AF_INET;
server_addr.sin_port = htons(8888);
if (inet_pton(AF_INET,
"127.0.0.1",
&server_addr.sin_addr) != 1) {
std::cerr << "invalid IPv4 address\\n";
close(socket_fd);
return 1;
}
if (connect(socket_fd,
reinterpret_cast<sockaddr*>(&server_addr),
sizeof(server_addr)) == –1) {
std::perror("connect");
close(socket_fd);
return 1;
}
std::string payload(message_size, 'X');
if (!SendAll(socket_fd,
payload.data(),
payload.size())) {
close(socket_fd);
return 1;
}
if (shutdown(socket_fd, SHUT_WR) == –1) {
std::perror("shutdown");
close(socket_fd);
return 1;
}
std::string received;
received.reserve(payload.size());
while (true) {
char buffer[16 * 1024];
ssize_t n = recv(socket_fd,
buffer,
sizeof(buffer),
0);
if (n > 0) {
received.append(
buffer,
static_cast<std::size_t>(n));
continue;
}
if (n == 0) {
break;
}
if (errno == EINTR) {
continue;
}
std::perror("recv");
close(socket_fd);
return 1;
}
close(socket_fd);
bool ok = received == payload;
std::cout << "sent=" << payload.size()
<< ", received=" << received.size()
<< ", verify=" << (ok ? "OK" : "FAILED")
<< '\\n';
return ok ? 0 : 1;
}
33.1 测试 1 MiB 回显
./bin/05_epoll_et_server
另一个终端:
./bin/stress_client 1048576
预期:
sent=1048576, received=1048576, verify=OK
33.2 为什么不能用一次 write 作为压力测试
错误写法:
write(fd, big_message.data(), big_message.size());
如果返回值小于 big_message.size(),剩余数据没有被发送。
压力测试本身也必须正确处理部分发送,否则无法区分:
- 服务端错误;
- 客户端发送不完整。
第十八部分:编译与运行
三十四、编译全部版本
make
使用的编译参数:
g++ -std=c++17 -O2 -Wall -Wextra -Wpedantic
三十五、运行指定层次
第 0 层
./bin/00_select_server
./bin/00_select_client
第 1 层
./bin/01_epoll_lt_server
./bin/01_epoll_lt_client
第 2 层
./bin/02_epoll_lt_server
./bin/02_epoll_lt_client
第 3 层
./bin/03_epoll_et_server
./bin/03_epoll_et_client
第 4 层
./bin/04_epoll_et_server
./bin/04_epoll_et_client
第 5 层
./bin/05_epoll_et_server
./bin/05_epoll_et_client
客户端输入结束:
Ctrl+D
客户端会半关闭写方向,但继续读取服务端回显。
第十九部分:调试方法
三十六、使用 strace 观察事件循环
strace -f \\
-e trace=epoll_create1,epoll_ctl,epoll_wait,accept,recvfrom,sendto,close \\
./bin/05_epoll_et_server
重点观察:
epoll_wait(…)
accept(…)
recvfrom(…)
sendto(…)
epoll_ctl(… EPOLL_CTL_MOD …)
三十七、使用 ss 查看连接
ss -ltnp 'sport = :8888'
查看监听状态:
LISTEN
查看已建立连接:
ss -tnp 'sport = :8888 or dport = :8888'
三十八、模拟慢客户端
可以让客户端发送大量数据,却暂时不读取回显,从而观察:
- 服务端发送队列增长;
- 达到高水位后暂停 EPOLLIN;
- 发送队列下降后恢复读取;
- 超过硬上限后的断开策略。
第二十部分:常见错误清单
三十九、错误:ET 只加 EPOLLET,不改读取方式
events = EPOLLIN | EPOLLET;
read(fd, buf, sizeof(buf)); // 只读一次
后果:未读数据可能滞留,处理停滞。
四十、错误:ET 使用阻塞 fd 循环读
while (read(fd, buf, sizeof(buf)) > 0) {
}
后果:缓冲区读空后阻塞整个事件循环。
四十一、错误:把 EAGAIN 当成连接失败
if (recv(...) == –1) {
close(fd);
}
后果:正常的非阻塞控制流被误判为错误。
四十二、错误:永久监听 EPOLLOUT
events = EPOLLIN | EPOLLOUT;
后果:没有待发送数据时仍可能持续被可写事件唤醒。
四十三、错误:忽略部分发送
send(fd, data, size, 0);
后果:只发送前半部分,应用却以为全部完成。
四十四、错误:发送队列没有上限
后果:慢客户端可以使进程内存持续增长。
四十五、错误:ET 设置处理预算但不重新激活
后果:预算耗尽时缓冲区可能仍然可读,但没有新边缘触发。
四十六、错误:EPOLLRDHUP 立即 close
后果:可能丢弃 FIN 之前已经到达但尚未读取的数据,或丢弃尚未发送的响应。
四十七、错误:在 epoll_event.data.ptr 中保存裸指针却不管理生命周期
如果连接对象被释放,而 ready event 中仍保留旧指针,就可能出现 use-after-free。
本文使用 data.fd 和 unordered_map<int, ClientState>,降低了示例中的生命周期复杂度。生产代码可以使用指针,但必须结合:
- 稳定对象地址;
- 引用计数;
- generation 标识;
- 延迟销毁;
- 线程同步。
第二十一部分:进一步工程化方向
四十八、应用层协议拆包
Echo 示例把收到的字节原样返回,不代表一次 recv() 就是一条业务消息。
生产服务通常需要:
- 固定长度协议;
- 分隔符协议;
- 长度字段 + Payload;
- HTTP/WebSocket 等成熟协议解析器。
长度头示意:
┌──────────────┬─────────────────────┐
│ 4 字节长度 N │ N 字节 Payload │
└──────────────┴─────────────────────┘
每连接需要独立接收缓冲区:
struct ClientState {
std::string input;
std::string output;
};
解析原则:
input 不够一帧:
保留,等待下次 recv
input 包含完整一帧:
取出并处理
input 包含多帧:
循环解析
四十九、定时器
可以使用:
- timerfd;
- 时间轮;
- 最小堆;
- 周期性 epoll_wait 超时。
用于处理:
- 空闲连接超时;
- 握手超时;
- 请求处理超时;
- 心跳。
五十、信号
可以使用 signalfd 把信号也变成 epoll 事件,避免传统异步信号处理函数中的限制。
五十一、多线程
常见模型包括:
单 acceptor + 多 event loop
每线程一个 epoll
主从 Reactor
EPOLLONESHOT + 工作线程
SO_REUSEPORT 多监听实例
多线程下要额外解决:
- 连接归属;
- 同一 fd 是否会并发处理;
- 对象销毁时机;
- 跨线程唤醒;
- 发送队列锁;
- fd 复用;
- 惊群。
五十二、跨线程唤醒
常使用:
- eventfd;
- pipe;
- socketpair。
工作线程向 event loop 提交任务后,通过 eventfd 唤醒 epoll_wait()。
第二十二部分:FAQ
五十三、close 前必须 EPOLL_CTL_DEL 吗
简单单线程程序中,关闭文件描述符的最终引用后,内核会清理相关 epoll 监控关系。
但显式 DEL 仍有价值:
- 用户态生命周期清晰;
- 错误更容易定位;
- 复杂的 dup()/共享 open file description 场景中更容易控制;
- 避免连接状态容器和 epoll 兴趣集合长时间不一致。
不能简单总结成“close 一定会产生幽灵事件”或“DEL 永远必须”。
五十四、accepted socket 会继承 O_NONBLOCK 吗
在 Linux 上,accept() 返回的新套接字不会自动继承监听套接字的 O_NONBLOCK 文件状态标志。
因此需要:
fcntl(client_fd, F_SETFL, ... | O_NONBLOCK);
或者:
accept4(..., SOCK_NONBLOCK | SOCK_CLOEXEC);
五十五、EAGAIN 和 EWOULDBLOCK 是否相同
Linux 上通常相同,但可移植代码应同时检查:
errno == EAGAIN ||
errno == EWOULDBLOCK
五十六、LT 可以使用阻塞套接字吗
从接口语义上可以,但事件驱动服务器通常仍建议使用非阻塞套接字,以保证任何 I/O 操作都不会意外长时间阻塞整个事件循环。
五十七、ET 是否每次都必须读到 EAGAIN
最稳定、最容易验证的通用规则是:
非阻塞读取到 EAGAIN
对流式 fd,某些情况下短读也可以说明当前数据已经取尽,但把所有 ET 路径统一写成“循环到 EAGAIN”更不容易出错。
如果为了公平预算提前停止,就必须使用用户态重新调度或 EPOLLONESHOT 重装等机制,不能直接回到普通 epoll_wait() 后假设一定会有新边缘。
五十八、epoll 可以监控普通文件吗
常规磁盘文件通常不适合通过 epoll 获取就绪通知,epoll_ctl() 可能返回 EPERM。epoll 主要面向支持轮询语义的对象,例如 socket、pipe、eventfd、timerfd、signalfd 等。
五十九、epoll 能替代线程池吗
不能。
epoll 解决的是:
大量 I/O 对象的就绪等待
线程池解决的是:
CPU 密集型或可能阻塞的业务任务并发执行
两者经常组合使用。
第二十三部分:六层演进总结
六十、每一层到底学到了什么
第 0 层:select
理解:
- fd 集合;
- select() 破坏性修改;
- 全量扫描;
- FD_SETSIZE。
第 1 层:epoll LT 最小迁移
理解:
- epoll 实例;
- interest list;
- ready events;
- epoll_create1/ctl/wait。
第 2 层:非阻塞 LT
理解:
- 非阻塞套接字;
- EINTR;
- EAGAIN;
- 排空连接队列;
- LT 也应避免阻塞事件循环。
第 3 层:ET 正确读取
理解:
- EPOLLET;
- 单次 read 的停滞问题;
- 阻塞循环读的事件循环卡顿;
- 非阻塞 + 处理到 EAGAIN。
第 4 层:写事件状态机
理解:
- 部分发送;
- 发送偏移;
- EPOLLOUT 按需订阅;
- 半关闭;
- MSG_NOSIGNAL。
第 5 层:工程边界
理解:
- 背压;
- 高低水位;
- 内存硬上限;
- 每事件预算;
- EPOLLONESHOT;
- 公平性。
六十一、最终核心原则
可以把整篇文章压缩成六句话:
1. epoll 只报告就绪,不替你执行 I/O。
2. 非阻塞 I/O 的 EAGAIN 是正常边界,不是错误。
3. ET 收到事件后,要把当前就绪状态处理完整。
4. TCP send 可能部分完成,必须保存发送状态。
5. EPOLLOUT 只在存在待发送数据时监听。
6. 没有背压和资源上限的发送队列不是完整设计。
附录 A:源码目录
epoll_分层升级技术博客/
├── 00_select_original/
│ ├── server.cpp
│ └── client.cpp
├── 01_epoll_lt_minimal/
│ ├── server.cpp
│ └── client.cpp
├── 02_epoll_lt_nonblocking/
│ ├── server.cpp
│ └── client.cpp
├── 03_epoll_et_correct_read/
│ ├── server.cpp
│ └── client.cpp
├── 04_epoll_et_write_queue/
│ ├── server.cpp
│ └── client.cpp
├── 05_epoll_et_backpressure/
│ ├── server.cpp
│ └── client.cpp
├── pitfalls/
│ ├── et_single_read_server.cpp
│ └── et_blocking_loop_server.cpp
├── tools/
│ └── stress_client.cpp
├── Makefile
├── README.md
└── 从_select_到_epoll_C++17_完整技术博客.md
附录 B:官方资料名称
进一步核对接口语义时,应优先查阅 Linux man-pages 中的:
- epoll(7)
- epoll_create(2)
- epoll_ctl(2)
- epoll_wait(2)
- accept(2)
- recv(2)
- send(2)
- socket(7)
- fcntl(2)
本文的重点不是记住某份最终代码,而是建立可以迁移到真实网络项目的状态机思维:
每个连接都不是一个“fd 数字”,而是一个持续变化的 I/O 状态对象;epoll 只是让这些状态对象在合适的时机获得处理机会。
网硕互联帮助中心



评论前必须登录!
注册