Linux Reactor 反应堆实战:用 epoll 构建可扩展的 TCP 事件循环

简介: 本文实现了一个基于Linux epoll的非阻塞TCP回显服务器,演示Reactor模式核心原理:通过事件循环+I/O多路复用管理海量连接,避免线程模型的调度与内存开销。代码涵盖水平触发、部分读写、缓冲区管理及连接生命周期控制,适合作为高性能网络编程入门参考。(239字)

一个最简单的 TCP 服务器通常为每个客户端创建一个线程,线程阻塞在 acceptrecvsend 上。这种写法容易理解,但连接数增加后会带来线程栈、上下文切换、锁竞争和调度开销。线程数量与连接数量绑定,也使超时、广播和连接回收变得复杂。

Reactor 模式采用另一种组织方式:I/O 多路复用器负责观察大量文件描述符,事件循环在某个或多个就绪事件到达后,再调用对应连接的处理逻辑。它并不会让单个连接的网络操作变快,而是减少了无效等待,使有限的执行线程能够管理更多连接。

本文实现一个基于 Linux epoll 的非阻塞 TCP 回显服务器。它只承担演示职责,不包含 TLS、认证、协议解析和生产级限流;这些能力应在明确需求后单独设计。

核心原理

Reactor 的职责划分

一个基本 Reactor 通常包含四类对象:

  1. 事件循环:调用 epoll_wait,取得已经就绪的文件描述符。
  2. 监听器:处理监听套接字上的可读事件,并循环调用 accept
  3. 连接对象:保存客户端套接字、输入缓冲区、输出缓冲区和当前关注的事件。
  4. 事件分派器:根据文件描述符找到连接对象,调用读、写或错误处理函数。

关键点是回调函数必须尽量短,并且不能在非阻塞套接字上执行可能长时间等待的操作。否则,事件循环仍然会被一个连接拖住。

水平触发与边缘触发

epoll 支持水平触发和边缘触发。

水平触发模式下,只要接收缓冲区仍有数据,后续的 epoll_wait 仍可能报告可读事件,代码更容易验证。边缘触发模式只在状态从“无数据”变为“有数据”等变化时通知,因此处理函数必须一直读到 EAGAIN,否则剩余数据可能没有新的通知。

本文使用水平触发。生产环境若改用 EPOLLET,必须同时满足以下条件:套接字设置为非阻塞;accept 循环直到 EAGAINrecv 循环直到 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,写缓冲区非空时才关注可写事件,关闭路径统一注销并释放资源。

在此基础上,下一步应补充协议解析、缓冲区上限、超时、优雅退出、指标和并发测试。只有这些边界被明确处理,事件循环才具备从示例走向长期运行服务的基础。

相关文章
|
1月前
|
缓存 人工智能 监控
Qwen3.8-Max 深度使用实战:从 2.4 万亿参数到生产级智能体落地
Qwen3.8-Max 是阿里云通义千问 2026 年 8 月最新发布的旗舰基座模型,2.4 万亿参数 MoE 架构、1M 上下文窗口、原生多模态(文本+图像+视频),具备"自主编程十数天交付完整项目"的长程闭环能力。本文不是又一篇"怎么调 API"的入门教程,而是一线团队将 Qwen3.8-Max 从 PoC 推向生产的深度实践记录:百炼平台开通与 API Key 管理、OpenAI 兼容协议接入、多模态与 Function Calling 进阶、思考模式与上下文缓存调优、Token Plan 订阅选型、生产环境避坑实录。
|
1月前
|
数据采集 人工智能 算法
为什么你的品牌在AI里查无此人?GEO优化的五个关键动作
本文为技术实践分享,介绍AI搜索时代企业亟需的GEO(生成式引擎优化)——不同于SEO,GEO聚焦让AI“认识、信任并推荐”品牌。文章提炼罗小军提出的五大关键动作:诊断AI可见度、结构化知识资产、建设权威信源、布局场景词矩阵、建立持续监测机制,助力企业抢占AI决策入口。(239字)
|
1月前
|
数据采集 人工智能 数据挖掘
他山科研 Skills 上架 Qoder:18 个 Skills,覆盖从调研到答辩全流程
他山团队推出20项AI科研Skill,覆盖文献检索、实验设计、学术写作到论文审查全流程,已在Qoder CN技能市场上线。含论文检索、深度研究、假设生成、统计分析、科研绘图等高频工具,助力科研提效。
283 0
|
1月前
|
人工智能 自然语言处理 数据可视化
阿里千问办公上线!一句话产出 PPT、数据报表、完整网页,注册送2000积分
千问办公是阿里巴巴推出的AI原生办公平台,基于Qwen3.8大模型,支持一句话交付PPT、文档、视频、网页等成果;深度集成钉钉,打通本地文件与浏览器,覆盖桌面端、网页端及企业协同场景。千问办公官网:https://t.aliyun.com/U/805n7O
398 0
|
1月前
|
Web App开发 iOS开发 Docker
在局域网自建跨设备文件传输站:PairDrop、WebRTC 与 HTTPS 部署实践
PairDrop是一款基于WebRTC的跨平台浏览器文件传输工具,无需安装客户端,支持iPhone与Windows在局域网内直传照片等文件。它通过信令服务器协助设备发现,利用WebRTC DataChannel实现点对点传输,兼顾便捷性与安全性,适合家庭或小型办公场景临时交换文件。(239字)
317 1
|
1月前
|
数据采集 人工智能 搜索推荐
4层进阶路径:从Schema到交叉印证的内容重构方法
本文提出AI时代内容可见度提升的4层进阶路径:基础标记、答案前置与信息块化、平台分发适配、数据交叉验证迭代。核心是将内容重构为AI可直接摘录的高密度信息块,而非传统文章。
95 2
|
1月前
|
存储 人工智能 自然语言处理
企业网盘和AI知识库打通后,企业知识管理发生了什么变化?
企业网盘与AI知识库打通,彻底改变知识管理:从“存文件”升级为“懂内容、连关系、主动服务”。支持自然语言搜索、跨文档智能关联、语义精准检索,并保障数据隔离、权限继承与全链路审计。知识不再沉睡,真正活起来。(239字)
79 1
|
1月前
|
存储 人工智能 数据处理
基于 YOLO11 的车辆品牌 Logo 检测:从数据标注到云上训练工程化实践
本文介绍基于YOLO11的车辆品牌Logo检测工程实践:涵盖31类、5697张图像的数据集构建,从Label Studio标注、云上OSS存储与版本管理,到YOLO11训练调优、小目标增强及多维度评估,实现端到端可复现、可迭代的AI落地流程。(239字)
基于 YOLO11 的车辆品牌 Logo 检测:从数据标注到云上训练工程化实践
|
1月前
|
安全 算法 测试技术
医院电子病历越权与下载攻防全解析:从水平越权、垂直提权到JWT防线与目录穿越
本文以医院电子病历系统为案例,完整复盘一次安全测试:从水平越权(篡改ID查他人病历)、垂直越权(伪造角色提权)、JWT鉴权原理,到目录穿越与文件名枚举攻击及层层防御演进,揭示越权本质是“客户端身份不可信”的悖论,强调服务端校验、签名防护与随机化命名等关键防御实践。(239字)
100 2
医院电子病历越权与下载攻防全解析:从水平越权、垂直提权到JWT防线与目录穿越
|
1月前
|
人工智能 Android开发 iOS开发
阿里云JVS Claw官网:免费体验7天,不用下载就能用!阿里版OpenClaw龙虾AI助手
阿里云JVS Claw是基于OpenClaw打造的AI龙虾助手,免配置、开箱即用,支持网页/手机/电脑多端访问(下载或不下载均可)。官网直达:https://t.aliyun.com/U/2gB4Og 免费体验7天,付费套餐低至29元/月。
547 0