Reactor模式:poll单线程事件分发
在前两讲中,我们分别使用 fork() 进程和 pthread 线程处理并发连接。这两种技术实现简单,但性能随并发数增长而快速下降,无法满足极端高并发需求。本文引入 I/O 事件分发(I/O event dispatch)机制,阐述 Reactor 模式(反应堆模式)的核心思想,并实现基于 poll() 的单线程事件分发服务器。
完整代码可在 GitHub 获取。
事件驱动模型
事件驱动的基本原理
事件驱动(event-driven)模型具有资源占用少、效率高、可扩展性强的特点,是支持高性能高并发的核心范式。
在 GUI 编程中,程序为按钮点击、文本输入等控件操作注册回调函数(callback function)。后台运行一个无限循环的事件分发线程(event dispatch thread),持续从事件队列(event queue)中取出事件,查找并执行对应的回调函数。Web 前端编程的 JavaScript 事件机制遵循相同原理。
网络编程中,通过 poll()、epoll() 等 I/O 多路复用(I/O multiplexing)技术,可将套接字事件抽象为事件驱动模型,实现高性能并发处理。
Reactor 模式
事件驱动模型在网络编程中被称为 Reactor 模式(反应堆模式)或 Event Loop 模式。其核心包含两个要素:
- 事件分发线程(Reactor 线程 / Event Loop 线程):运行无限循环,通过
poll()/epoll()等 I/O 分发技术检测就绪事件 - 事件回调机制:每个 I/O 事件注册对应的回调函数,事件就绪时由 Reactor 线程调用回调处理
典型事件类型包括:
- Acceptor 上的连接建立事件
- 已连接套接字上的数据可读事件
- 已连接套接字上的发送缓冲区可写事件
- 通信管道(pipe)上的数据到达事件
I/O 模型与线程模型对比
网络程序的处理流程可抽象为五个阶段:
| 阶段 | 说明 | 类型 |
|---|---|---|
| read | 从套接字接收数据 | I/O 密集 |
| decode | 解析接收数据 | CPU 密集 |
| compute | 业务逻辑计算 | CPU 密集 |
| encode | 编码处理结果 | CPU 密集 |
| send | 通过套接字发送结果 | I/O 密集 |
其中 read 和 send 与套接字直接相关,decode/compute/encode 属于业务逻辑。以下对比几种并发模型对这些阶段的处理方式。
fork 进程模型
第 25 讲中,每个连接由独立子进程处理。随着连接数增长,子进程数量线性增加,即使连接空闲也持续占用进程资源。fork() 的创建开销大,且进程间无法共享内存,通信成本高。
pthread 线程模型
第 26 讲中,每个连接由独立线程处理。线程比进程更轻量,线程池进一步降低了创建开销。但空闲连接仍占用线程资源,阻塞 I/O 模型下线程无法释放给其他连接使用。
单 Reactor 线程模型
Reactor 模式通过事件分发解决空闲连接占用资源的问题。单 Reactor 线程同时负责 Acceptor 连接建立事件和已连接套接字的 I/O 事件分发。
单 Reactor 线程 + Worker 线程池
单 Reactor 线程模型中,业务逻辑(decode/compute/encode)与 I/O 事件分发在同一线程执行。XML 解析、数据库查询、文件传输等 CPU 密集型操作会阻塞 Reactor 线程,影响事件分发效率。
解决方案:将业务逻辑拆分为独立任务,提交至 Worker 线程池执行,与 Reactor 线程解耦。Reactor 线程仅负责 I/O 相关工作,计算结果完成后交回 Reactor 线程通过套接字发送。
样例程序
以下使用本课程定制的网络编程框架,实现基于 poll() 的单线程 Reactor 服务器:
#include <lib/acceptor.h>
#include "lib/common.h"
#include "lib/event_loop.h"
#include "lib/tcp_server.h"
char rot13_char(char c) {
if ((c >= 'a' && c <= 'm') || (c >= 'A' && c <= 'M'))
return c + 13;
else if ((c >= 'n' && c <= 'z') || (c >= 'N' && c <= 'Z'))
return c - 13;
else
return c;
}
// 连接建立完成回调
int onConnectionCompleted(struct tcp_connection *tcpConnection) {
printf("connection completed\n");
return 0;
}
// 数据读取完成回调
int onMessage(struct buffer *input, struct tcp_connection *tcpConnection) {
printf("get message from tcp connection %s\n", tcpConnection->name);
printf("%s", input->data);
struct buffer *output = buffer_new();
int size = buffer_readable_size(input);
for (int i = 0; i < size; i++) {
buffer_append_char(output, rot13_char(buffer_read_char(input)));
}
tcp_connection_send_buffer(tcpConnection, output);
return 0;
}
// 数据发送完成回调
int onWriteCompleted(struct tcp_connection *tcpConnection) {
printf("write completed\n");
return 0;
}
// 连接关闭回调
int onConnectionClosed(struct tcp_connection *tcpConnection) {
printf("connection closed\n");
return 0;
}
int main(int c, char **v) {
// 创建 event_loop(Reactor 对象)
struct event_loop *eventLoop = event_loop_init();
// 初始化 acceptor,监听指定端口
struct acceptor *acceptor = acceptor_init(SERV_PORT);
// 创建 TCPServer,线程数设为 0 表示单线程模式
// 单线程同时负责 Acceptor 连接处理和已连接套接字 I/O 处理
struct TCPserver *tcpServer = tcp_server_init(eventLoop, acceptor,
onConnectionCompleted, onMessage,
onWriteCompleted, onConnectionClosed, 0);
tcp_server_start(tcpServer);
// 运行 event_loop 无限循环,等待并分发事件
event_loop_run(eventLoop);
return 0;
}程序结构分析
event_loop:Reactor 核心对象,与线程绑定,内部运行无限循环完成事件分发。底层使用 poll() 作为事件分发机制。
acceptor:封装监听套接字,负责接受新连接。
TCPServer:整合 event_loop、acceptor 及回调函数。线程数参数为 0 时,采用单线程模式——同一线程既处理 Acceptor 连接事件,也处理已连接套接字的 I/O 事件。
回调函数:框架通过四个回调函数暴露业务接口:
| 回调 | 触发时机 |
|---|---|
onConnectionCompleted | 新连接建立完成 |
onMessage | 数据读取至缓冲区 |
onWriteCompleted | 数据发送完成 |
onConnectionClosed | 连接关闭 |
运行验证
启动服务器程序:
$ ./poll-server-onethread
[msg] set poll as dispatcher
[msg] add channel fd == 4, main thread
[msg] poll added channel fd==4
[msg] add channel fd == 5, main thread
[msg] poll added channel fd==5
[msg] event loop run, main thread连接两个 telnet 客户端,观察服务器输出:
[msg] get message channel i==1, fd==5
[msg] activate channel fd == 5, revents=2, main thread
[msg] new connection established, socket == 6
connection completed
[msg] add channel fd == 6, main thread
[msg] poll added channel fd==6
[msg] get message channel i==2, fd==6
[msg] activate channel fd == 6, revents=2, main thread
get message from tcp connection connection-6
afadsfaf
[msg] get message channel i==1, fd==5
[msg] activate channel fd == 5, revents=2, main thread
[msg] new connection established, socket == 7
connection completed
[msg] add channel fd == 7, main thread
[msg] poll added channel fd==7
[msg] get message channel i==3, fd==7
[msg] activate channel fd == 7, revents=2, main thread
get message from tcp connection connection-7
sfasggwqe
[msg] get message channel i==3, fd==7
[msg] activate channel fd == 7, revents=2, main thread
[msg] poll delete channel fd==7
connection closed
[msg] get message channel i==2, fd==6
[msg] activate channel fd == 6, revents=2, main thread
[msg] poll delete channel fd==6
connection closed整个运行过程中仅 main thread 在工作,单线程 Reactor 即可高效处理多个并发连接。
模型评估
| 维度 | 评估 |
|---|---|
| 实现复杂度 | 中等,需理解事件驱动范式 |
| 资源利用率 | 高,空闲连接不占用线程资源 |
| 单线程瓶颈 | 业务逻辑阻塞将影响所有连接的事件分发 |
| 可扩展性 | 受限于单核 CPU 处理能力 |
| 适用场景 | I/O 密集型、业务逻辑轻量的服务 |
单线程 Reactor 模式解决了空闲连接占用资源的问题,但业务逻辑与 I/O 分发在同一线程执行,CPU 密集型操作会阻塞事件循环。下一部分将介绍主从 Reactor 模式,将 Acceptor 连接事件与已连接套接字 I/O 事件分离至不同线程处理。
思考题
- 尝试修改
onMessage回调,实现简单的命令处理(如cd、ls等文件操作命令)。 - 当前程序中 decode-compute-encode 在哪个执行上下文中运行?如何将业务逻辑与 I/O 逻辑解耦?
版本信息
| 项目 | 说明 |
|---|---|
| 更新日期 | 2026-06-09 |
| 目标内核 | Linux 7.0 |
| 关键 API | poll() (POSIX.1-2001, Linux 2.0+); epoll_create()/epoll_ctl()/epoll_wait() (Linux 2.5.44+); accept() (POSIX.1) |
| 备注 | poll() 在 Linux 7.0 中仍为标准 POSIX 接口;epoll 在高并发场景下性能优于 poll(O(1) vs O(n));Linux 6.x 引入的 epoll_pwait2() 支持纳秒级超时精度;Linux 7.0 中 io_uring 可作为 I/O 多路复用的替代方案,提供异步 I/O 能力 |