{T}

主从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 服务器初始化代码:

java
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 模式:

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
    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 密集型任务提交至独立线程池:

c
// 数据读取完成回调 — 集成 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;
}

运行验证

启动服务器,观察多线程初始化过程:

plaintext
$ ./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

客户端连接后,观察事件分发行为:

plaintext
[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-1Thread-4 分别处理不同连接的 I/O 事件
  • 新连接被负载均衡分配至不同的从 Reactor 线程

模型评估

维度单 Reactor 线程主从 Reactor 多线程
连接处理能力受限于单线程主 Reactor 专注 accept,吞吐量高
I/O 并行度单线程串行多线程并行处理
CPU 利用率单核多核充分利用
锁开销低(同套接字仅一个线程处理)
实现复杂度中等较高(需处理跨线程通信)
适用场景I/O 密集、轻量业务高并发、多核服务器

思考题

  1. 日志中 main thread 首先加入了 fd==4(监听套接字),随后又加入了 fd==7fd==7 的用途是什么?它与跨线程通信有何关系?
  2. 尝试修改服务器代码,将 decode-compute-encode 部分使用线程池实现,并与从 Reactor I/O 线程分离。

版本信息

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