{T}

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 模式。其核心包含两个要素:

  1. 事件分发线程(Reactor 线程 / Event Loop 线程):运行无限循环,通过 poll()/epoll() 等 I/O 分发技术检测就绪事件
  2. 事件回调机制:每个 I/O 事件注册对应的回调函数,事件就绪时由 Reactor 线程调用回调处理

典型事件类型包括:

  • Acceptor 上的连接建立事件
  • 已连接套接字上的数据可读事件
  • 已连接套接字上的发送缓冲区可写事件
  • 通信管道(pipe)上的数据到达事件

I/O 模型与线程模型对比

网络程序的处理流程可抽象为五个阶段:

阶段说明类型
read从套接字接收数据I/O 密集
decode解析接收数据CPU 密集
compute业务逻辑计算CPU 密集
encode编码处理结果CPU 密集
send通过套接字发送结果I/O 密集

其中 readsend 与套接字直接相关,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 服务器:

c
#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连接关闭

运行验证

启动服务器程序:

plaintext
$ ./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 客户端,观察服务器输出:

plaintext
[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 事件分离至不同线程处理。

思考题

  1. 尝试修改 onMessage 回调,实现简单的命令处理(如 cdls 等文件操作命令)。
  2. 当前程序中 decode-compute-encode 在哪个执行上下文中运行?如何将业务逻辑与 I/O 逻辑解耦?

版本信息

项目说明
更新日期2026-06-09
目标内核Linux 7.0
关键 APIpoll() (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 能力