news 2026/8/6 14:15:05

从零实现一个最小可用的 Reactor:用 epoll 拆开连接、事件与业务逻辑

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
从零实现一个最小可用的 Reactor:用 epoll 拆开连接、事件与业务逻辑

从零实现一个最小可用的 Reactor:用 epoll 拆开连接、事件与业务逻辑

问题背景

一个阻塞式 TCP 服务器通常从accept得到连接,然后在当前线程中调用read。只要某个客户端迟迟不发送完整数据,线程就可能停在读取操作上。为每个连接创建线程能够绕开这个问题,但连接数增加后,线程栈、上下文切换和生命周期管理也会成为额外负担。

Reactor 的思路不是让线程依次等待每个连接,而是把等待工作交给操作系统:程序先注册自己关心的文件描述符和事件,等内核报告“某个描述符已经可以读或写”后,再调用对应处理函数。这样,一个事件循环就可以管理多个连接。

不过,仅仅把epoll_wait写进main并不等于形成了清晰的 Reactor。一个可维护的实现至少要回答四个问题:

  1. 谁负责注册、修改和删除事件?
  2. 文件描述符与连接状态如何关联?
  3. 一次readwrite没有处理完数据怎么办?
  4. 回调执行期间连接被关闭,如何避免继续访问失效对象?

下面实现一个最小的 TCP 回显服务器。它不追求生产级功能,而是把事件循环、连接抽象和数据收发边界完整串起来,适合作为理解 Reactor 的实验骨架。

Reactor 的核心分工

这个实现分为三类对象:

  • Reactor:持有 epoll 实例,负责事件注册和事件分派。
  • Acceptor:监听新连接,把已连接套接字设置为非阻塞,并创建Connection
  • Connection:保存单条连接的输入、输出状态,处理可读、可写和错误事件。

操作系统只认识文件描述符及其事件,不认识 C++ 对象。因此需要建立fd -> EventHandler的映射。事件到达后,Reactor 根据文件描述符找到处理对象,再调用onEvent。这里使用shared_ptr管理处理对象,映射持有对象的所有权;删除映射后,如果当前分派过程还保留一份局部引用,对象会在回调结束后再析构,从而降低回调中关闭连接导致悬空访问的风险。

为什么必须使用非阻塞套接字

事件就绪只表示某次操作现在“有机会推进”,并不承诺业务需要的全部数据已经到齐。例如,可读事件到达后,第一次recv可能只读到部分请求;继续读取时也可能得到EAGAIN。如果套接字仍是阻塞模式,事件循环就可能卡在某个连接上,其他已就绪连接无法得到处理。

因此 Reactor 通常需要配合非阻塞 I/O:

  • recv > 0:消费已经到达的数据。
  • recv == 0:对端完成发送并关闭连接。
  • recv < 0 && errno == EAGAIN:当前数据已经读完,返回事件循环。
  • send只写出部分数据:保留剩余内容,并订阅可写事件。

水平触发与边缘触发

示例采用 epoll 默认的水平触发模式,没有设置EPOLLET。只要描述符仍然满足条件,后续epoll_wait仍可能报告该事件。这种模式更容易验证,也更适合作为第一版实现。

即使采用水平触发,读写循环仍然要正确处理EINTREAGAIN和部分写入。切换为边缘触发后要求更严格:一次通知中通常需要持续读写到EAGAIN,否则剩余数据可能无法及时触发下一次通知。不要只添加EPOLLET标志,而不检查处理函数是否满足这个约束。

可运行实现

准备一台支持 epoll 的 Linux 环境,并确保编译器支持 C++17。创建reactor_echo.cpp,内容如下:

#include<arpa/inet.h>#include<errno.h>#include<fcntl.h>#include<netinet/in.h>#include<sys/epoll.h>#include<sys/socket.h>#include<unistd.h>#include<array>#include<cstring>#include<iostream>#include<memory>#include<stdexcept>#include<string>#include<unordered_map>classReactor;classEventHandler{public:virtual~EventHandler()=default;virtualintfd()const=0;virtualvoidonEvent(uint32_tevents)=0;};staticvoidsetNonBlocking(intfd){intflags=fcntl(fd,F_GETFL,0);if(flags==-1||fcntl(fd,F_SETFL,flags|O_NONBLOCK)==-1){throwstd::runtime_error("fcntl failed");}}classReactor{public:Reactor(){epfd_=epoll_create1(EPOLL_CLOEXEC);if(epfd_==-1)throwstd::runtime_error("epoll_create1 failed");}~Reactor(){close(epfd_);}voidadd(conststd::shared_ptr<EventHandler>&handler,uint32_tevents){epoll_event ev{};ev.events=events;ev.data.fd=handler->fd();if(epoll_ctl(epfd_,EPOLL_CTL_ADD,handler->fd(),&ev)==-1){throwstd::runtime_error("epoll add failed");}handlers_[handler->fd()]=handler;}voidmodify(intfd,uint32_tevents){epoll_event ev{};ev.events=events;ev.data.fd=fd;if(epoll_ctl(epfd_,EPOLL_CTL_MOD,fd,&ev)==-1){throwstd::runtime_error("epoll modify failed");}}voidremove(intfd){epoll_ctl(epfd_,EPOLL_CTL_DEL,fd,nullptr);handlers_.erase(fd);}voidrun(){std::array<epoll_event,64>events{};while(true){intcount=epoll_wait(epfd_,events.data(),events.size(),-1);if(count==-1){if(errno==EINTR)continue;throwstd::runtime_error("epoll_wait failed");}for(inti=0;i<count;++i){intfd=events[i].data.fd;autoit=handlers_.find(fd);if(it==handlers_.end())continue;autohandler=it->second;handler->onEvent(events[i].events);}}}private:intepfd_=-1;std::unordered_map<int,std::shared_ptr<EventHandler>>handlers_;};classConnection:publicEventHandler{public:Connection(Reactor&reactor,intfd):reactor_(reactor),fd_(fd){}~Connection()override{if(fd_!=-1)close(fd_);}intfd()constoverride{returnfd_;}voidonEvent(uint32_tevents)override{if(events&EPOLLIN)readAvailable();if(closed_)return;if(events&EPOLLOUT)writeAvailable();if(closed_)return;if(events&(EPOLLERR|EPOLLHUP))shutdown();}private:voidreadAvailable(){std::array<char,4096>buffer{};while(true){ssize_t n=recv(fd_,buffer.data(),buffer.size(),0);if(n>0){output_.append(buffer.data(),static_cast<size_t>(n));continue;}if(n==0){peerClosed_=true;break;}if(errno==EINTR)continue;if(errno==EAGAIN||errno==EWOULDBLOCK)break;shutdown();return;}if(!output_.empty()){reactor_.modify(fd_,EPOLLIN|EPOLLOUT|EPOLLRDHUP);}elseif(peerClosed_){shutdown();}}voidwriteAvailable(){while(!output_.empty()){ssize_t n=send(fd_,output_.data(),output_.size(),MSG_NOSIGNAL);if(n>0){output_.erase(0,static_cast<size_t>(n));continue;}if(n==-1&&errno==EINTR)continue;if(n==-1&&(errno==EAGAIN||errno==EWOULDBLOCK))break;shutdown();return;}if(output_.empty()){if(peerClosed_)shutdown();elsereactor_.modify(fd_,EPOLLIN|EPOLLRDHUP);}}voidshutdown(){if(closed_)return;closed_=true;intoldFd=fd_;fd_=-1;reactor_.remove(oldFd);close(oldFd);}Reactor&reactor_;intfd_;boolclosed_=false;boolpeerClosed_=false;std::string output_;};classAcceptor:publicEventHandler{public:Acceptor(Reactor&reactor,uint16_tport):reactor_(reactor){fd_=socket(AF_INET,SOCK_STREAM|SOCK_CLOEXEC,0);if(fd_==-1)throwstd::runtime_error("socket failed");intenabled=1;setsockopt(fd_,SOL_SOCKET,SO_REUSEADDR,&enabled,sizeof(enabled));setNonBlocking(fd_);sockaddr_in address{};address.sin_family=AF_INET;address.sin_addr.s_addr=htonl(INADDR_ANY);address.sin_port=htons(port);if(bind(fd_,reinterpret_cast<sockaddr*>(&address),sizeof(address))==-1){throwstd::runtime_error("bind failed");}if(listen(fd_,SOMAXCONN)==-1){throwstd::runtime_error("listen failed");}}~Acceptor()override{close(fd_);}intfd()constoverride{returnfd_;}voidonEvent(uint32_tevents)override{if(!(events&EPOLLIN))return;while(true){intclient=accept4(fd_,nullptr,nullptr,SOCK_NONBLOCK|SOCK_CLOEXEC);if(client>=0){reactor_.add(std::make_shared<Connection>(reactor_,client),EPOLLIN|EPOLLRDHUP);continue;}if(errno==EINTR)continue;if(errno==EAGAIN||errno==EWOULDBLOCK)break;std::cerr<<"accept failed: "<<std::strerror(errno)<<'\n';break;}}private:Reactor&reactor_;intfd_=-1;};intmain(intargc,char**argv){uint16_tport=8080;if(argc==2)port=static_cast<uint16_t>(std::stoul(argv[1]));try{Reactor reactor;autoacceptor=std::make_shared<Acceptor>(reactor,port);reactor.add(acceptor,EPOLLIN);std::cout<<"listening on 0.0.0.0:"<<port<<'\n';reactor.run();}catch(conststd::exception&ex){std::cerr<<ex.what()<<'\n';return1;}}

编译与验证

使用以下命令编译:

g++-std=c++17-O2-Wall-Wextra-pedanticreactor_echo.cpp-oreactor_echo ./reactor_echo8080

在另一个终端建立连接:

nc127.0.0.18080

输入任意文本并回车,服务器应返回相同字节。还可以并行启动多个客户端,确认某个空闲连接不会阻止其他连接收发:

foriin1234;do(printf'client-%s\n'"$i"|nc-N127.0.0.18080)&donewait

不同nc实现的退出选项可能不同。如果不支持-N,可尝试其帮助信息中用于“标准输入结束后关闭连接”的选项,或者直接交互测试。不要据此假定不同系统上的命令行参数完全一致。

可用strace观察事件循环是否按预期调用 epoll 和套接字系统调用:

strace-f-etrace=epoll_wait,epoll_ctl,accept4,recvfrom,sendto ./reactor_echo8080

系统调用名称及展示形式可能随工具和平台而异,但正常情况下应看到监听描述符被注册、连接被接受,以及epoll_wait在没有事件时阻塞。

实现中的关键边界

部分写入不能丢弃

即使send返回正数,也不表示整个缓冲区都已发送。示例通过output_保存待发送字节,只删除已经成功写出的前缀。缓冲区未清空时继续订阅EPOLLOUT,清空后取消该事件,避免套接字长期可写导致事件循环频繁被唤醒。

std::string::erase(0, n)会移动剩余数据,适合教学示例,但大流量场景下可能产生额外开销。工程实现通常维护读取偏移量,或采用分块缓冲区、环形缓冲区及writev。替换存储结构时,部分写入的语义不能被省略。

对端半关闭不等于立即丢弃输出

recv返回零或收到EPOLLRDHUP,说明对端不再发送数据,但本端可能仍有已经生成、尚未写完的响应。示例用peerClosed_记录状态,先尝试清空输出缓冲区,再关闭连接。如果一看到半关闭就直接销毁对象,最后一段响应可能被截断。

回调应保持短小

当前回显逻辑只复制字节,不执行耗时任务。真实服务若在事件循环中进行数据库查询、磁盘读取或重计算,整个线程仍会被阻塞。常见做法是把耗时任务提交到工作线程,完成后通过eventfd、任务队列或其他线程安全通知机制唤醒 Reactor,再由事件循环更新连接状态。

这也带来新的生命周期问题:工作结果返回时,原连接可能已经关闭,文件描述符甚至可能被系统复用。因此异步任务不能只记住一个整数 fd;还应使用连接对象的弱引用、不可复用的连接编号或代际标识进行校验。

必须设置输出缓冲上限

示例为了聚焦事件模型,没有限制output_大小。如果客户端持续发送但不读取响应,缓冲区会不断增长。在对外服务中,应配置单连接高水位,例如达到上限后暂停订阅EPOLLIN、拒绝新请求或关闭连接,并记录可观测的原因。这是背压机制的一部分,具体阈值需要结合协议消息大小、并发规模和内存预算确定,不能脱离业务给出通用数值。

常见问题

为什么accept也要循环到EAGAIN

一次监听事件可能对应多个已完成握手的连接。循环调用accept4可以取走当前已经排队的连接。这个处理方式也让代码更容易切换到边缘触发模式。遇到EINTR应重试,遇到EAGAIN才表示当前队列已处理完。

为什么不能始终监听EPOLLOUT

大多数正常 TCP 连接在发送缓冲区有空间时都处于可写状态。如果始终监听,epoll_wait可能持续返回可写事件,即使应用没有数据需要发送。正确做法是仅在输出缓冲区非空时开启可写事件,写空后立即取消。

EPOLLERR到达后还需要调用close吗?

需要。错误事件只是通知,文件描述符生命周期仍由应用管理。若要记录具体套接字错误,可在关闭前调用getsockopt(fd, SOL_SOCKET, SO_ERROR, ...)。示例直接关闭,是为了保持主流程清晰。

多线程可以共同调用同一个 Reactor 吗?

当前实现没有同步保护,只适用于单事件循环线程。跨线程调用addmodifyremove会引入映射并发访问和对象生命周期竞争。若要支持跨线程操作,应把变更封装成任务放入线程安全队列,并用eventfd唤醒事件循环,由 Reactor 所在线程统一执行。

这个回显服务器能直接作为生产服务器吗?

不能。它缺少协议帧解析、输出高水位、空闲超时、优雅停机、资源限制、指标采集和完整错误日志。代码展示的是 Reactor 的最小闭环,而不是某种生产能力承诺。扩展时应优先补上缓冲区限制、定时器、信号处理和压力条件下的故障测试。

总结

Reactor 的本质不是某个特定类名,而是一组明确的责任划分:内核负责报告就绪事件,事件循环负责分派,连接对象负责维护协议和收发状态。非阻塞 I/O 让单个连接无法长期占住事件线程,动态订阅可写事件避免无效唤醒,输出缓冲则承接部分写入和背压控制。

完成这个最小实现后,下一步不应急于增加复杂框架,而应沿着真实风险扩展:先加入单连接缓冲上限和空闲超时,再实现长度字段协议解析,最后引入工作线程与跨线程唤醒。每增加一种并发机制,都应重新验证连接关闭、任务返回和文件描述符复用三个边界。只有这些状态转换可解释、可测试,Reactor 才真正从演示代码变成可靠的网络程序基础。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/6 14:14:29

不写一行接口,让 DBeaver 直连你的指标层——背后只用了一个端口

先问一个扎心的问题&#xff1a;你们公司的"月活"到底有几种算法&#xff1f; 我打赌不止一种。运营那份报表按登录去重&#xff0c;产品那份按事件去重&#xff0c;财务对账时又是另一套口径。同一个词&#xff0c;三个数字&#xff0c;开会时谁也说服不了谁——最后…

作者头像 李华
网站建设 2026/8/6 14:14:00

SSL证书详细教程:选型、申请、签发和安装全流程解析

在SSL部署过程中&#xff0c;不少企业因为流程不熟或配置疏忽&#xff0c;遇到了签发失败、站点打不开、浏览器报错等问题。本文国科云结合多年SSL证书服务经验&#xff0c;系统整理SSL证书选型、申报部署、日常维护相关问题&#xff0c;为企业申请安装SSL证书提供借鉴参考。一…

作者头像 李华
网站建设 2026/8/6 14:13:51

新品还没到货,如何用易元 AI 提前制作商品宣传图

一、新品未到货&#xff0c;无法拍摄宣传图阻碍预热运营不少电商商家会开启新品预售模式&#xff0c;商品样品还未入库、没有实物可供拍摄。没有实拍图&#xff0c;店铺无法上架商品链接&#xff0c;也不能发布短视频、图文种草预热。等待实物到货再拍摄&#xff0c;会压缩预热…

作者头像 李华
网站建设 2026/8/6 14:12:29

2026年企业大模型聚合平台选型指南:5类AI Gateway能力对比与实践评测

摘要随着企业开始规模化使用大模型&#xff0c;AI应用已经从单一模型调用进入多模型协同阶段。企业可能同时使用不同类型的大模型&#xff1a;通用大模型用于企业办公和知识问答&#xff1b;推理模型用于复杂分析&#xff1b;多模态模型用于图片、视频生成&#xff1b;行业模型…

作者头像 李华
网站建设 2026/8/6 14:11:37

Nacos 1.x 到 2.x 客户端升级实战:避坑指南与最佳实践

1. 项目概述&#xff1a;一次必要的“心脏搭桥”手术最近在负责的一个微服务项目里&#xff0c;我们决定将Nacos客户端从经典的1.x版本升级到2.x。这个决定并非一时兴起&#xff0c;而是随着服务规模扩大&#xff0c;1.x客户端在长连接管理、服务发现性能和配置推送效率上逐渐显…

作者头像 李华
网站建设 2026/8/6 14:11:28

终极二维码修复指南:QRazyBox让损坏二维码起死回生的完整教程

终极二维码修复指南&#xff1a;QRazyBox让损坏二维码起死回生的完整教程 【免费下载链接】qrazybox QR Code Analysis and Recovery Toolkit 项目地址: https://gitcode.com/gh_mirrors/qr/qrazybox 你是否曾面对一个模糊不清、部分损坏的二维码束手无策&#xff1f;QR…

作者头像 李华