一个最简单的 TCP 服务器通常为每个客户端创建一个线程,线程阻塞在 accept、recv 或 send 上。这种写法容易理解,但连接数增加后会带来线程栈、上下文切换、锁竞争和调度开销。线程数量与连接数量绑定,也使超时、广播和连接回收变得复杂。
Reactor 模式采用另一种组织方式:I/O 多路复用器负责观察大量文件描述符,事件循环在某个或多个就绪事件到达后,再调用对应连接的处理逻辑。它并不会让单个连接的网络操作变快,而是减少了无效等待,使有限的执行线程能够管理更多连接。
本文实现一个基于 Linux epoll 的非阻塞 TCP 回显服务器。它只承担演示职责,不包含 TLS、认证、协议解析和生产级限流;这些能力应在明确需求后单独设计。
核心原理
Reactor 的职责划分
一个基本 Reactor 通常包含四类对象:
- 事件循环:调用
epoll_wait,取得已经就绪的文件描述符。 - 监听器:处理监听套接字上的可读事件,并循环调用
accept。 - 连接对象:保存客户端套接字、输入缓冲区、输出缓冲区和当前关注的事件。
- 事件分派器:根据文件描述符找到连接对象,调用读、写或错误处理函数。
关键点是回调函数必须尽量短,并且不能在非阻塞套接字上执行可能长时间等待的操作。否则,事件循环仍然会被一个连接拖住。
水平触发与边缘触发
epoll 支持水平触发和边缘触发。
水平触发模式下,只要接收缓冲区仍有数据,后续的 epoll_wait 仍可能报告可读事件,代码更容易验证。边缘触发模式只在状态从“无数据”变为“有数据”等变化时通知,因此处理函数必须一直读到 EAGAIN,否则剩余数据可能没有新的通知。
本文使用水平触发。生产环境若改用 EPOLLET,必须同时满足以下条件:套接字设置为非阻塞;accept 循环直到 EAGAIN;recv 循环直到 EAGAIN;发送缓冲区耗尽或返回 EAGAIN 时正确调整关注事件。
为什么必须处理部分读写
TCP 是字节流,不保留应用层消息边界。一次 recv 可能只得到半条消息,也可能得到多条消息;一次 send 也不保证把全部数据交给内核。因此,服务器需要维护输入和输出缓冲区。
本例使用回显协议:收到多少字节就回送多少字节。这样可以集中展示部分读写问题,而不会把示例混入复杂的业务协议。
环境准备
示例假定环境满足以下条件:
- Linux 系统,内核提供
epoll。 - C++17 编译器。
- 本机允许监听
9000端口。
创建文件 reactor_echo.cpp,内容如下:
#include <arpa/inet.h>
#include <cerrno>
#include <cstring>
#include <fcntl.h>
#include <iostream>
#include <netinet/in.h>
#include <sys/epoll.h>
#include <sys/socket.h>
#include <unistd.h>
#include <string>
#include <unordered_map>
namespace {
constexpr int kPort = 9000;
constexpr int kMaxEvents = 64;
struct Connection {
int fd;
std::string output;
};
bool set_nonblocking(int fd) {
int flags = fcntl(fd, F_GETFL, 0);
return flags >= 0 && fcntl(fd, F_SETFL, flags | O_NONBLOCK) == 0;
}
bool update_events(int epoll_fd, int fd, uint32_t events) {
epoll_event event{
};
event.events = events;
event.data.fd = fd;
return epoll_ctl(epoll_fd, EPOLL_CTL_MOD, fd, &event) == 0;
}
void close_connection(int epoll_fd,
std::unordered_map<int, Connection>& connections,
int fd) {
epoll_ctl(epoll_fd, EPOLL_CTL_DEL, fd, nullptr);
close(fd);
connections.erase(fd);
}
bool accept_clients(int listen_fd, int epoll_fd,
std::unordered_map<int, Connection>& connections) {
while (true) {
sockaddr_in peer{
};
socklen_t length = sizeof(peer);
int client_fd = accept4(listen_fd, reinterpret_cast<sockaddr*>(&peer),
&length, SOCK_NONBLOCK | SOCK_CLOEXEC);
if (client_fd < 0) {
if (errno == EAGAIN || errno == EWOULDBLOCK) return true;
if (errno == EINTR) continue;
return false;
}
epoll_event event{
};
event.events = EPOLLIN;
event.data.fd = client_fd;
if (epoll_ctl(epoll_fd, EPOLL_CTL_ADD, client_fd, &event) != 0) {
close(client_fd);
continue;
}
connections.emplace(client_fd, Connection{
client_fd, {
}});
}
}
bool read_client(int epoll_fd,
std::unordered_map<int, Connection>& connections,
int fd) {
char buffer[4096];
while (true) {
ssize_t count = recv(fd, buffer, sizeof(buffer), 0);
if (count > 0) {
connections.at(fd).output.append(buffer, static_cast<size_t>(count));
continue;
}
if (count == 0) return false;
if (errno == EINTR) continue;
if (errno == EAGAIN || errno == EWOULDBLOCK) break;
return false;
}
if (!connections.at(fd).output.empty()) {
return update_events(epoll_fd, fd, EPOLLIN | EPOLLOUT);
}
return true;
}
bool write_client(int epoll_fd,
std::unordered_map<int, Connection>& connections,
int fd) {
auto& output = connections.at(fd).output;
while (!output.empty()) {
ssize_t count = send(fd, output.data(), output.size(), MSG_NOSIGNAL);
if (count > 0) {
output.erase(0, static_cast<size_t>(count));
continue;
}
if (count < 0 && errno == EINTR) continue;
if (count < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) return true;
return false;
}
return update_events(epoll_fd, fd, EPOLLIN);
}
} // namespace
int main() {
int listen_fd = socket(AF_INET, SOCK_STREAM | SOCK_CLOEXEC, 0);
if (listen_fd < 0 || !set_nonblocking(listen_fd)) return 1;
int reuse = 1;
setsockopt(listen_fd, SOL_SOCKET, SO_REUSEADDR, &reuse, sizeof(reuse));
sockaddr_in address{
};
address.sin_family = AF_INET;
address.sin_addr.s_addr = htonl(INADDR_ANY);
address.sin_port = htons(kPort);
if (bind(listen_fd, reinterpret_cast<sockaddr*>(&address), sizeof(address)) != 0 ||
listen(listen_fd, SOMAXCONN) != 0) return 1;
int epoll_fd = epoll_create1(EPOLL_CLOEXEC);
if (epoll_fd < 0) return 1;
epoll_event listen_event{
};
listen_event.events = EPOLLIN;
listen_event.data.fd = listen_fd;
epoll_ctl(epoll_fd, EPOLL_CTL_ADD, listen_fd, &listen_event);
std::unordered_map<int, Connection> connections;
epoll_event events[kMaxEvents];
while (true) {
int ready = epoll_wait(epoll_fd, events, kMaxEvents, -1);
if (ready < 0) {
if (errno == EINTR) continue;
break;
}
for (int i = 0; i < ready; ++i) {
int fd = events[i].data.fd;
uint32_t flags = events[i].events;
if (fd == listen_fd) {
if (!accept_clients(listen_fd, epoll_fd, connections)) return 1;
continue;
}
bool keep = !(flags & (EPOLLERR | EPOLLHUP | EPOLLRDHUP));
if (keep && (flags & EPOLLIN)) keep = read_client(epoll_fd, connections, fd);
if (keep && (flags & EPOLLOUT)) keep = write_client(epoll_fd, connections, fd);
if (!keep) close_connection(epoll_fd, connections, fd);
}
}
close(epoll_fd);
close(listen_fd);
return 0;
}
编译与运行
使用编译器构建:
g++ -std=c++17 -O2 -Wall -Wextra -pedantic reactor_echo.cpp -o reactor_echo
./reactor_echo
另开终端测试:
printf 'hello reactor\n' | nc 127.0.0.1 9000
应能看到服务器返回相同内容。这里的验证只说明基本收发链路可用,并不能证明它已经具备生产环境所需的容量、延迟或安全性。
代码中的关键路径
监听套接字使用非阻塞模式,并注册到 epoll。收到可读事件后,accept_clients 会反复调用 accept4,直到返回 EAGAIN。这样即使一次事件中有多个新连接,也不会只处理一个。
客户端连接收到数据时,read_client 持续读取,直到暂时没有更多数据。数据进入 output 后,代码通过 EPOLL_CTL_MOD 增加 EPOLLOUT 关注。只有输出缓冲区非空时才关注可写事件,避免套接字始终可写导致事件循环空转。
写操作可能只发送一部分数据,所以每次成功发送后都从字符串前端移除已发送内容;遇到 EAGAIN 则保留剩余内容,等待下一次可写通知。连接关闭、对端半关闭或出现错误时,统一从 epoll 和连接表中删除文件描述符。
生产化时需要补齐的边界
设置输出上限
示例没有限制 output 大小。慢客户端持续发送时,服务器可能积累大量待发送数据。实际服务应设置单连接上限,超过上限后选择暂停读取、丢弃、限速或关闭连接。策略应与协议的可靠性要求一致。
增加协议帧解析
回显字节不等于处理消息。若协议使用换行、固定长度或长度前缀,应在输入缓冲区中循环提取完整帧,并保留未完成的半帧。长度字段必须校验上限,不能直接依据客户端提供的长度分配任意大小内存。
处理信号与优雅退出
可以通过 signalfd 将信号转成 epoll 事件,也可以采用自管道方案。收到退出信号后停止接收新连接,关闭或排空已有连接,再释放资源。不要依赖进程被强制终止来完成清理。
评估线程模型
单线程 Reactor 适合 I/O 逻辑简单的服务。如果业务计算耗时,应该把计算任务投递到工作线程池,完成后再把结果安全地交回事件线程。多个线程直接操作同一个连接,必须明确所有权和唤醒机制,否则容易产生竞态。
常见问题
为什么 recv 返回零?
这表示对端进行了有序关闭。服务器应停止继续读取,并释放该连接。若业务需要半关闭语义,也可以先关闭写方向,但必须明确协议状态。
为什么 send 会触发 SIGPIPE?
对端已经关闭时,向某些套接字写入可能产生 SIGPIPE。示例使用 MSG_NOSIGNAL 抑制该信号,并通过返回值处理错误。也可以在进程级别忽略 SIGPIPE,但进程级设置会影响其他代码,使用前应评估。
什么时候需要边缘触发?
当事件数量较大、希望减少重复通知时,可以考虑边缘触发。但它要求所有读写函数严格循环到 EAGAIN,并正确管理状态。若团队还没有稳定的缓冲区和状态机测试,水平触发通常更容易维护。
EPOLLHUP 是否意味着可以忽略剩余数据?
不能简单忽略。错误和挂起事件的具体表现取决于连接状态;某些情况下仍应先尝试读取内核中已经排队的数据。工程代码应根据协议和实际错误码制定关闭顺序,而不是只看一个标志位。
总结
Reactor 的核心不是“使用 epoll”这一行代码,而是把连接生命周期、事件关注、部分读写和缓冲区状态组织成一个可验证的状态机。最小实现可以从水平触发开始:所有套接字非阻塞,监听和读取循环到 EAGAIN,写缓冲区非空时才关注可写事件,关闭路径统一注销并释放资源。
在此基础上,下一步应补充协议解析、缓冲区上限、超时、优雅退出、指标和并发测试。只有这些边界被明确处理,事件循环才具备从示例走向长期运行服务的基础。