【Linux网络】从0手写Reactor反应堆(二):完善核心细节——ET非阻塞读写、分层架构与回调机制
🎬 博主简介:
文章目录
- 前言:
- 一. 前置准备:基础组件的补充修改
- 1.1 Poller 与 Reactor 的结构优化
- 1.2 非阻塞工具与 Socket 改造
- 二. 主线一:Listener Recver 完善 —— 接收新连接的全流程
- 2.1 为什么 ET 模式下 accept 必须循环?
- 2.2 accept 的错误码分级处理
- 2.3 新连接的封装:从 fd 到 IOHandler
- 2.4 关键设计:回指指针 _R
- 三. 主线二:IOHandler Recver 完善 ——ET 循环读与应用层缓冲区
- 3.1 ET 模式下的循环读取
- 3.2 应用层缓冲区:解决粘包问题的基础
- 3.3 职责分离:用回调上抛协议处理
- 四. 分层解耦:协议层与业务层的接入
- 4.1 整体分层架构
- 4.2 协议层 Protocol
- 4.3 业务层 Calculator
- 4.4 回调的传递链路
- 五. Sender 发送逻辑:写事件的按需开启
- 5.1 为什么写事件不能常设?
- 5.2 循环发送逻辑
- 5.3 动态事件修改接口
- 六. 入口组装:Main.cpp 完整流程
- 结尾:
前言:
上一篇我们搭好了 Reactor 的整体骨架:完成了 Connection 基类抽象、Listener 与 IOHandler 派生、Poller 多路复用封装、Reactor 事件派发核心。框架能跑起来,但核心逻辑大多是空壳 ——Listener 只占了坑没真正接收连接,IOHandler 只声明了接口没有读写实现,更没有协议解析和业务处理能力。本篇我们沿着两条主线把核心细节补实:一条以Listener 的 Recver为核心,打通 ET 模式下接收新连接的完整流程;另一条以IOHandler 的 Recver+Sender为核心,完善非阻塞循环读写与应用层缓冲区。同时引入协议层与业务层,通过回调机制实现网络 IO 与业务逻辑的彻底解耦,最终让 Reactor 形成完整的请求处理链路。
一. 前置准备:基础组件的补充修改
在完善核心逻辑之前,我们先对底层组件做几处必要的改造,为后续 ET 模式和非阻塞 IO 铺路。
1.1 Poller 与 Reactor 的结构优化
首先给 Poller 的等待接口补充超时日志,方便调试观察;再把 Reactor 的单次事件派发抽成独立的LoopOnce方法,主循环只负责循环调用,结构更清晰。
// Poller.hpp 节选intWaitEvents(structepoll_eventrevs[],intmaxevents,inttimeout){intn=epoll_wait(_epfd,revs,maxevents,timeout);if(n<0){LOG(LogLevel::FATAL)<<"epoll_wait error";}elseif(n==0){LOG(LogLevel::INFO)<<"epoll_wait time out";}returnn;}// Reactor.hpp 节选voidLoopOnce(inttimeout){intn=_epoll->WaitEvents(revs,gnum,timeout);for(inti=0;i<n;i++){intsockfd=revs[i].data.fd;uint32_tevents=revs[i].events;// 统一异常转读写if((events&EPOLLHUP)||(events&EPOLLERR))events=EPOLLIN|EPOLLOUT;// 事件派发if((events&EPOLLIN)&&IsConnectionExists(sockfd))_connections[sockfd]->Recver();if((events&EPOLLOUT)&&IsConnectionExists(sockfd))_connections[sockfd]->Sender();}}voidDisPatcher(){inttimeout=-1;// -1:阻塞等待,无事件不占用CPUwhile(true){DebugPrint();// 调试:打印当前管理的所有fdLoopOnce(timeout);}}这里补充一个容易被忽略的点:timeout 设为 0 是非阻塞轮询,会疯狂打超时日志占满 CPU;设为 - 1 是阻塞等待,有事件才唤醒,这是服务器的常规配置。
1.2 非阻塞工具与 Socket 改造
ET 模式的硬性要求:所有被监听的 fd 必须设置为非阻塞。否则最后一次循环读 / 写时,没有数据会导致进程挂起,整个事件循环卡死。
我们先在 Common.hpp 中封装一个通用的非阻塞设置函数:
// Common.hpp 节选voidSetNonBlcok(intfd){intflags=fcntl(fd,F_GETFL);if(flags<0){LOG(LogLevel::ERROR)<<"fcntl set non block failed, fd: "<<fd;return;}fcntl(fd,F_SETFL,flags|O_NONBLOCK);}然后改造 TcpSocket:创建套接字时自动设置非阻塞;同时修改 Accepter 接口,通过输出型参数把 errno 带出来,供上层判断是真错误还是无数据。
// Socket.hpp TcpSocket 节选voidCreateSocketOrDie()override{_sockfd=socket(AF_INET,SOCK_STREAM,0);// ... 错误处理SetNonBlcok(_sockfd);// 创建即设为非阻塞// ... 设置地址端口复用}intAccepter(InetAddr*clientaddress,int*code)override{structsockaddr_inpeer;socklen_t len=sizeof(peer);intsockfd=accept(_sockfd,CONV(&peer),&len);*code=errno;// 带出错误码if(sockfd<0)return-1;*clientaddress=peer;returnsockfd;}二. 主线一:Listener Recver 完善 —— 接收新连接的全流程
基础工作做好了,我们从 Listener 的读事件处理开始,把接收新连接的逻辑补全。这是整个 Reactor 接收客户端的入口。
2.1 为什么 ET 模式下 accept 必须循环?
很多同学刚接触 ET 时会写错:accept 只调用一次。在 LT 模式下没问题,因为没处理完的连接会一直通知;但 ET 模式只在状态变化时通知一次 —— 如果同一时间有 10 个客户端完成三次握手,只 accept 一次就会漏掉 9 个,而且再也不会收到通知。
所以正确做法是:while 循环 accept,直到返回 - 1 且错误码为 EAGAIN,表示全连接队列已经空了。
2.2 accept 的错误码分级处理
accept 返回 - 1 不代表真的出错了,我们要根据 errno 区分处理:
// Listener.hpp Recver 节选voidRecver()override{LOG(LogLevel::INFO)<<"Listener event ready, sockfd: "<<_listensockfd->Socketfd();while(true){interrcode=0;InetAddr clientaddr;intsockfd=_listensockfd->Accepter(&clientaddr,&errcode);if(sockfd>=0){// 拿到新连接,后续封装处理}else{if(errcode==EAGAIN||errcode==EWOULDBLOCK){LOG(LogLevel::INFO)<<"accept finish, no more connections";break;// 没有新连接了,本轮结束}elseif(errcode==EINTR){continue;// 被信号中断,重试}else{LOG(LogLevel::ERROR)<<"accept error";break;// 真错误,退出}}}}2.3 新连接的封装:从 fd 到 IOHandler
拿到新的 sockfd 后,绝对不能直接用它 recv/send。按照 Reactor 的设计思想,每一个 fd 都要封装成 Connection 对象,交由 Reactor 统一管理。
步骤很清晰:
- 给新 fd 设置非阻塞(accept 返回的 fd 默认继承监听套接字的非阻塞属性,但显式设置更稳妥);
- 构造 IOHandler 对象,传入 fd 和业务回调;
- 设置该连接关心的事件:EPOLLIN | EPOLLET;
- 填充客户端地址信息;
- 加入 Reactor,注册到内核 epoll。
// Listener.hpp Recver 成功分支if(sockfd>=0){LOG(LogLevel::INFO)<<"accept success, new sockfd: "<<sockfd;SetNonBlcok(sockfd);std::shared_ptr<Connection>conn=std::make_shared<IOHandler>(sockfd,_on_Message);conn->SetEvents(EPOLLIN|EPOLLET);conn->SetClientAddress(clientaddr);_R->AddConnection(conn);// 加入反应堆}2.4 关键设计:回指指针 _R
上面代码里的_R是个很巧妙的设计,这里单独讲一下。
问题:Listener 要把新连接加入 Reactor,就需要调用 Reactor 的 AddConnection。但 Reactor 包含 Connection,如果 Connection 再包含 Reactor,就会出现循环头文件依赖,编译不通过。
解决方案:
- 在 Connection.hpp 中前向声明
class Reactor,只声明不包含头文件; - 给 Connection 增加一个公有成员
Reactor* _R,用原生指针指向所属的 Reactor; - Reactor 在 AddConnection 时,反向给 conn->_R 赋值。
// Connection.hpp 节选classReactor;// 前向声明classConnection{public:Connection():_events(0),_R(nullptr){}// ... 其他接口public:Reactor*_R;// 回指指针};// Reactor.hpp AddConnection 节选voidAddConnection(std::shared_ptr<Connection>&conn){intsockfd=conn->Sockfd();_epoll->AddEvents(sockfd,conn->Events());_connections[sockfd]=conn;conn->_R=this;// 回指赋值}为什么用原生指针不用智能指针?因为 Reactor 用 shared_ptr 管理 Connection,如果 Connection 再用 shared_ptr 指回 Reactor,就会形成循环引用,导致引用计数永远无法归零,内存泄漏。用原生指针只做访问,不管理生命周期,是最简单稳妥的方案。
补充一个编译坑:即便有前向声明,在调用
_R->AddConnection的地方,编译器也必须看到 Reactor 的完整定义。所以头文件包含顺序要注意:Main 里先包含 Reactor.hpp,再包含 Listener.hpp,否则会报invalid use of incomplete type错误。
三. 主线二:IOHandler Recver 完善 ——ET 循环读与应用层缓冲区
接收到的新连接最终都会走到 IOHandler,它负责真正的数据收发。读逻辑和 accept 的 ET 处理思路高度一致。
3.1 ET 模式下的循环读取
同样的道理:ET 模式下读事件只通知一次,必须循环调用 recv,把内核接收缓冲区的数据全部读完,直到返回 EAGAIN。
// IOHandler.hpp Recver 节选voidRecver()override{LOG(LogLevel::INFO)<<"IOHandler event ready, sockfd: "<<_sockfd;charbuffer[gbuffersize];while(true){intn=recv(_sockfd,buffer,sizeof(buffer)-1,0);if(n>0){buffer[n]=0;_inbuffer+=buffer;// 追加到输入缓冲区}elseif(n==0){LOG(LogLevel::INFO)<<"client quit, addr: "<<_clientaddr.StringAddress();Excepter();return;// 对端关闭,直接结束函数}else{if(errno==EAGAIN||errno==EWOULDBLOCK){break;// 数据读完了}elseif(errno==EINTR){continue;// 被信号打断,重试}else{LOG(LogLevel::ERROR)<<"recv error, sockfd: "<<_sockfd;Excepter();return;}}}// ... 报文处理, 就在下面}注意两个细节:
- 对端关闭和读错误时用return而不是 break,因为出错了就没必要继续后面的业务处理,直接退出函数;
- 数据不是处理完就丢,而是追加到
_inbuffer里,这就是应用层接收缓冲区。
3.2 应用层缓冲区:解决粘包问题的基础
TCP 是字节流协议,没有报文边界,一次 recv 不一定读到完整报文。如果 buffer 是局部变量,函数返回数据就丢了,根本没法处理半包。
每个连接独立的_inbuffer就是解决方案:没读完的、不完整的报文都暂存在里面,下次读到新数据再拼接,直到凑齐完整报文再处理。这就是 “先接收,再解包” 的思路。
3.3 职责分离:用回调上抛协议处理
IOHandler 只负责数据的读取和发送,不应该关心报文怎么解析、业务怎么处理。如果把 JSON 解析、加减乘除都写在 IOHandler 里,代码又会耦合回原生 epoll 的样子。
我们用回调函数把协议处理上抛:IOHandler 只负责把读到的缓冲区交给回调,处理完拿回应答结果,放到发送缓冲区里。
// 回调类型定义usingOnMessage_t=std::function<std::string(std::string&inbuffer,int*code)>;// IOHandler 读取完成后intcode=0;std::string result=_on_Message(_inbuffer,&code);if(code==0){_outbuffer+=result;// 应答放入输出缓冲区}else{Excepter();return;}// 有数据就尝试发送if(!_outbuffer.empty())Sender();这样 IOHandler 就保持了单一职责:只做 IO,不碰业务。
四. 分层解耦:协议层与业务层的接入
有了回调机制,我们就可以在 IO 层之上,再搭建协议层和业务层,形成清晰的三层架构。
4.1 整体分层架构
从上到下:
- 业务层(Calculator):纯业务计算,完全不知道网络和协议的存在;
- 协议层(Protocol):负责封包、解包、序列化、反序列化,解决粘包;
- IO 层(IOHandler/Reactor):负责非阻塞数据收发、事件派发。
层与层之间通过回调交互,下层不知道上层的具体实现,只知道接口规范,完美符合开闭原则。
4.2 协议层 Protocol
Protocol 的核心工作就是处理字节流:
- 解包 UnPack:从字节流中按 “长度 + 分隔符” 提取完整的 JSON 报文;
- 封包 Pack:给 JSON 报文加上长度头和分隔符,发往网络;
- HandlerRequest:循环解包 → 反序列化 → 调用业务回调 → 序列化封包 → 返回所有应答。
// Protocol.hpp HandlerRequest 核心逻辑std::stringHandlerRequest(std::string&streamstr,int*code){std::string resp_package;while(true){std::string jsonstring;intn=UnPack(streamstr,&jsonstring);if(n==0){// 报文不完整,等下次*code=0;returnresp_package;}elseif(n==-1){// 协议损坏*code=-1;exit(1);// 直接断开不守规矩的客户端}Request request;request.Deserialize(jsonstring);// 反序列化Response response=_cb(request);// 调用业务回调std::string respjsonstr;response.Serialize(&respjsonstr);// 序列化resp_package+=Pack(respjsonstr);// 封包拼接}}这里用 while 循环也是为了处理粘包:一次读到的字节流里可能包含多个完整报文,要全部处理完。
4.3 业务层 Calculator
Calculator 是最纯粹的一层:输入 Request,输出 Response,只做加减乘除计算。它完全感知不到网络、epoll、JSON,甚至不知道自己运行在 Reactor 里。
// Calculator.hpp Exec 节选ResponseExec(constRequest&req){Response resp;switch(req._oper){case'+':resp._result=req._x+req._y;break;case'-':resp._result=req._x-req._y;break;case'*':resp._result=req._x*req._y;break;case'/':if(req._y==0)resp._exitcode=-1;// 防御性编程:防止除零崩溃elseresp._result=req._x/req._y;break;// ... 其他操作符default:resp._exitcode=3;break;}returnresp;}这种设计的好处是:业务逻辑可以单独测试、单独替换,比如明天想把计算器换成聊天服务,只需要换个业务回调,网络层和协议层一行都不用改。
4.4 回调的传递链路
回调是怎么从 Main 一路传到 IOHandler 的?我们梳理一下:
- Main 中创建 Calculator,把计算函数作为回调传给 Protocol;
- Main 中把 Protocol 的请求处理函数作为回调传给 Listener;
- Listener 每次 accept 到新连接,把回调传给新创建的 IOHandler;
- IOHandler 读到数据后,调用回调,一路向上回到业务层。
整个过程像搭积木一样逐层组装,没有任何硬编码依赖。
五. Sender 发送逻辑:写事件的按需开启
讲完读,我们再看写。发送逻辑比读取多了一个非常关键的设计点:写事件不能常设。
5.1 为什么写事件不能常设?
很多同学会想当然地给 fd 同时加上 EPOLLIN 和 EPOLLOUT,这是典型的错误。
- 读事件:默认不满足(内核缓冲区没数据),所以 EPOLLIN 常设,等数据来了再通知,没问题;
- 写事件:默认满足(内核发送缓冲区为空),如果常设 EPOLLOUT,epoll 会一直返回写就绪,事件循环疯狂空转,CPU 直接打满 100%。
最佳实践:EPOLLOUT 按需开启。默认只开 EPOLLIN;只有当发送缓冲区满了、数据没发完的时候,才开启 EPOLLOUT,等下次写就绪了继续发;发完了立刻关闭 EPOLLOUT。
5.2 循环发送逻辑
和读一样,ET 模式下写也要循环,直到把数据发完或者发送缓冲区满。
// IOHandler.hpp Sender 节选voidSender()override{while(true){intn=send(_sockfd,_outbuffer.c_str(),_outbuffer.size(),0);if(n>=0){_outbuffer.erase(0,n);// 移除已发送的部分if(_outbuffer.empty())break;}else{if(errno==EAGAIN||errno==EWOULDBLOCK){break;// 发送缓冲区满了,写条件不满足}elseif(errno==EINTR){continue;// 被信号打断,重试}else{LOG(LogLevel::ERROR)<<"send error, sockfd: "<<_sockfd;Excepter();return;}}}// 根据发送结果控制写事件开关if(_outbuffer.empty())_R->EnableReadWrite(_sockfd,true,false);// 发完了,关闭写事件else_R->EnableReadWrite(_sockfd,true,true);// 没发完,开启写事件}5.3 动态事件修改接口
上面的EnableReadWrite是 Reactor 提供的接口,作用是动态修改某个 fd 在 epoll 中关心的事件。目前我们先把接口留出来,具体的epoll_ctl MOD实现会在下一篇完善,顺带完成连接的释放与生命周期管理。
六. 入口组装:Main.cpp 完整流程
最后我们看 Main 函数,把所有层像搭积木一样拼起来,就能直观感受到分层架构的清晰:
// Main.cppintmain(){// 1. 业务层:计算器std::unique_ptr<Calculator>cal=std::make_unique<Calculator>();// 2. 协议层:绑定业务回调std::unique_ptr<Protocol>protocol=std::make_unique<Protocol>([&cal](constRequest&req){returncal->Exec(req);});// 3. 监听连接:绑定协议回调,设置ET读事件std::shared_ptr<Connection>connection=std::make_shared<Listener>(gport,[&protocol](std::string&inbuffer,int*code){returnprotocol->HandlerRequest(inbuffer,code);});connection->SetEvents(EPOLLIN|EPOLLET);// 4. 反应堆std::unique_ptr<Reactor>reactor=std::make_unique<Reactor>();reactor->AddConnection(connection);// 5. 启动事件循环reactor->DisPatcher();return0;}总结
到这里,Reactor 的核心细节就完善得差不多了。我们回顾一下本篇的核心内容:
- ET 模式适配:监听套接字和普通 IO 套接字都实现了非阻塞循环读写,正确处理 EAGAIN、EINTR 等错误码,符合 ET 工作规范;
- 新连接闭环:Listener 从 accept 到封装 IOHandler,再加入 Reactor 管理,整条链路打通;
- 应用层缓冲区:每个连接独立的输入输出缓冲区,为解决粘包半包提供了基础;
- 三层解耦架构:IO 层、协议层、业务层通过回调分离,职责单一,扩展性极强;
- 写事件按需开启:明确了读写事件的不同处理策略,避免 CPU 空转。
结尾:
🍓 我是草莓熊 Lotso!若这篇技术干货帮你打通了学习中的卡点: 👀 【关注】跟我一起深耕技术领域,从基础到进阶,见证每一次成长 ❤️ 【点赞】让优质内容被更多人看见,让知识传递更有力量 ⭐ 【收藏】把核心知识点、实战技巧存好,需要时直接查、随时用 💬 【评论】分享你的经验或疑问(比如曾踩过的技术坑?),一起交流避坑 🗳️ 【投票】用你的选择助力社区内容方向,告诉大家哪个技术点最该重点拆解 技术之路难免有困惑,但同行的人会让前进更有方向~愿我们都能在自己专注的领域里,一步步靠近心中的技术目标!结语:当然目前还有收尾工作没完成:EnableReadWrite接口只有声明没有实现,无法动态修改 epoll 事件;Excepter只打了日志,还没有真正的连接移除、fd 关闭、资源释放逻辑等;下一篇我们会重点完善连接的生命周期管理:实现 EnableReadWrite、完善异常处理、连接安全移除,让整个 Reactor 的连接从创建到销毁形成完整闭环,并且会最终展示完整的代码。
✨把这些内容吃透超牛的!放松下吧✨ ʕ˘ᴥ˘ʔ づきらど