主从Reactor模式:多线程事件分发
在前文中,我们引入了 Reactor 模式,并通过单线程 Reactor 同时分发 Acceptor 连接建立事件和已连接套接字的 I/O 事件。该模式在高并发场景下存在瓶颈:单 Reactor 线程同时承担连接分发与 I/O 分发职责,当客户端连接请求密集时,连接成功率会显著下降。此外,现代硬件普遍配备多核多路 CPU,单线程模型无法充分利用计算资源。
本文将 Acceptor 连接建立事件与已连接套接字 I/O 事件分离,形成主从 Reactor 模式(Master-Slave Reactor Pattern),并通过多线程实现高效的事件分发。
主从 Reactor 模式
架构设计
主从 Reactor 模式的核心思想是职责分离:
- 主 Reactor(Main Reactor):仅负责分发 Acceptor 上的连接建立事件,通过
accept()获取已连接套接字 - 从 Reactor(Sub Reactor):负责分发已连接套接字上的 I/O 事件(可读、可写等)
从 Reactor 的数量可根据 CPU 核数灵活配置,实现多核并行处理。同一套接字的 I/O 事件始终由同一个从 Reactor 线程处理,避免了多线程并发访问同一套接字的锁开销。
跨线程连接分配
主 Reactor 线程获取已连接套接字后,需将其分配至某个从 Reactor 线程。这涉及跨线程数据传递问题:从 Reactor 线程可能正处于 poll()/epoll_wait() 的阻塞状态。
常见的解决方案包括:
- self-pipe trick:从 Reactor 线程创建一个管道(pipe),主线程通过写入管道字节唤醒从 Reactor
- eventfd:Linux 2.6.22+ 提供的轻量级事件通知机制,比 pipe 更高效
- epoll 的 EPOLL_CTL_MOD:通过修改 fd 的事件掩码触发从 Reactor 唤醒
这些机制属于高性能网络框架的核心实现,后续实战篇将详细展开。
主从 Reactor + Worker 线程池
模式演进
主从 Reactor 模式解决了 I/O 分发的高效率问题,但未解决业务逻辑与 I/O 分发之间的耦合。将 Worker 线程池与主从 Reactor 结合,形成生产环境中广泛采用的完整架构:
该架构的核心设计:
- 主 Reactor:专注连接建立,性能瓶颈从单线程模型中释放
- 从 Reactor:专注 I/O 读写,多线程并行处理
- Worker 线程池:处理 CPU 密集型的 decode/compute/encode 操作,与 I/O 线程解耦
Netty 的实现参考
Netty 是该模式的典型实现。以下为 Netty 服务器初始化代码:
public final class TelnetServer {
static final int PORT = Integer.parseInt(System.getProperty("port", "8023"));
public static void main(String[] args) throws Exception {
// 主 Reactor 线程组:仅负责 Acceptor 连接处理
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
// 从 Reactor 线程组:负责已连接套接字的 I/O 事件分发
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap b = new ServerBootstrap();
b.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.handler(new LoggingHandler(LogLevel.INFO))
.childHandler(new TelnetServerInitializer());
b.bind(PORT).sync().channel().closeFuture().sync();
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
}注意:Netty 文档中的 workerGroup 对应本文的从 Reactor 线程,负责 I/O 事件分发,而非业务逻辑处理。业务逻辑线程通常由应用开发者自行设计,通过 ChannelPipeline 中的 ChannelHandler 配置,并建议从 I/O 线程中分离以支持更高并发度。
样例程序
使用课程定制的网络编程框架,仅需修改线程数参数即可实现主从 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
struct event_loop *eventLoop = event_loop_init();
// 初始化 acceptor
struct acceptor *acceptor = acceptor_init(SERV_PORT);
// 创建 TCPServer,线程数设为 4
// 1 个主 Reactor 线程 + 4 个从 Reactor 线程
// 每个从 Reactor 线程绑定独立的 event_loop
struct TCPserver *tcpServer = tcp_server_init(eventLoop, acceptor,
onConnectionCompleted, onMessage,
onWriteCompleted, onConnectionClosed, 4);
tcp_server_start(tcpServer);
// 主线程运行 event_loop,处理 Acceptor 事件
event_loop_run(eventLoop);
return 0;
}与单线程版本的对比
本程序与第 27 讲的唯一区别在于 tcp_server_init() 的线程数参数:第 27 讲为 0(单线程),本文为 4(1 个主 Reactor + 4 个从 Reactor)。框架内部根据线程数参数自动创建对应数量的从 Reactor 线程,每个线程绑定独立的 event_loop。
Worker 线程的集成
框架通过回调函数暴露业务接口,应用开发者可在 onMessage 回调中将 CPU 密集型任务提交至独立线程池:
// 数据读取完成回调 — 集成 Worker 线程池
int onMessage(struct buffer *input, struct tcp_connection *tcpConnection) {
printf("get message from tcp connection %s\n", tcpConnection->name);
printf("%s", input->data);
// 将 decode-compute-encode 提交至 Worker 线程池
struct buffer *output = thread_handle(input);
// 处理完成后通过从 Reactor I/O 线程发送数据
tcp_connection_send_buffer(tcpConnection, output);
return 0;
}运行验证
启动服务器,观察多线程初始化过程:
$ ./poll-server-multithreads
[msg] set poll as dispatcher
[msg] add channel fd == 4, main thread
[msg] poll added channel fd==4
[msg] set poll as dispatcher
[msg] add channel fd == 7, main thread
[msg] poll added channel fd==7
[msg] event loop thread init and signal, Thread-1
[msg] event loop run, Thread-1
[msg] event loop thread started, Thread-1
[msg] set poll as dispatcher
[msg] add channel fd == 9, main thread
[msg] poll added channel fd==9
[msg] event loop thread init and signal, Thread-2
[msg] event loop run, Thread-2
[msg] event loop thread started, Thread-2
[msg] set poll as dispatcher
[msg] add channel fd == 11, main thread
[msg] event loop thread init and signal, Thread-3
[msg] event loop thread started, Thread-3
[msg] set poll as dispatcher
[msg] event loop run, Thread-3
[msg] add channel fd == 13, main thread
[msg] poll added channel fd==13
[msg] event loop thread init and signal, Thread-4
[msg] event loop run, Thread-4
[msg] event loop thread started, Thread-4
[msg] add channel fd == 5, main thread
[msg] event loop run, main thread客户端连接后,观察事件分发行为:
[msg] get message channel i==1, fd==5
[msg] activate channel fd == 5, revents=2, main thread
[msg] new connection established, socket == 14
connection completed
[msg] get message channel i==0, fd==7
[msg] activate channel fd == 7, revents=2, Thread-1
[msg] wakeup, Thread-1
[msg] add channel fd == 14, Thread-1
[msg] poll added channel fd==14
[msg] get message channel i==1, fd==14
[msg] activate channel fd == 14, revents=2, Thread-1
get message from tcp connection connection-14
fasfas关键观察:
main thread仅处理新连接建立(fd==5为监听套接字上的可读事件)Thread-1至Thread-4分别处理不同连接的 I/O 事件- 新连接被负载均衡分配至不同的从 Reactor 线程
模型评估
| 维度 | 单 Reactor 线程 | 主从 Reactor 多线程 |
|---|---|---|
| 连接处理能力 | 受限于单线程 | 主 Reactor 专注 accept,吞吐量高 |
| I/O 并行度 | 单线程串行 | 多线程并行处理 |
| CPU 利用率 | 单核 | 多核充分利用 |
| 锁开销 | 无 | 低(同套接字仅一个线程处理) |
| 实现复杂度 | 中等 | 较高(需处理跨线程通信) |
| 适用场景 | I/O 密集、轻量业务 | 高并发、多核服务器 |
思考题
- 日志中
main thread首先加入了fd==4(监听套接字),随后又加入了fd==7。fd==7的用途是什么?它与跨线程通信有何关系? - 尝试修改服务器代码,将 decode-compute-encode 部分使用线程池实现,并与从 Reactor I/O 线程分离。
版本信息
| 项目 | 说明 |
|---|---|
| 更新日期 | 2026-06-09 |
| 目标内核 | Linux 7.0 |
| 关键 API | poll() (POSIX.1-2001); pthread_create() (POSIX.1); eventfd() (Linux 2.6.22+); epoll_ctl() (Linux 2.5.44+) |
| 备注 | Linux 7.0 中 io_uring(Linux 5.1+ 引入)提供了内核侧异步 I/O 能力,可作为 Reactor 模式的替代或补充;pidfd + pidfd_send_signal()(Linux 5.4+)可用于线程级信号管理;SO_INCOMING_CPU 套接字选项(Linux 3.19+)支持连接与 CPU 核的亲和性绑定,优化主从 Reactor 的负载均衡 |