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

📅 2026/8/6 14:15:07
从零实现一个最小可用的 Reactor:用 epoll 拆开连接、事件与业务逻辑
从零实现一个最小可用的 Reactor用 epoll 拆开连接、事件与业务逻辑问题背景一个阻塞式 TCP 服务器通常从accept得到连接然后在当前线程中调用read。只要某个客户端迟迟不发送完整数据线程就可能停在读取操作上。为每个连接创建线程能够绕开这个问题但连接数增加后线程栈、上下文切换和生命周期管理也会成为额外负担。Reactor 的思路不是让线程依次等待每个连接而是把等待工作交给操作系统程序先注册自己关心的文件描述符和事件等内核报告“某个描述符已经可以读或写”后再调用对应处理函数。这样一个事件循环就可以管理多个连接。不过仅仅把epoll_wait写进main并不等于形成了清晰的 Reactor。一个可维护的实现至少要回答四个问题谁负责注册、修改和删除事件文件描述符与连接状态如何关联一次read或write没有处理完数据怎么办回调执行期间连接被关闭如何避免继续访问失效对象下面实现一个最小的 TCP 回显服务器。它不追求生产级功能而是把事件循环、连接抽象和数据收发边界完整串起来适合作为理解 Reactor 的实验骨架。Reactor 的核心分工这个实现分为三类对象Reactor持有 epoll 实例负责事件注册和事件分派。Acceptor监听新连接把已连接套接字设置为非阻塞并创建Connection。Connection保存单条连接的输入、输出状态处理可读、可写和错误事件。操作系统只认识文件描述符及其事件不认识 C 对象。因此需要建立fd - EventHandler的映射。事件到达后Reactor 根据文件描述符找到处理对象再调用onEvent。这里使用shared_ptr管理处理对象映射持有对象的所有权删除映射后如果当前分派过程还保留一份局部引用对象会在回调结束后再析构从而降低回调中关闭连接导致悬空访问的风险。为什么必须使用非阻塞套接字事件就绪只表示某次操作现在“有机会推进”并不承诺业务需要的全部数据已经到齐。例如可读事件到达后第一次recv可能只读到部分请求继续读取时也可能得到EAGAIN。如果套接字仍是阻塞模式事件循环就可能卡在某个连接上其他已就绪连接无法得到处理。因此 Reactor 通常需要配合非阻塞 I/Orecv 0消费已经到达的数据。recv 0对端完成发送并关闭连接。recv 0 errno EAGAIN当前数据已经读完返回事件循环。send只写出部分数据保留剩余内容并订阅可写事件。水平触发与边缘触发示例采用 epoll 默认的水平触发模式没有设置EPOLLET。只要描述符仍然满足条件后续epoll_wait仍可能报告该事件。这种模式更容易验证也更适合作为第一版实现。即使采用水平触发读写循环仍然要正确处理EINTR、EAGAIN和部分写入。切换为边缘触发后要求更严格一次通知中通常需要持续读写到EAGAIN否则剩余数据可能无法及时触发下一次通知。不要只添加EPOLLET标志而不检查处理函数是否满足这个约束。可运行实现准备一台支持 epoll 的 Linux 环境并确保编译器支持 C17。创建reactor_echo.cpp内容如下#includearpa/inet.h#includeerrno.h#includefcntl.h#includenetinet/in.h#includesys/epoll.h#includesys/socket.h#includeunistd.h#includearray#includecstring#includeiostream#includememory#includestdexcept#includestring#includeunordered_mapclassReactor;classEventHandler{public:virtual~EventHandler()default;virtualintfd()const0;virtualvoidonEvent(uint32_tevents)0;};staticvoidsetNonBlocking(intfd){intflagsfcntl(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_ptrEventHandlerhandler,uint32_tevents){epoll_event ev{};ev.eventsevents;ev.data.fdhandler-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.eventsevents;ev.data.fdfd;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::arrayepoll_event,64events{};while(true){intcountepoll_wait(epfd_,events.data(),events.size(),-1);if(count-1){if(errnoEINTR)continue;throwstd::runtime_error(epoll_wait failed);}for(inti0;icount;i){intfdevents[i].data.fd;autoithandlers_.find(fd);if(ithandlers_.end())continue;autohandlerit-second;handler-onEvent(events[i].events);}}}private:intepfd_-1;std::unordered_mapint,std::shared_ptrEventHandlerhandlers_;};classConnection:publicEventHandler{public:Connection(Reactorreactor,intfd):reactor_(reactor),fd_(fd){}~Connection()override{if(fd_!-1)close(fd_);}intfd()constoverride{returnfd_;}voidonEvent(uint32_tevents)override{if(eventsEPOLLIN)readAvailable();if(closed_)return;if(eventsEPOLLOUT)writeAvailable();if(closed_)return;if(events(EPOLLERR|EPOLLHUP))shutdown();}private:voidreadAvailable(){std::arraychar,4096buffer{};while(true){ssize_t nrecv(fd_,buffer.data(),buffer.size(),0);if(n0){output_.append(buffer.data(),static_castsize_t(n));continue;}if(n0){peerClosed_true;break;}if(errnoEINTR)continue;if(errnoEAGAIN||errnoEWOULDBLOCK)break;shutdown();return;}if(!output_.empty()){reactor_.modify(fd_,EPOLLIN|EPOLLOUT|EPOLLRDHUP);}elseif(peerClosed_){shutdown();}}voidwriteAvailable(){while(!output_.empty()){ssize_t nsend(fd_,output_.data(),output_.size(),MSG_NOSIGNAL);if(n0){output_.erase(0,static_castsize_t(n));continue;}if(n-1errnoEINTR)continue;if(n-1(errnoEAGAIN||errnoEWOULDBLOCK))break;shutdown();return;}if(output_.empty()){if(peerClosed_)shutdown();elsereactor_.modify(fd_,EPOLLIN|EPOLLRDHUP);}}voidshutdown(){if(closed_)return;closed_true;intoldFdfd_;fd_-1;reactor_.remove(oldFd);close(oldFd);}Reactorreactor_;intfd_;boolclosed_false;boolpeerClosed_false;std::string output_;};classAcceptor:publicEventHandler{public:Acceptor(Reactorreactor,uint16_tport):reactor_(reactor){fd_socket(AF_INET,SOCK_STREAM|SOCK_CLOEXEC,0);if(fd_-1)throwstd::runtime_error(socket failed);intenabled1;setsockopt(fd_,SOL_SOCKET,SO_REUSEADDR,enabled,sizeof(enabled));setNonBlocking(fd_);sockaddr_in address{};address.sin_familyAF_INET;address.sin_addr.s_addrhtonl(INADDR_ANY);address.sin_porthtons(port);if(bind(fd_,reinterpret_castsockaddr*(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(!(eventsEPOLLIN))return;while(true){intclientaccept4(fd_,nullptr,nullptr,SOCK_NONBLOCK|SOCK_CLOEXEC);if(client0){reactor_.add(std::make_sharedConnection(reactor_,client),EPOLLIN|EPOLLRDHUP);continue;}if(errnoEINTR)continue;if(errnoEAGAIN||errnoEWOULDBLOCK)break;std::cerraccept failed: std::strerror(errno)\n;break;}}private:Reactorreactor_;intfd_-1;};intmain(intargc,char**argv){uint16_tport8080;if(argc2)portstatic_castuint16_t(std::stoul(argv[1]));try{Reactor reactor;autoacceptorstd::make_sharedAcceptor(reactor,port);reactor.add(acceptor,EPOLLIN);std::coutlistening on 0.0.0.0:port\n;reactor.run();}catch(conststd::exceptionex){std::cerrex.what()\n;return1;}}编译与验证使用以下命令编译g-stdc17-O2-Wall-Wextra-pedanticreactor_echo.cpp-oreactor_echo ./reactor_echo8080在另一个终端建立连接nc127.0.0.18080输入任意文本并回车服务器应返回相同字节。还可以并行启动多个客户端确认某个空闲连接不会阻止其他连接收发foriin1234;do(printfclient-%s\n$i|nc-N127.0.0.18080)donewait不同nc实现的退出选项可能不同。如果不支持-N可尝试其帮助信息中用于“标准输入结束后关闭连接”的选项或者直接交互测试。不要据此假定不同系统上的命令行参数完全一致。可用strace观察事件循环是否按预期调用 epoll 和套接字系统调用strace-f-etraceepoll_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 吗当前实现没有同步保护只适用于单事件循环线程。跨线程调用add、modify或remove会引入映射并发访问和对象生命周期竞争。若要支持跨线程操作应把变更封装成任务放入线程安全队列并用eventfd唤醒事件循环由 Reactor 所在线程统一执行。这个回显服务器能直接作为生产服务器吗不能。它缺少协议帧解析、输出高水位、空闲超时、优雅停机、资源限制、指标采集和完整错误日志。代码展示的是 Reactor 的最小闭环而不是某种生产能力承诺。扩展时应优先补上缓冲区限制、定时器、信号处理和压力条件下的故障测试。总结Reactor 的本质不是某个特定类名而是一组明确的责任划分内核负责报告就绪事件事件循环负责分派连接对象负责维护协议和收发状态。非阻塞 I/O 让单个连接无法长期占住事件线程动态订阅可写事件避免无效唤醒输出缓冲则承接部分写入和背压控制。完成这个最小实现后下一步不应急于增加复杂框架而应沿着真实风险扩展先加入单连接缓冲上限和空闲超时再实现长度字段协议解析最后引入工作线程与跨线程唤醒。每增加一种并发机制都应重新验证连接关闭、任务返回和文件描述符复用三个边界。只有这些状态转换可解释、可测试Reactor 才真正从演示代码变成可靠的网络程序基础。