Reactor模式
1.Reactor 模式定义
Reactor 模式(反应器模式,也称 Dispatcher 分发器模式) 是一种事件驱动的设计模式: 它通过 I/O 多路复用(select / poll / epoll)同时监听多个句柄,当某个句柄上的事件就绪时,把事件分派给事先注册好的处理程序(回调)去处理。

2.Reactor模式的角色构成
| 角色 | 解释 |
|---|---|
| Handle(句柄) | 用于标识不同的资源,本质就是一个文件描述符 |
| Synchronous Event Demultiplexer(同步事件分离器) | 本质就是一个系统调用,用于等待事件的发生。对于 Linux 来说,指的就是 I/O 多路复用,比如 select、poll、epoll 等 |
| Event Handler(事件处理器) | 由多个回调方法构成,这些回调方法构成了与应用相关的、对于某个事件的处理反馈 |
| Concrete Event Handler(具体事件处理器) | 事件处理器中各个回调方法的具体实现 |
| Initiation Dispatcher(初始分发器) | 实际上就是 Reactor 角色。它会通过同步事件分离器等待事件发生,当对应事件就绪时调用事件处理器,最后调用对应的回调方法来处理这个事件 |
3.Reactor 模式的工作流程
- 注册:应用向初始分发器注册具体事件处理器,同时声明它关心哪些事件类型(如"可读")。每个事件处理器都关联着一个 Handle。
- 交出 Handle:初始分发器要求每个事件处理器提供它内部的 Handle。这个 Handle 向操作系统标识了要监视的资源;在 Reactor 中通常一个处理器对应一个 Handle,通过 Handle 即可反查处理器。
- 启动事件循环:所有处理器注册完成后,启动事件循环。分发器把所有注册的 Handle 汇成一份监视列表,交给同步事件分离器等待事件发生。
- 事件就绪:某个 Handle 变为 Ready 时,同步事件分离器通知初始分发器(在 Linux 上就是
epoll_wait返回就绪集合)。- 查找处理器:分发器以 Ready 的 Handle 作为 key,找到对应的事件处理器。
- 回调:调用该处理器中对应的回调方法来响应事件。
- 回到:继续下一轮等待 —— 这就是"事件循环"。
4.epoll ET 服务器(Reactor 模式)
(1)设计思路
A. epoll ET 服务器
在实现 ET 模式之前,必须先明确三条前提:
- 所有文件描述符必须设置为非阻塞(fcntl(fd, F_SETFL, O_NONBLOCK));
- 读事件就绪时必须循环读取,直到 recv 返回 EAGAIN;
- 监听套接字的读事件就绪时必须循环 accept,直到返回 EAGAIN。
原因在于 ET(边缘触发)只在状态"由无到有"时通知一次。本次若不把数据读完,内核不会因为"还有剩余数据"而再次通知,剩余数据将永远无法读到;同样,若只 accept 一个连接,剩余排队连接也可能再也收不到通知。
- 读事件:若是监听套接字读事件就绪,则循环调用 accept 函数获取底层连接,直到返回 EAGAIN;若是其他套接字读事件就绪,则循环调用 recv 函数读取客户端发来的数据,直到返回 EAGAIN。其中 recv 返回 0 表示对端已关闭连接,返回 -1 且 errno == EAGAIN 表示本轮数据已读完,两种情况都应退出循环。
- 写事件:写事件就绪表示内核发送缓冲区有空间可写,此时应该从应用层发送缓冲区中取出待发送的数据,调用 send 写入内核发送缓冲区,并循环直到数据发完或返回 EAGAIN。
- 异常事件:当某个套接字的异常事件(EPOLLERR / EPOLLHUP)就绪时不做过多处理,直接关闭该套接字并注销。需要注意的是,这两个事件总会被上报,无需在 events 中主动注册。
当 epoll ET 服务器监测到某一事件就绪后,就会将该事件交给对应的服务处理程序进行处理
B. Reactor 模式的五个角色
在这个 epoll ET 服务器中,Reactor 模式中的五个角色对应如下:
- 句柄:文件描述符。
- 同步事件分离器:I/O 多路复用 epoll,即 epoll_wait 函数。
- 事件处理器:读回调、写回调和异常回调的函数签名,即事件处理的接口定义。
- 具体事件处理器:Connection 对象。它把一个文件描述符与该描述符对应的读、写、异常回调封装在一起,从而在事件就绪时能够通过文件描述符找到并执行相应的回调方法。
- 初始分发器:Reactor 类中的 Dispatcher 函数。Dispatcher 函数的工作即为:调用 epoll_wait 函数等待事件发生,有事件发生后将就绪的事件派发给对应的服务处理程序即可。
C. Connection 类
在 Reactor 的工作流程中说到,在注册事件处理器时需要将其与 Handle 关联,本质上就是需要将读回调、写回调和异常回调与某个文件描述符关联起来。这样做的目的就是为了当某个文件描述符上的事件就绪时可以找到其对应的各种回调函数,进而执行对应的回调方法来处理该事件。
可以设计一个 Connection 类,该类中的成员包括:
- 一个文件描述符;
- 该文件描述符对应的读回调、写回调和异常回调;
- 接收缓冲区与发送缓冲区。ET 模式下一次 recv 可能只读到半条消息,需要缓冲区攒够一条完整报文再处理;send 未发送完时,剩余数据也需要暂存,等写事件就绪后再继续发送;
- 连接状态等其他成员。
D. TcpServer 类
在 Reactor 的工作流程中说到,当所有事件处理器注册完毕后,会使用同步事件分离器等待这些事件发生,当某个事件处理器的 Handle 变为 Ready 状态时,同步事件分离器会通知初始分发器,然后初始分发器会将 Ready 状态的 Handle 作为 key 来寻找其对应的事件处理器,并调用该事件处理器中对应的回调方法来响应该事件。
本质就是当事件注册完毕后,会调用 epoll_wait 函数来等待这些事件发生,当某个事件就绪时 epoll_wait 函数会告知调用方,然后调用方就根据就绪的文件描述符来找到其对应的各种回调函数,并调用对应的回调函数进行事件处理。
对此可以设计一个 Reactor 类:
该类当中有一个成员函数 Dispatcher,即初始分发器,在该函数内部会调用 epoll_wait 函数等待事件的发生,epoll_wait 返回后即可得到本轮所有就绪的事件。
当事件就绪后需要根据就绪的文件描述符来找到其对应的各种回调函数。由于会将每个文件描述符及其对应的各种回调都封装到一个 Connection 结构中,所以可以根据文件描述符找到其对应的 Connection 结构。
使用 C++ STL 中的 unordered_map,来建立各个文件描述符与其对应的 Connection 结构之间的映射。这个 unordered_map 应当作为 Reactor 类的成员变量——因为"按文件描述符查找处理器"正是初始分发器自身的职责,它与 Dispatcher 必须在同一个类中,否则 Dispatcher 将无法完成事件分派。
Reactor 类中还需要提供成员函数 AddConnection,用于向初始分发器中注册事件:调用 epoll_ctl 以 EPOLL_CTL_ADD 将该文件描述符加入内核监视列表,同时把该文件描述符与 Connection 的映射插入 unordered_map。
此外还应提供成员函数 DelConnection,用于注销事件并释放连接,其内部需要依次完成三件事:
- epoll_ctl(_epfd, EPOLL_CTL_DEL, fd, nullptr) 从内核摘除;
- close(fd) 关闭文件描述符;
- _connections.erase(fd) 从映射表中删除。
三者缺一不可,尤其是最后一步——若忘记从映射表中删除,当该文件描述符的编号被新连接复用时,会查找到已经失效的旧 Connection,从而产生极难定位的错误。
E. epoll ET 服务器的工作流程
初始化:创建、绑定、监听套接字,将监听套接字设为非阻塞,创建 epoll 模型。
注册监听套接字:为监听套接字创建对应的 Connection 结构,调用 Reactor 类的 AddConnection 函数将其加入 epoll 模型,并建立"监听套接字 → Connection"的映射关系。
事件循环:调用 Reactor 类的 Dispatcher 函数开始派发。Dispatcher 内部是一个 while 循环,每轮调用 epoll_wait 等待事件,事件就绪后按文件描述符找到对应的 Connection,调用相应回调处理。
动态注册与注销:处理过程中会不断新增和移除被监视的文件描述符:
- 新连接到来 --> accept 循环到 EAGAIN --> 创建 Connection --> AddConnection 注册;
- 连接断开或出错 --> DelConnection 注销(从 epoll 摘除、close、从映射表删除)。
退出:Dispatcher 是死循环,需用 _stop 标志控制退出,退出后统一关闭监听套接字与 epoll 句柄。
(2)Connection 结构
- Connection 结构中除了包含文件描述符和其对应的读回调、写回调和异常回调外,还包含一个输入缓冲区 _inBuffer、一个输出缓冲区 _outBuffer 以及一个回指指针 _svrPtr。
- 当某个文件描述符的读事件就绪时,调用 recv 函数读取客户端发来的数据,但一次 recv 既可能只读到半个报文,也可能读到多个报文,因此需要将读取到的数据暂时存放到该文件描述符对应的 _inBuffer 中,循环从中分离出完整的报文再进行数据处理,_inBuffer 本质就是用来解决粘包与半包问题的。
- 当处理完一个报文请求后,需将响应数据发送给客户端。先把响应数据追加到 _outBuffer,再尝试调用 send 发送;若底层 TCP 发送缓冲区空间不足,未发完的部分就留在 _outBuffer 中,等写事件就绪时再依次发送。
- Connection 结构中设置回指指针 _svrPtr,便于快速找到 TcpServer 对象,因为后续需要根据 Connection 结构找到这个 TcpServer 对象。如上层业务处理函数 NetCal 向 _outBuffer 递交数据后,需通过 Connection 中的回指指针"提醒"上层发起发送。需要注意:真正要做的动作是"注册写事件",而这需要操作 epoll——epoll 由 Reactor 持有,因此这一步应通过 TcpServer 转发给 Reactor 完成。
Connection 结构中需提供一个管理回调的成员函数,便于外部对回调进行设置:
为了能够正常工作,常规的sock必须是要有自己独立的接收缓冲区&&发送缓冲区
class Connection
{
public:
Connection(int sock = -1)
: _sock(sock)
, _tsvr(nullptr)
{}
void SetCallBack(func_t recv_cb, func_t send_cb, func_t except_cb)
{
_recv_cb = recv_cb;
_send_cb = send_cb;
_except_cb = except_cb;
}
~Connection()
{}
};
private:
int _listensock;
int _port;
Epoll _poll;
std::unordered_map<int, Connection *> _connections;
struct epoll_event *_revs;
int _revs_num;
// 这个是上层的业务处理
callback_t _cb;
};
(3)TcpServer 类
TcpServer 类中有两个关键成员:
- 一个 unordered_map,用于建立文件描述符与其对应的 Connection 结构之间的映射;
- 一个 Epoll 对象成员 _epoll,它是同步事件分离器的封装。 由于 _epoll 是值成员,Epoll 对象会随 TcpServer 的构造自动创建。
在 TcpServer 的构造函数中调用 Epoll 类提供的 EpollCreate 函数,创建内核中的 epoll 实例,并把 epoll_create 返回的文件描述符记录到 Epoll 对象的成员变量 _epollFd 中,便于后续 epoll_ctl 与 epoll_wait 使用。
TcpServer 对象析构时,_epoll 作为成员会自动析构,其析构函数中调用 close 关闭 epoll 模型对应的文件描述符,无需手动处理。
此外,析构时还需要释放 _connections 中所有的 Connection 对象——该映射表存放的是裸指针,不会自动释放,应逐个 delete,否则每处理一个连接就泄漏一个对象。
A. 封装 Epoll 类
Epoll.hpp
#pragma once
#include <iostream>
#include <sys/epoll.h>
class Epoll
{
const static int gnum = 128;
const static int gtimeout = 5000;
public:
Epoll(int timeout = gtimeout) : _timeout(timeout)
{}
void CreateEpoll()
{
_epfd = epoll_create(gnum);
if(_epfd < 0) exit(5);
}
bool DelFromEpoll(int sock)
{
int n = epoll_ctl(_epfd, EPOLL_CTL_DEL, sock, nullptr);
return n == 0;
}
bool CtrlEpoll(int sock, uint32_t events)
{
events |= EPOLLET;
struct epoll_event ev;
ev.events = events;
ev.data.fd = sock;
int n = epoll_ctl(_epfd, EPOLL_CTL_MOD, sock, &ev);
return n == 0;
}
bool AddSockToEpoll(int sock, uint32_t events)
{
struct epoll_event ev;
ev.events = events;
ev.data.fd = sock;
int n = epoll_ctl(_epfd, EPOLL_CTL_ADD, sock, &ev);
return n == 0;
}
int WaitEpoll(struct epoll_event revs[], int num)
{
return epoll_wait(_epfd, revs, num, _timeout);
}
~Epoll()
{}
private:
int _epfd;
int _timeout;
};
B. TcpServer 类部分代码
// 这个网络服务器坚决不要和上层业务强耦合
class TcpServer
{
const static int gport = 8080;
const static int gnum = 128;
public:
TcpServer(int port = gport) : _port(port), _revs_num(gnum)
{
// 1. 创建listensock
_listensock = Sock::Socket();
Sock::Bind(_listensock, _port);
Sock::Listen(_listensock);
// 2. 创建多路转接对象
_poll.CreateEpoll();
// 3. 添加listensock到服务器中
AddConnection(_listensock,
std::bind(&TcpServer::Accepter, this, std::placeholders::_1),
nullptr, nullptr);
// 4. 构建一个获取就绪事件的缓冲区
_revs = new struct epoll_event[_revs_num];
}
~TcpServer()
{
if (_listensock >= 0) close(_listensock);
if (_revs) delete[] _revs;
}
private:
int _listensock;
int _port;
Epoll _poll;
std::unordered_map<int, Connection *> _connections;
struct epoll_event *_revs;
int _revs_num;
// 这个是上层的业务处理
callback_t _cb;
};
a. AddConnection 函数
TcpServer 类中的 AddConnection 函数用于进行事件注册。
在注册事件时需要传入一个文件描述符和三个回调函数,表示当该文件描述符上的事件(默认只关心读事件)就绪后应该执行的回调方法。
在 AddConnection 函数内部要做的就是:将套接字设置为非阻塞(ET 模式的要求),把套接字和回调函数等属性封装为一个 Connection 对象,再将该套接字添加到 epoll 模型中,并建立文件描述符与 Connection 之间的映射关系,由 TcpServer 统一管理。
// 这里专门针对任意sock进行添加TcpServer
void AddConnection(int sock, func_t recv_cb, func_t send_cb, func_t except_cb)
{
Sock::SetNonBlock(sock);
// 除了 _listensock,以后会存在大量的socket,每一个sock都必须被封装成为一个Connection
// 当服务器中存在大量的Connection时,TcpServer需要将所有的Connection要进行管理:先描述,再组织
// 1. 构建conn对象,封装sock
Connection *conn = new Connection(sock);
conn->SetCallBack(recv_cb, send_cb, except_cb);
conn->_tsvr = this;
//conn->_lasttimestamp = time();
// 2. 添加sock到epoll中
_poll.AddSockToEpoll(sock, EPOLLIN | EPOLLET); // 任何多路转接的服务器,一般默认只会打开对读取事件的关心,写入事件会按需进行打开
// 3. 还要将对应的Connection*对象指针添加到Connections映射表中
_connections.insert(std::make_pair(sock, conn));
}
b. Dispatcher 函数(初始分发器)
TcpServer 中的 Dispatcher 函数即初始分发器。其内部是一个 while 循环,每轮调用一次 epoll_wait 等待事件发生。
epoll_wait 返回就绪事件的个数 n,就绪事件存放在 _revs 数组中。
随后遍历这 n 个事件,对每一个事件依次做三件事:
- 取出该事件的文件描述符 fd 与事件类型 revents;
- 通过 unordered_map 找到该文件描述符对应的 Connection 结构;
- 根据 revents 判断该调用哪个回调,并调用它:
- 含 EPOLLERR / EPOLLHUP --> 调用异常回调(应优先判断)
- 含 EPOLLIN --> 调用读回调
- 含 EPOLLOUT --> 调用写回调 本轮所有事件处理完毕后,回到循环开头继续 epoll_wait。
LoopOnce() —— 单轮事件处理
void LoopOnce()
{
int n = _poll.WaitEpoll(_revs, _revs_num);
for (int i = 0; i < n; i++)
{
int sock = _revs[i].data.fd;
uint32_t revents = _revs[i].events;
// 将所有的异常全部交给read或者write来统一处理
if(revents & EPOLLERR) revents |= (EPOLLIN | EPOLLOUT);
if(revents & EPOLLHUP) revents |= (EPOLLIN | EPOLLOUT);
if (revents & EPOLLIN)
{
if (IsConnectionExists(sock) && _connections[sock]->_recv_cb != nullptr)
_connections[sock]->_recv_cb(_connections[sock]);
}
if (revents & EPOLLOUT)
{
if (IsConnectionExists(sock) && _connections[sock]->_send_cb != nullptr)
_connections[sock]->_send_cb(_connections[sock]);
}
}
}
Dispatcher() —— 事件循环
// 根据就绪的事件,进行特定事件的派发
void Dispatcher(callback_t cb)
{
_cb = cb;
while (true)
{
ConnectAliveCheck();
LoopOnce(); // 将epoll当做定时器来使用
}
}
private 成员(与前一致)
private:
int _listensock;
int _port;
Epoll _poll;
std::unordered_map<int, Connection *> _connections;
struct epoll_event *_revs;
int _revs_num;
// 这个是上层的业务处理
callback_t _cb;
};
这段代码没有通过 switch /if 分支判断 epoll_wait 的返回结果,而是利用 for 循环条件本身完成返回值判定:
- epoll_wait 返回 -1:代表调用出错,循环条件不成立,不会进入循环体执行事件逻辑;
- epoll_wait 返回 0:代表等待超时,循环条件同样不成立,跳过事件处理;
- epoll_wait 返回大于 0:代表成功捕获到就绪事件,进入 for 循环,执行对应事件的回调函数。
在事件处理流程里,会优先处理异常事件,并交由对应的回调函数完成处理。
c. EnableReadWrite 函数
TcpServer 类中的 EnableReadWrite 函数,用于动态修改某个文件描述符的监听事件。 调用该函数时需要传入三个参数:
- 一个文件描述符,表示需要设置的是哪个文件描述符对应的事件;
- 两个 bool 值,分别表示是否关心读事件以及是否关心写事件。
函数内部会调用 Epoll 类封装好的 CtrlEpoll 函数(即 epoll_ctl(EPOLL_CTL_MOD))修改该文件描述符的监听事件。
为什么需要这个函数?
因为写事件不能长期注册:只要底层发送缓冲区不满,写事件就会一直处于就绪状态,导致 epoll_wait 每次都立即返回、CPU 空转。因此正确做法是"按需开关":
- 数据发不完时 --> EnableReadWrite(fd, true, true) 打开写事件;
- 数据全部发完后 --> EnableReadWrite(fd, true, false) 立即关闭写事件。
读事件通常是长期关心的(连接活着就要读),所以这个函数主要用来开关写事件。
void EnableReadWrite(Connection *conn, bool readable, bool writeable)
{
uint32_t events = ((readable ? EPOLLIN : 0) | (writeable ? EPOLLOUT : 0));
bool res = _poll.CtrlEpoll(conn->_sock, events);
assert(res);
}
(4)回调函数
- Accepter:连接事件触发时,执行这个回调函数,获取底层已经建立完成的连接。
- Recver:读事件就绪时,执行该回调函数,读取客户端传输的数据并做业务处理。
- Sender:写事件就绪时,执行该回调函数,向客户端回传响应数据。
- Excepter:异常事件发生时,调用该回调函数,完成各类资源的释放工作。
在为文件描述符创建 Connection 实例时,可以调用 Connection 类的 SetCallBack 方法,把上面这几组回调函数注册到 Connection 对象中。
- 监听套接字所绑定的 Connection 实例,它的_recvCb 回调会设置成 Accepter。原因是监听套接字的读事件就绪,本质就代表有新连接到来;监听套接字只需要关注读事件,所以它的_sendCb 与_exceptCb 可以直接置空(nullptr)。
- 当 Dispatcher 检测到监听套接字的读事件就绪后,就会执行该 Connection 内的_recvCb 回调,也就是 Accepter,以此获取新建立的底层连接。
- 而和客户端通信的普通套接字,它对应的 Connection 实例会把_recvCb、_sendCb、_exceptCb 分别绑定为 Recver、Sender、Excepter。
- 一旦 Dispatcher 捕获到这类套接字的就绪事件,就会执行 Connection 对象上对应的回调函数,也就是 Recver、Sender、Excepter。
A. Accepter
Accepter 回调专门负责处理新连接事件,执行逻辑如下:
- 调用封装好的 Accept 函数,拿到底层完成建立的客户端连接。
- 调用 AddConnection 函数,把新获取的套接字包装成 Connection 对象,交由服务器统一管理。
- 完成这一步后,套接字以及它需要监听的事件,就注册到 Dispatcher 事件分发器里了。
后续 Dispatcher 做事件派发时,会持续监听这个套接字的事件;一旦事件就绪,就自动执行该套接字所属 Connection 对象绑定的回调函数。
void Accepter(Connection *conn)
{
// logMessage(DEBUG, "Accepter been called");
// 一定是listensock已经就绪了,此次读取不会阻塞
while (true)
{
std::string clientip;
uint16_t clientport;
int accept_errno = 0;
// sock一定是非常规的IO sock
int sock = Sock::Accept(conn->_sock, &clientip, &clientport, &accept_errno);
if (sock < 0)
{
if (accept_errno == EAGAIN || accept_errno == EWOULDBLOCK) break;
else if (accept_errno == EINTR) continue; // 概率非常低
else
{
// accept失败
logMessage(WARNING, "accept error, %d : %s", accept_errno, strerror(accept_errno));
break;
}
}
// 将sock托管给TcpServer
if (sock >= 0)
{
AddConnection(sock, std::bind(&TcpServer::Recver, this, std::placeholders::_1),
std::bind(&TcpServer::Sender, this, std::placeholders::_1),
std::bind(&TcpServer::Excepter, this, std::placeholders::_1));
logMessage(DEBUG, "accept client %s:%d success, add to epoll&TcpServer success, sock: %d",
clientip.c_str(), clientport, sock);
}
}
}
这里实现的是 ET 模式的 epoll 服务器,所以在获取新连接时,需要循环调用 accept,同时监听套接字必须配置为非阻塞。
- ET 模式的特点是:仅当就绪事件状态发生变化(连接从无到有、数量由少变多)才会向上通知。如果没有一次性把内核中全部已建立的连接读取干净,并且后续不再产生新连接,那么内核里残留的连接就会永久丢失。a
- 循环调用 accept 会带来一个问题:当内核所有连接都读取完毕后,继续调用 accept 就会阻塞等待新连接。因此监听套接字要设置成非阻塞,内核无连接时 accept 会直接返回而不会阻塞。 同时,accept 返回的新客户端套接字也需要设置为非阻塞,防止后续循环调用 recv、send 时发生阻塞。
- 上述非阻塞配置,统一由 AddConnection 内部的 SetNonBlock 函数完成。
a. 设置非阻塞
设置文件描述符为非阻塞时,需先调用 fcntl 函数获取该文件描述符对应的文件状态标记,然后在该文件状态标记的基础上添加非阻塞标记 O_NONBLOCK,最后调用 fcntl 函数对该文件描述符的状态标记进行设置即可。
static bool SetNonBlock(int sock)
{
int fl = fcntl(sock, F_GETFL);
if(fl < 0) return false;
fcntl(sock, F_SETFL, fl | O_NONBLOCK);
return true;
}
监听套接字开启非阻塞之后,如果当前没有就绪连接,accept 会以错误的形式返回。所以当 accept 返回值小于 0 时,我们需要进一步判断错误码:
- 如果错误码是EAGAIN或EWOULDBLOCK,代表内核已经没有待接收的连接,所有连接都读取完成,直接返回 0,代表本次 Accepter 回调执行成功。
- 如果错误码是EINTR,说明 accept 调用过程被系统信号打断,需要继续循环调用 accept 尝试获取连接。
- 其余错误码则代表 accept 发生了真正的异常,返回 - 1,表示本次 Accepter 回调执行失败。
accept、recv、send 这类 IO 系统调用为什么会被信号中断?
当 IO 系统调用出错返回,并且错误码被置为EINTR,就代表本次 IO 读写操作在完成前被信号打断。简单来说:系统调用陷入内核之后、还没返回用户态之前,内核收到了信号。
内核准备从内核态切回用户态前,会检查未决信号集(pending 位图)。如果存在还未处理的信号,内核就会优先去处理这个信号。
EINTR属于一种特殊场景。IO 操作分为等待就绪和数据拷贝两个阶段,其中等待阶段耗时通常更长,此时进程处于阻塞等待的闲置状态。如果在这个等待阶段收到信号,内核会立刻暂停本次 IO 调用,转而去处理信号,于是就产生了EINTR中断错误。
b. 写事件按需开启
Accepter 接收新连接套接字,并把它注册到 Dispatcher 时,仅注册EPOLLIN与EPOLLET,也就是只让 epoll 监听该套接字的读事件。
不注册写事件的原因很简单:此时暂无待发送数据,没必要让 epoll 持续监听写就绪事件。通常读事件默认开启,写事件采用按需开启的策略:只有存在待发送数据时,才向 epoll 注册写事件;一旦所有数据发送完成,就立刻关闭该套接字的写事件监听。
B. Recver
Recver 回调负责处理读事件,执行流程如下: 循环调用 recv 函数读取数据,并把读到的数据存入当前套接字对应的 Connection 对象的_inBuffer缓冲区。 对_inBuffer内的数据做报文解析拆分,取出完整报文,未读完的残留数据继续保留在缓冲区中。 调用业务处理函数,处理解析出来的完整报文。
当 recv 返回值小于 0 时,需要判断错误码:
- 错误码为EAGAIN或EWOULDBLOCK:内核缓冲区数据已经全部读取完成;
- 错误码为EINTR:本次读取被信号中断,需要继续循环调用 recv 读取数据;
- 其余情况:代表读操作发生错误。
一旦读操作出错,直接触发该套接字绑定的_exceptCb异常回调,在异常回调内部关闭套接字。
a. 报文切割
报文切割的核心目的是解决 TCP 粘包问题,而粘包的处理方式和自定义协议强相关。 我们需要按照协议规则拆分报文,例如 UDP 常采用定长报头加自描述字段的方案实现报文分离。 本文重点在于演示完整的数据处理流程,为降低复杂度,不设计复杂协议,直接使用字符X作为报文分隔标记,每一条报文的末尾都以X标识报文结束。
报文切割逻辑就是以X作为分隔符,对_inBuffer内的数据进行拆分。 SplitMessage函数负责完成这个工作:把_inBuffer中切分出的完整报文存入 vector 容器;那些不足以组成完整报文的剩余数据,则继续保留在_inBuffer缓冲区里。
// 要把传入进来的缓冲区进行切分
// 1. buffer被切走的,同时也要将其从buffer中移除
// 2. 可能会存在多个报文,多个报文依次放入out
// buffer: 输入输出型参数 out: 输出型参数
void SpliteMessage(std::string &buffer, std::vector<std::string> *out)
{
while (true)
{
auto pos = buffer.find(SEP);
if (std::string::npos == pos) break;
std::string message = buffer.substr(0, pos);
buffer.erase(0, pos + SEP_LEN);
out->push_back(message);
}
}
b. 业务处理函数
对分割得到的完整报文执行反序列化操作。 执行业务逻辑。 业务处理完成后,组装生成响应报文。 把组装好的响应报文写入当前 Connection 对象的_outBuffer输出缓冲区,同时开启套接字的写事件监听。
后续 Dispatcher 进行事件分发时,就会监听该套接字的写事件;一旦写事件就绪,就会触发 Connection 绑定的写回调函数,把_outBuffer缓冲区里的响应数据发送给客户端。
void NetCal(Connection *conn, std::string &request)
{
logMessage(DEBUG, "NetCal been called, get request: %s", request.c_str());
// 1. 反序列化
Request req;
if(!req.Deserialized(request)) return;
// 2. 业务处理
Response resp = calculator(req);
// 3. 先序列化,再构建应答
std::string sendstr = resp.Serialize();
sendstr = Encode(sendstr);
// 4. 交给服务器conn
conn->_outbuffer += sendstr;
// 5. 想办法让底层的TcpServer,让它开始发送
// // a. 需要有完整的发送逻辑
// // b. 触发发送的动作,一旦开启EPOLLOUT,epoll会自动立马触发一次发送事件就绪,
// // 如果后续保持发送的开启,epoll会一直发送
conn->_tsvr->EnableReadWrite(conn, true, true);
}
c. 协议定制
协议部分:Encode / Request / Response
std::string Encode(std::string &s)
{
return s + SEP;
}
class Request
{
public:
std::string Serialize()
{
std::string str;
str = std::to_string(x_);
str += SPACE;
str += op_;
str += SPACE;
str += std::to_string(y_);
return str;
}
bool Deserialized(const std::string &str)
{
std::size_t left = str.find(SPACE);
if (left == std::string::npos) return false;
std::size_t right = str.rfind(SPACE);
if (right == std::string::npos) return false;
x_ = atoi(str.substr(0, left).c_str());
y_ = atoi(str.substr(right + SPACE_LEN).c_str());
if (left + SPACE_LEN > str.size()) return false;
else op_ = str[left + SPACE_LEN];
return true;
}
public:
Request(){}
Request(int x, int y, char op)
: x_(x)
, y_(y)
, op_(op)
{}
~Request(){}
public:
int x_;
int y_;
char op_; // '+', '-', '*', '/', '%'
};
class Response
{
public:
std::string Serialize()
{
std::string s;
s = std::to_string(code_);
s += SPACE;
s += std::to_string(result_);
return s;
}
bool Deserialized(const std::string &s)
{
std::size_t pos = s.find(SPACE);
if (pos == std::string::npos) return false;
code_ = atoi(s.substr(0, pos).c_str());
result_ = atoi(s.substr(pos + SPACE_LEN).c_str());
return true;
}
public:
Response(){}
Response(int result, int code)
: result_(result)
, code_(code)
{}
~Response(){}
public:
int result_; // 计算结果
int code_; // 计算结果的状态码
};
代码解析
这段代码实现简易的请求、响应报文序列化与反序列化逻辑,配合前面 Reactor 消息队列的报文切割模块使用。
1.Encode 函数:在序列化后的报文末尾追加SEP分隔符,标记报文边界,用于SplitMessage切割报文,解决 TCP 粘包。
2.Request 类:封装客户端发来的四则运算请求,成员包含两个操作数 x_、y_和运算符 op_。
- Serialize:把对象成员拼接成x op y格式字符串,完成序列化,用于网络发送。
- Deserialized:接收报文字符串,利用find和rfind查找空格,拆分解析出两个数字和运算符,填充到对象成员,解析失败返回 false。
3.Response 类:封装服务器返回给客户端的应答数据,包含运算结果 result_与状态码 code_。
- Serialize:将状态码、运算结果拼接为code result字符串,序列化输出。
- Deserialized:查找空格分割字符串,解析得到状态码和结果,还原响应对象。
整体业务链路:Recver 读取数据 --> SplitMessage 切出完整报文 --> Request 反序列化 --> 业务计算 --> Response 对象组装 --> Response 序列化 --> Encode 追加报文分隔符 --> 写入输出缓冲区_outBuffer,开启写事件 --> Dispatcher 触发 Sender 回调发送数据。
C. Sender
- 循环调用 send 函数发送数据,并将发送出去的数据从该套接字对应 Connection 结构的 _outBuffer 中删除。
- 若循环调用 send 函数后该套接字对应的 _outBuffer 中的数据被全部发送,此时就需要将该套接字对应的写事件关闭,因为已没有要发送的数据了,若 _outBuffer 中的数据还有剩余,那么该套接字对应的写事件就应继续打开。
// 最开始时,conn是没有被触发的
void Sender(Connection *conn)
{
while(true)
{
ssize_t n = send(conn->_sock, conn->_outbuffer.c_str(), conn->_outbuffer.size(), 0);
if(n > 0)
{
conn->_outbuffer.erase(0, n);
if(conn->_outbuffer.empty()) break;
}
else
{
if(errno == EAGAIN || errno == EWOULDBLOCK) break;
else if(errno == EINTR) continue;
else
{
logMessage(ERROR, "send error, %d : %s", errno, strerror(errno));
conn->_except_cb(conn);
break;
}
}
}
// 虽然不能确定是否发完了,但可以保证:如果没有出错,要么发完了,要么发送条件不满足,下次再发送
if(conn->_outbuffer.empty()) EnableReadWrite(conn, true, false);
else EnableReadWrite(conn, true, true);
}
Sender 是写事件回调函数,负责将_outbuffer缓冲区的数据通过非阻塞 send 循环发送。发送成功则从缓冲区移除已发送字节;遇到EAGAIN/EWOULDBLOCK代表内核缓冲区已满,停止发送等待下次写事件;EINTR代表被信号中断,继续重试。发生其他错误时触发异常回调关闭连接。函数末尾根据缓冲区是否为空,动态调整 epoll 监听:缓冲区为空就关闭写事件,防止空轮询;缓冲区还有数据,则继续监听写事件,等待下次就绪继续发送。
D. Excepter
- 对于异常事件就绪的套接字不做过多处理,调用 close 函数将该套接字关闭即可。
- 但在关闭该套接字前,需先将该套接字从 epoll 模型中删除,并取消该套接字与其对应的 Connection 结构的映射关系。
- 释放 Connection 对象。
void Excepter(Connection *conn)
{
if(!IsConnectionExists(conn->_sock)) return;
// 1. 从epoll中移除
bool res = _poll.DelFromEpoll(conn->_sock);
assert(res); // 要进行判断
// 2. 从unorder_map中移除
_connections.erase(conn->_sock);
// 3. close(sock);
close(conn->_sock);
// 4. delete conn;
delete conn;
logMessage(DEBUG, "Excepter 回收完毕,所有的异常情况");
}
代码解析:
- Excepter 是异常回调函数,当连接发生读写错误时,用来安全回收连接资源。
- 先判断连接是否还存在,如果不存在直接返回,防止重复回收。
- 将套接字从 epoll 实例中删除,epoll 不再监听该 fd 上的事件;assert用来调试阶段校验删除是否成功。
- 从保存所有连接的unordered_map(_connections)中擦除这条连接记录。
- 关闭 socket 文件描述符,释放操作系统层面的套接字资源。
- delete conn释放堆上创建的 Connection 对象内存,防止内存泄漏。
- 打印调试日志,标记本次连接资源回收完成。
作用:不管是读异常、写异常还是对端关闭连接,都会走到这个回调,统一完成资源清理,避免 fd 泄露、内存泄露。
(5)Socket 套接字
#pragma once
#include <iostream>
#include <string>
#include <cstring>
#include <cerrno>
#include <cassert>
#include <unistd.h>
#include <memory>
#include <sys/types.h>
#include <sys/socket.h>
#include <arpa/inet.h>
#include <netinet/in.h>
#include <ctype.h>
#include <fcntl.h>
class Sock
{
private:
const static int gbacklog = 10;
public:
Sock() {}
static int Socket()
{
int listensock = socket(AF_INET, SOCK_STREAM, 0);
if (listensock < 0)
{
exit(2);
}
int opt = 1;
setsockopt(listensock, SOL_SOCKET, SO_REUSEADDR | SO_REUSEPORT, &opt, sizeof(opt));
return listensock;
}
static void Bind(int sock, uint16_t port, std::string ip = "0.0.0.0")
{
struct sockaddr_in local;
memset(&local, 0, sizeof local);
local.sin_family = AF_INET;
local.sin_port = htons(port);
inet_pton(AF_INET, ip.c_str(), &local.sin_addr);
if (bind(sock, (struct sockaddr *)&local, sizeof(local)) < 0)
{
exit(3);
}
}
static void Listen(int sock)
{
if (listen(sock, gbacklog) < 0)
{
exit(4);
}
}
static int Accept(int listensock, std::string *ip, uint16_t *port, int *accept_errno)
{
struct sockaddr_in src;
socklen_t len = sizeof(src);
*accept_errno = 0;
int servicesock = accept(listensock, (struct sockaddr *)&src, &len);
if (servicesock < 0)
{
*accept_errno = errno;
return -1;
}
if(port) *port = ntohs(src.sin_port);
if(ip) *ip = inet_ntoa(src.sin_addr);
return servicesock;
}
static bool Connect(int sock, const std::string &server_ip, const uint16_t &server_port)
{
struct sockaddr_in server;
memset(&server, 0, sizeof(server));
server.sin_family = AF_INET;
server.sin_port = htons(server_port);
server.sin_addr.s_addr = inet_addr(server_ip.c_str());
if(connect(sock, (struct sockaddr*)&server, sizeof(server)) == 0) return true;
else return false;
}
static bool SetNonBlock(int sock)
{
int fl = fcntl(sock, F_GETFL);
if(fl < 0) return false;
fcntl(sock, F_SETFL, fl | O_NONBLOCK);
return true;
}
~Sock() {}
};
sock是封装 socket 基础操作的工具类,全部为静态方法,统一封装 TCP 网络编程的基础接口,提供创建套接字、绑定、监听、接受连接、客户端连接、设置非阻塞等能力。
- Socket():创建 TCP 监听套接字,并且设置端口复用与地址复用,服务端重启可以快速复用端口。失败直接退出程序。
- Bind():填充本机地址结构体,调用 bind 将 socket 绑定到指定 IP 与端口。默认绑定0.0.0.0,监听本机所有网卡。
- Listen():开启监听,gbacklog 为 10,设置 TCP 半连接队列长度。
- Accept():阻塞接受新连接,输出对端 IP、端口,同时带回 accept 调用的错误码;返回新连接的 socket 文件描述符。
- Connect():客户端使用,向指定 IP 端口发起 TCP 连接,连接成功返回 true。
- SetNonBlock():通过fcntl把指定 socket 设置为非阻塞模式,这是 Reactor 高并发模型的基础。
说明:这个工具类是整个 Reactor 消息队列底层网络基础,Sock::SetNonBlock把所有连接 socket 设置非阻塞,配合 epoll + 四个回调函数(Accepter、Recver、Sender、Excepter)实现事件驱动服务器
(6)代码演示








(7)服务器测试
客户端建立连接后,服务端会为该连接分配 5 号文件描述符,4 号文件描述符由 epoll 占用。客户端可向服务端提交多项简单计算任务,先进行运行:

这个是一条条的发送,但在真实的业务场景中,肯定很多条数据一起发送,我们继续来演示:

一次发很对条,演示粘包问题:

这里我们会发现9+这个计算表达式子没有显示出来
日志可以看到,服务端第一次 recv 只读到 1 条报文(1 + 1),第二次却一口气读到了 2 条(3 * 8 和 8 / 0)——这就是"粘包"。根本原因在于 TCP 是字节流,不保证消息边界:recv 一次能读到多少,只取决于调用它的那一刻内核接收缓冲区里积了多少字节,与发送方发了几次无关。而服务端并没有把这两条当成一条处理——_inbuffer 先把数据攒下来,SpliteMessage 再按 \r\n 逐条切分,最终 NetCal 被准确调用了 2 次。这正是 _inbuffer + 分隔符循环切分存在的意义:缓冲区负责"攒",循环负责"切"。
由于使用了多路转接技术,虽然 epoll 服务器是一个单进程的服务器,但却可同时为多个客户端提供服务。

(8)总结
基于 IO 多路转接,当检测到事件就绪后,通过回调函数执行业务逻辑的编程模型,叫做反应堆(Reactor)模式。上面代码中的 TcpServer 就是一个反应堆实例,而每一个 Connection 对象代表一个事件。每个事件内部包含:
- 文件描述符
- 独立缓冲区
- 回调函数
- 指向反应堆的指针
反应堆内部存在事件分发函数,一旦 epoll 监测到某个事件就绪,分发函数就会触发该事件绑定的回调函数执行。
特性
- 单进程模型:由同一个进程完成事件分发与 IO 读写操作;
- 半同步半异步:异步体现在事件的到来是随机、不可预测的;
- 同步体现在 IO 操作由当前线程直接执行。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/2501_93351213/article/details/167223811




