找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖

4546

积分

0

好友

593

主题
发表于 1 小时前 | 查看: 2| 回复: 0

适用范围:Java 17/21/25、Linux、TCP 长连接、自定义二进制协议、网关、即时通信、设备接入与 RPC 通信层。

摘要

Java NIO 的价值并不只是“一个线程管理很多连接”。真正决定系统吞吐、尾延迟与稳定性的,是围绕 SelectorSocketChannel、缓冲区、任务调度和反压机制建立的一组并发约束:

  • 一个连接的 I/O 状态必须由固定 EventLoop 串行维护
  • 跨线程操作必须先进入 EventLoop 任务队列,再由所属线程执行
  • TCP 读取必须处理拆包与粘包,写出必须处理半写
  • OP_WRITE 只能在确有待发送数据时开启
  • 业务线程池、连接级积压和全局内存都必须有上限
  • 读写缓冲区的生命周期必须清晰,禁止在异步线程之间无约束共享
  • 高并发指标必须建立在可复现压测、容量预算和可观测性之上

本文从模型选择开始,完整讲解主从 Reactor 的线程划分、连接归属、协议编解码、异步业务调度、读写反压、内存治理、优雅停机、Kubernetes 部署与压测方法,并给出一套可落地的 Java NIO 参考实现。如果你正在设计或维护一个高并发通信层,希望这里对线程模型、背压机制和内存治理的讨论能给你一些新的视角。也欢迎在云栈社区与更多开发者交流架构设计和实战经验。

目录

  1. 一次典型的通信层扩容事故
  2. 先做技术选型:BIO、虚拟线程、NIO 与 Netty
  3. Java NIO 的核心对象与操作系统映射
  4. Reactor 模型的三种形态
  5. 生产级主从 Reactor 总体架构
  6. 必须遵守的八条并发不变量
  7. 协议设计:先解决边界,再谈业务
  8. 核心源码实现
  9. 读路径:从字节流到业务消息
  10. 写路径:半写、队列与 OP_WRITE
  11. 背压与过载保护
  12. 缓冲区与堆外内存治理
  13. 线程数与任务调度策略
  14. 连接生命周期与异常处理
  15. 网络参数与 Kubernetes 部署
  16. 可观测性设计
  17. 优雅启动、摘流与停机
  18. 压测方法与容量规划
  19. 常见错误与修复方式
  20. 自研 NIO、Netty 与虚拟线程的决策边界
  21. 生产上线检查清单
  22. 总结

1. 一次典型的通信层扩容事故

某次大促开始后,核心接入网关出现以下现象:

  • 新建连接速率突然上升
  • 活跃连接数持续增长
  • CPU 使用率达到 95%
  • P99 延迟由个位数毫秒恶化到数百毫秒
  • 线程数、上下文切换和堆内存同步上涨
  • 部分实例触发 Full GC,随后被健康检查摘除

旧架构采用“一个连接对应一个平台线程”的阻塞式处理模型。大量连接处于空闲或慢读状态时,线程仍然占用栈空间和调度资源。请求高峰到来后,线程池、连接队列与业务队列相互放大,最终形成如下故障链:

连接突增
   ↓
平台线程数上升
   ↓
线程栈与上下文切换成本增加
   ↓
业务队列等待时间上升
   ↓
超时重试进一步放大流量
   ↓
CPU 饱和、P99 恶化、实例失稳

这里真正需要解决的并不是“把 BIO 改成 NIO”这么简单,而是重新设计整个通信层:

  1. 如何让少量 I/O 线程管理大量连接
  2. 如何确保慢业务不阻塞 I/O 线程
  3. 如何解决 TCP 字节流的拆包、粘包与半写
  4. 如何在业务消费能力不足时限制读入速度
  5. 如何控制单连接和全局待发送内存
  6. 如何在容器环境中完成摘流与优雅停机
  7. 如何用指标证明系统是否真的达到目标

2. 先做技术选型:BIO、虚拟线程、NIO 与 Netty

2.1 不要把 NIO 写成“高并发唯一解”

在现代 Java 中,高并发网络服务至少存在四类可行方案:

方案 编程模型 优点 主要代价 典型场景
阻塞 I/O + 平台线程 同步 最简单 线程成本高 低并发内部工具
阻塞 I/O + 虚拟线程 同步 代码直观,适合大量阻塞任务 仍需控制外部资源与内存,精细 I/O 调度能力较弱 HTTP/RPC 业务服务、连接数较高但协议处理简单
Java NIO + Reactor 异步事件驱动 连接密度高,I/O 调度和内存可精细控制 实现复杂,容易出现并发与缓冲区错误 网关、IM、设备接入、自定义协议
Netty 成熟事件驱动框架 编解码、内存池、TLS、原生传输、可观测性生态完善 需要理解 Netty 线程模型和引用计数 绝大多数生产级异步网络系统

虚拟线程显著降低了“一任务一线程”模型的成本,但它并不会自动解决以下问题:

  • 协议帧边界
  • 单连接有序性
  • 写队列无限增长
  • 对端慢读
  • 全局内存预算
  • 连接洪泛与慢速攻击
  • 下游资源容量限制

因此,技术选择应由协议复杂度、连接密度、延迟目标、团队能力和生态需求共同决定。

2.2 什么时候应该直接使用 Netty

满足以下任意条件时,通常优先使用 Netty:

  • 需要 TLS、HTTP/2、WebSocket、MQTT、Protobuf 等成熟协议支持
  • 需要池化 ByteBuf、引用计数和零拷贝切片
  • 需要 epoll、kqueue 等原生传输
  • 需要完善的 ChannelPipeline、IdleState、流量整形与编解码器
  • 项目没有长期维护底层通信框架的专职团队

自研 Java NIO 的合理目标通常是:

  • 学习和验证 Reactor 核心机制
  • 构建极窄场景的专用传输层
  • 对消息布局、内存与调度有特殊要求
  • 对第三方依赖或运行环境有严格限制

3. Java NIO 的核心对象与操作系统映射

3.1 Channel

ServerSocketChannel 负责监听端口并接收连接,SocketChannel 表示一个 TCP 连接。Channel 必须设置为非阻塞模式,才能注册到 Selector

serverChannel.configureBlocking(false);
socketChannel.configureBlocking(false);

3.2 Selector

Selector 是就绪事件的聚合器。一个 EventLoop 通常持有一个 Selector,循环执行:

执行跨线程任务
   ↓
select 等待事件
   ↓
处理 selectedKeys
   ↓
执行定时任务或维护任务
   ↓
进入下一轮

Linux 上,JDK 默认实现通常利用 epoll。Java API 不保证应用必须感知具体实现,因此业务代码不应依赖内部类名。

3.3 SelectionKey

Channel 注册到 Selector 后会得到 SelectionKey,常用事件包括:

事件 含义
OP_ACCEPT 监听 Channel 有新连接可接收
OP_CONNECT 非阻塞连接建立过程完成
OP_READ 连接当前可能可读
OP_WRITE 连接当前可能可写

“就绪”并不等于“一定能完成完整读写”。例如:

  • OP_READ 触发后,单次 read() 可能只读到一部分消息
  • OP_WRITE 触发后,单次 write() 可能只写出部分缓冲区
  • Channel 可能已被另一条关闭路径取消

3.4 ByteBuffer

ByteBuffer 的核心状态是:

  • position:下一次读或写的位置
  • limit:当前可访问边界
  • capacity:容量上限

典型读流程:

int read = channel.read(buffer); // 写入 buffer
buffer.flip();                   // 切换为读取模式
consume(buffer);
buffer.compact();                // 保留未消费字节,切回写入模式

处理 TCP 累积缓冲区时,通常应使用 compact(),而不是无条件 clear()clear() 会丢弃尚未消费的半包数据。

4. Reactor 模型的三种形态

4.1 单 Reactor 单线程

所有 I/O 和业务处理均由一个线程完成。实现简单,但任何慢操作都会阻塞整个事件循环。

适合:

  • 教学示例
  • 连接数很少的工具
  • 业务逻辑近乎零成本的代理原型

不适合:

  • 数据库访问
  • 远程调用
  • 压缩、加密、复杂序列化
  • 延迟不可控的业务逻辑

4.2 单 Reactor 多线程

I/O 事件由一个 Reactor 线程处理,业务任务提交给线程池。该模型解决了业务阻塞问题,但单个 Selector 仍承担全部连接 I/O。

4.3 主从 Reactor 多线程

职责划分:

  • Main Reactor / Acceptor:只处理连接接入和基础参数设置
  • SubReactor / EventLoop:负责连接注册、读事件、写事件、连接关闭和状态变更
  • Business Executor:执行可能阻塞或耗时的业务逻辑

需要特别强调:

主从 Reactor 并不要求多个监听 Socket 同时绑定同一端口。最常见的实现是一个 Acceptor 接收连接,再将连接轮询分配到多个 SubReactor。

SO_REUSEPORT 可以在支持该选项的系统上实现多个监听 Socket 绑定同一地址与端口,但其语义依赖操作系统,不能作为跨平台默认前提。

5. 生产级主从 Reactor 总体架构

5.1 分层结构

5.2 连接归属模型

每个已连接 SocketChannel 在生命周期内固定归属于一个 EventLoop:

connectionId → owner EventLoop → Selector → SelectionKey → ConnectionContext

除非实现复杂的连接迁移机制,否则不要在运行中把 Channel 从一个 Selector 迁移到另一个 Selector。

5.3 每个连接应维护的状态

final class ConnectionContext {
    final SocketChannel channel;
    final EventLoop owner;
    final ByteBuffer readBuffer;
    final Deque<ByteBuffer> outboundQueue;

    SelectionKey key;
    long queuedOutboundBytes;
    boolean readSuspended;
    boolean closed;
    long lastReadNanos;
    long lastWriteNanos;
}

实际项目还可加入:

  • 连接 ID
  • 对端地址
  • 登录态或会话信息
  • 协议版本
  • 解码器状态
  • 单连接请求序号
  • 流控窗口
  • TLS 状态
  • 心跳与空闲超时状态

6. 必须遵守的八条并发不变量

不变量 1:连接 I/O 状态只由所属 EventLoop 修改

SelectionKey.interestOps()、写队列、解码累积区、关闭状态等,应在 EventLoop 线程内串行修改。

不变量 2:跨线程操作必须先入队,再 wakeup

错误方式:

// 业务线程直接修改 SelectionKey
key.interestOps(key.interestOps() | SelectionKey.OP_WRITE);

推荐方式:

owner.execute(() -> enableWrite(key));

EventLoop 的 execute() 负责:

  1. 将任务放入线程安全队列
  2. 当调用方不是 EventLoop 自身时执行 selector.wakeup()
  3. 由 EventLoop 串行执行任务

不变量 3:注册 Channel 由目标 EventLoop 完成

错误方式:

// Acceptor 线程直接向另一个线程正在 select 的 Selector 注册
channel.register(subReactor.selector(), SelectionKey.OP_READ);
subReactor.selector().wakeup();

推荐方式:

subReactor.register(channel); // 内部入队并 wakeup

不变量 4:异步业务不能持有即将复用的读缓冲区

I/O 线程把 ByteBuffer 交给业务线程后,如果立即 clear() 并继续读取,就会产生数据覆盖和竞态。

正确策略包括:

  • 解码为不可变消息对象后再提交业务线程
  • 复制完整帧到独立缓冲区
  • 使用带引用计数的池化缓冲区,并严格管理生命周期

不变量 5:单次 read/write 不代表完整消息完成

必须处理:

  • 零字节读写
  • 半包
  • 粘包
  • 半写
  • 对端关闭

不变量 6:OP_WRITE 默认关闭

只要写队列为空,就必须移除 OP_WRITE。否则多数 Socket 长时间处于可写状态,Selector 会持续返回,造成空转和 CPU 飙升。

不变量 7:队列和内存必须有硬上限

至少需要以下限制:

  • 业务线程池队列上限
  • 单连接待发送字节上限
  • 全局待发送字节上限
  • 单帧最大长度
  • 单连接读累积区上限
  • 单轮 EventLoop 最大任务数或最大执行时长

不变量 8:关闭必须幂等

读异常、写异常、业务超时和主动断开可能同时触发关闭。关闭逻辑必须保证重复调用无副作用。

7. 协议设计:先解决边界,再谈业务

7.1 TCP 是字节流,不保留消息边界

发送方连续执行两次 write(),接收方可能:

  • 一次读到两条完整消息
  • 第一次只读到半条
  • 第一次读到一条半
  • 多次才能读完一条

因此,自定义协议必须定义明确的帧格式。

7.2 推荐帧格式

+---------+---------+---------+---------+----------+-------------+
| Magic   | Version | Flags   | MsgType | BodyLen  | Body        |
| 2 bytes | 1 byte  | 1 byte  | 2 bytes | 4 bytes  | N bytes     |
+---------+---------+---------+---------+----------+-------------+

建议字段:

字段 作用
Magic 快速识别非法协议流量
Version 支持协议兼容与灰度升级
Flags 压缩、加密、响应类型等标志
MsgType 消息类型或命令字
BodyLen 消息体长度
Body Protobuf、JSON 或自定义编码

7.3 长度校验

解码器必须在分配内存之前校验:

if (bodyLength < 0 || bodyLength > maxFrameLength) {
    throw new ProtocolException("invalid body length: " + bodyLength);
}

不能因为报文声明了 1 GB 长度,就真的申请 1 GB 缓冲区。

7.4 字节序

协议必须明确使用大端或小端。网络协议通常使用大端:

buffer.order(ByteOrder.BIG_ENDIAN);

8. 核心源码实现

下面给出一套精简但结构正确的参考实现。它重点展示主从 Reactor、线程安全注册、半写处理和背压入口。生产项目仍应补充 TLS、鉴权、监控、配置热更新与更完善的内存池。

8.1 服务配置

public record ReactorConfig(
        int port,
        int backlog,
        int ioThreads,
        int businessThreads,
        int businessQueueCapacity,
        int readBufferBytes,
        int maxFrameBytes,
        long outboundHighWatermark,
        long outboundLowWatermark,
        long outboundHardLimit
) {
    public ReactorConfig {
        if (port < 1 || port > 65535) {
            throw new IllegalArgumentException("invalid port");
        }
        if (ioThreads < 1 || businessThreads < 1) {
            throw new IllegalArgumentException("thread count must be positive");
        }
        if (readBufferBytes < Integer.BYTES + maxFrameBytes) {
            throw new IllegalArgumentException(
                    "read buffer must hold at least one complete frame"
            );
        }
        if (!(0 < outboundLowWatermark
                && outboundLowWatermark < outboundHighWatermark
                && outboundHighWatermark < outboundHardLimit)) {
            throw new IllegalArgumentException("invalid outbound watermarks");
        }
    }

    public static ReactorConfig defaults(int port) {
        int cpu = Runtime.getRuntime().availableProcessors();
        return new ReactorConfig(
                port,
                4096,
                Math.max(1, cpu),
                Math.max(4, cpu * 2),
                10_000,
                8 * 1024,
                4 * 1024,
                2L * 1024 * 1024,
                1L * 1024 * 1024,
                8L * 1024 * 1024
        );
    }
}

这里的默认值只是启动基线,不是生产环境固定答案。最终参数必须由消息大小、连接数、业务耗时、容器 CPU 配额和压测结果决定。

8.2 基础辅助类

public final class ProtocolException extends RuntimeException {
    private static final long serialVersionUID = 1L;

    public ProtocolException(String message) {
        super(message);
    }
}

public final class NamedThreadFactory implements ThreadFactory {
    private final String prefix;
    private final AtomicInteger sequence = new AtomicInteger();

    public NamedThreadFactory(String prefix) {
        this.prefix = Objects.requireNonNull(prefix, "prefix");
    }

    @Override
    public Thread newThread(Runnable task) {
        Thread thread = new Thread(task, prefix + "-" + sequence.incrementAndGet());
        thread.setUncaughtExceptionHandler((t, error) -> {
            // 接入统一日志与告警系统。
        });
        return thread;
    }
}

8.3 有界业务线程池

public final class BusinessExecutor implements AutoCloseable {
    private final ThreadPoolExecutor executor;

    public BusinessExecutor(int threads, int queueCapacity) {
        this.executor = new ThreadPoolExecutor(
                threads,
                threads,
                0L,
                TimeUnit.MILLISECONDS,
                new ArrayBlockingQueue<>(queueCapacity),
                new NamedThreadFactory("business"),
                new ThreadPoolExecutor.AbortPolicy()
        );
    }

    public boolean submit(Runnable task) {
        try {
            executor.execute(task);
            return true;
        } catch (RejectedExecutionException rejected) {
            return false;
        }
    }

    public int queueSize() {
        return executor.getQueue().size();
    }

    @Override
    public void close() {
        executor.shutdown();
    }
}

不要默认使用无界 LinkedBlockingQueue。无界队列会把短时过载转化为长时间延迟和内存风险。

8.4 EventLoop

public final class EventLoop implements Runnable, AutoCloseable {
    private static final int MAX_TASKS_PER_TICK = 1024;

    private final Selector selector;
    private final Queue<Runnable> taskQueue = new ConcurrentLinkedQueue<>();
    private final Thread thread;
    private final ReactorConfig config;
    private final BusinessExecutor businessExecutor;
    private final AtomicBoolean running = new AtomicBoolean(true);

    public EventLoop(
            String name,
            ReactorConfig config,
            BusinessExecutor businessExecutor
    ) throws IOException {
        this.selector = Selector.open();
        this.config = config;
        this.businessExecutor = businessExecutor;
        this.thread = new Thread(this, name);
    }

    public void start() {
        thread.start();
    }

    public boolean inEventLoop() {
        return Thread.currentThread() == thread;
    }

    public void execute(Runnable task) {
        Objects.requireNonNull(task, "task");
        taskQueue.offer(task);
        if (!inEventLoop()) {
            selector.wakeup();
        }
    }

    public void register(SocketChannel channel) {
        execute(() -> {
            try {
                ConnectionContext context = new ConnectionContext(
                        channel,
                        this,
                        config,
                        businessExecutor
                );
                SelectionKey key = channel.register(
                        selector,
                        SelectionKey.OP_READ,
                        context
                );
                context.attachKey(key);
            } catch (IOException e) {
                closeQuietly(channel);
            }
        });
    }

    @Override
    public void run() {
        while (running.get()) {
            try {
                runTasks(MAX_TASKS_PER_TICK);

                int selected = selector.select(1000);
                if (selected > 0) {
                    processSelectedKeys();
                }

                runTasks(MAX_TASKS_PER_TICK);
            } catch (IOException e) {
                // 记录 EventLoop 级异常。是否退出由故障策略决定。
            } catch (Throwable unexpected) {
                // 防止单连接异常杀死整个 EventLoop。
            }
        }

        shutdownChannels();
        closeQuietly(selector);
    }

    private void processSelectedKeys() {
        Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
        while (iterator.hasNext()) {
            SelectionKey key = iterator.next();
            iterator.remove();

            if (!key.isValid()) {
                continue;
            }

            ConnectionContext context = (ConnectionContext) key.attachment();
            try {
                if (key.isReadable()) {
                    context.onReadable();
                }
                if (key.isValid() && key.isWritable()) {
                    context.onWritable();
                }
            } catch (Throwable error) {
                context.close(error);
            }
        }
    }

    private void runTasks(int maxTasks) {
        for (int i = 0; i < maxTasks; i++) {
            Runnable task = taskQueue.poll();
            if (task == null) {
                return;
            }
            try {
                task.run();
            } catch (Throwable ignored) {
                // 单个任务失败不能终止 EventLoop。
            }
        }
    }

    private void shutdownChannels() {
        for (SelectionKey key : selector.keys()) {
            Object attachment = key.attachment();
            if (attachment instanceof ConnectionContext context) {
                context.close(null);
            } else {
                closeQuietly(key.channel());
            }
        }
    }

    @Override
    public void close() {
        if (running.compareAndSet(true, false)) {
            selector.wakeup();
        }
    }

    private static void closeQuietly(AutoCloseable closeable) {
        try {
            closeable.close();
        } catch (Exception ignored) {
        }
    }
}

设计重点:

  • register() 不在 Acceptor 线程直接操作目标 Selector
  • interestOps、写队列和连接状态都应通过 execute() 回到所属 EventLoop
  • 每轮限制任务数量,避免跨线程任务持续涌入导致 I/O 饥饿
  • select(timeout) 让 EventLoop 即使遗漏唤醒,也能周期性执行维护任务

8.5 EventLoopGroup

public final class EventLoopGroup implements AutoCloseable {
    private final EventLoop[] loops;
    private final AtomicInteger cursor = new AtomicInteger();

    public EventLoopGroup(
            int size,
            ReactorConfig config,
            BusinessExecutor businessExecutor
    ) throws IOException {
        this.loops = new EventLoop[size];
        for (int i = 0; i < size; i++) {
            loops[i] = new EventLoop(
                    "nio-event-loop-" + i,
                    config,
                    businessExecutor
            );
        }
    }

    public void start() {
        for (EventLoop loop : loops) {
            loop.start();
        }
    }

    public EventLoop next() {
        int index = Math.floorMod(cursor.getAndIncrement(), loops.length);
        return loops[index];
    }

    @Override
    public void close() {
        for (EventLoop loop : loops) {
            loop.close();
        }
    }
}

简单轮询足以作为初始分配策略。更复杂的实现可以根据以下指标加权选择:

  • 当前连接数
  • EventLoop 任务队列长度
  • 最近 I/O 延迟
  • 待发送字节数

但要避免为了“绝对均衡”引入全局锁和高频统计开销。

8.6 Acceptor

public final class Acceptor implements Runnable, AutoCloseable {
    private final ServerSocketChannel serverChannel;
    private final Selector selector;
    private final EventLoopGroup workers;
    private final AtomicBoolean running = new AtomicBoolean(true);
    private final Thread thread;

    public Acceptor(ReactorConfig config, EventLoopGroup workers) throws IOException {
        this.workers = workers;
        this.selector = Selector.open();
        this.serverChannel = ServerSocketChannel.open();
        this.serverChannel.configureBlocking(false);
        this.serverChannel.setOption(StandardSocketOptions.SO_REUSEADDR, true);
        this.serverChannel.bind(
                new InetSocketAddress(config.port()),
                config.backlog()
        );
        this.serverChannel.register(selector, SelectionKey.OP_ACCEPT);
        this.thread = new Thread(this, "nio-acceptor");
    }

    public void start() {
        thread.start();
    }

    @Override
    public void run() {
        while (running.get()) {
            try {
                selector.select(1000);
                Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
                while (iterator.hasNext()) {
                    SelectionKey key = iterator.next();
                    iterator.remove();

                    if (key.isValid() && key.isAcceptable()) {
                        acceptBatch();
                    }
                }
            } catch (IOException e) {
                if (running.get()) {
                    // 记录并根据策略决定是否继续。
                }
            }
        }

        closeQuietly(serverChannel);
        closeQuietly(selector);
    }

    private void acceptBatch() throws IOException {
        while (true) {
            SocketChannel channel = serverChannel.accept();
            if (channel == null) {
                return;
            }

            try {
                channel.configureBlocking(false);
                channel.setOption(StandardSocketOptions.TCP_NODELAY, true);
                channel.setOption(StandardSocketOptions.SO_KEEPALIVE, true);
                workers.next().register(channel);
            } catch (Throwable error) {
                closeQuietly(channel);
            }
        }
    }

    @Override
    public void close() {
        if (running.compareAndSet(true, false)) {
            selector.wakeup();
        }
    }

    private static void closeQuietly(AutoCloseable closeable) {
        try {
            closeable.close();
        } catch (Exception ignored) {
        }
    }
}

acceptBatch() 会持续接收,直到 accept() 返回 null。这样可在一次就绪通知中尽量排空待接入连接,但生产环境还可以增加单轮 Accept 上限,防止连接洪峰让已建立连接长时间得不到处理。

8.7 帧解码器

public final class LengthFieldFrameDecoder {
    private static final int LENGTH_FIELD_BYTES = Integer.BYTES;
    private final int maxFrameBytes;

    public LengthFieldFrameDecoder(int maxFrameBytes) {
        this.maxFrameBytes = maxFrameBytes;
    }

    public List<byte[]> decode(ByteBuffer input) {
        List<byte[]> frames = new ArrayList<>();
        input.flip();

        while (input.remaining() >= LENGTH_FIELD_BYTES) {
            input.mark();
            int bodyLength = input.getInt();

            if (bodyLength < 0 || bodyLength > maxFrameBytes) {
                throw new ProtocolException("invalid frame length: " + bodyLength);
            }

            if (input.remaining() < bodyLength) {
                input.reset();
                break;
            }

            byte[] body = new byte[bodyLength];
            input.get(body);
            frames.add(body);
        }

        input.compact();
        return frames;
    }
}

该实现把完整帧复制为独立 byte[],因此可以安全提交给异步业务线程。其代价是一次内存复制。追求更高性能时,可以改为池化缓冲区切片,但必须引入明确的引用计数和释放协议。

8.8 ConnectionContext

public final class ConnectionContext {
    private final SocketChannel channel;
    private final EventLoop owner;
    private final BusinessExecutor businessExecutor;
    private final ReactorConfig config;
    private final LengthFieldFrameDecoder frameDecoder;
    private final ByteBuffer readBuffer;
    private final Deque<ByteBuffer> outboundQueue = new ArrayDeque<>();

    private SelectionKey key;
    private long queuedOutboundBytes;
    private boolean readSuspended;
    private boolean closed;

    public ConnectionContext(
            SocketChannel channel,
            EventLoop owner,
            ReactorConfig config,
            BusinessExecutor businessExecutor
    ) {
        this.channel = channel;
        this.owner = owner;
        this.config = config;
        this.businessExecutor = businessExecutor;
        this.frameDecoder = new LengthFieldFrameDecoder(config.maxFrameBytes());
        this.readBuffer = ByteBuffer.allocateDirect(config.readBufferBytes());
    }

    public void attachKey(SelectionKey key) {
        assertOwnerThread();
        this.key = key;
    }

    public void onReadable() throws IOException {
        assertOwnerThread();
        if (closed || readSuspended) {
            return;
        }

        int totalRead = 0;
        while (readBuffer.hasRemaining()) {
            int read = channel.read(readBuffer);
            if (read > 0) {
                totalRead += read;
                continue;
            }
            if (read == 0) {
                break;
            }
            close(null);
            return;
        }

        if (totalRead == 0) {
            return;
        }

        List<byte[]> frames = frameDecoder.decode(readBuffer);
        for (byte[] frame : frames) {
            dispatch(frame);
            if (closed) {
                break;
            }
        }

        if (!readBuffer.hasRemaining()) {
            // 累积区被填满且仍无法形成完整帧,说明协议异常或缓冲区配置不足。
            close(new ProtocolException("read cumulation overflow"));
        }
    }

    private void dispatch(byte[] frame) {
        boolean accepted = businessExecutor.submit(() -> {
            try {
                byte[] response = handleBusiness(frame);
                if (response != null) {
                    send(encodeFrame(response));
                }
            } catch (Throwable error) {
                owner.execute(() -> close(error));
            }
        });

        if (!accepted) {
            // 精简实现选择立即关闭,避免已经解码的后续帧被静默丢弃。
            // 完整实现也可维护连接级 pendingInbound 队列,待业务队列恢复后继续提交。
            close(new RejectedExecutionException("business executor overloaded"));
        }
    }

    public void send(ByteBuffer data) {
        owner.execute(() -> enqueueWrite(data));
    }

    private void enqueueWrite(ByteBuffer data) {
        assertOwnerThread();
        if (closed) {
            return;
        }

        int bytes = data.remaining();
        long next = queuedOutboundBytes + bytes;
        if (next > config.outboundHardLimit()) {
            close(new IllegalStateException("outbound hard limit exceeded"));
            return;
        }

        outboundQueue.addLast(data);
        queuedOutboundBytes = next;
        enableWrite();

        if (queuedOutboundBytes >= config.outboundHighWatermark()) {
            suspendRead();
        }
    }

    public void onWritable() throws IOException {
        assertOwnerThread();

        while (!outboundQueue.isEmpty()) {
            ByteBuffer current = outboundQueue.peekFirst();
            int before = current.remaining();
            int written = channel.write(current);
            int consumed = before - current.remaining();
            queuedOutboundBytes -= consumed;

            if (current.hasRemaining()) {
                if (written == 0) {
                    break;
                }
                continue;
            }

            outboundQueue.removeFirst();
        }

        if (outboundQueue.isEmpty()) {
            disableWrite();
        }

        if (readSuspended
                && queuedOutboundBytes <= config.outboundLowWatermark()
                && businessExecutor.queueSize() < config.businessQueueCapacity() / 2) {
            resumeRead();
        }
    }

    private void suspendRead() {
        assertOwnerThread();
        if (!readSuspended && key != null && key.isValid()) {
            readSuspended = true;
            key.interestOps(key.interestOps() & ~SelectionKey.OP_READ);
        }
    }

    private void resumeRead() {
        assertOwnerThread();
        if (readSuspended && key != null && key.isValid()) {
            readSuspended = false;
            key.interestOps(key.interestOps() | SelectionKey.OP_READ);
        }
    }

    private void enableWrite() {
        if (key != null && key.isValid()) {
            key.interestOps(key.interestOps() | SelectionKey.OP_WRITE);
        }
    }

    private void disableWrite() {
        if (key != null && key.isValid()) {
            key.interestOps(key.interestOps() & ~SelectionKey.OP_WRITE);
        }
    }

    public void close(Throwable cause) {
        if (!owner.inEventLoop()) {
            owner.execute(() -> close(cause));
            return;
        }
        if (closed) {
            return;
        }

        closed = true;
        if (key != null) {
            key.cancel();
        }
        outboundQueue.clear();
        queuedOutboundBytes = 0;

        try {
            channel.close();
        } catch (IOException ignored) {
        }
    }

    private void assertOwnerThread() {
        if (!owner.inEventLoop()) {
            throw new IllegalStateException("must run in owner EventLoop");
        }
    }

    private static ByteBuffer encodeFrame(byte[] body) {
        ByteBuffer buffer = ByteBuffer.allocate(Integer.BYTES + body.length);
        buffer.putInt(body.length);
        buffer.put(body);
        return buffer.flip();
    }

    private static byte[] handleBusiness(byte[] request) {
        // 示例:生产环境调用真正的业务服务。
        return request;
    }
}

这段代码体现了几个生产级关键点:

  • 业务线程只处理独立消息对象
  • 业务完成后通过 owner.execute() 返回 EventLoop
  • 写队列只由 EventLoop 修改
  • 正确处理半写
  • 待发送字节达到高水位后暂停读
  • 下降到低水位后才恢复读,避免频繁抖动
  • 超过硬上限时关闭连接,避免单个慢客户端拖垮实例

8.9 Server Bootstrap

public final class ReactorServer implements AutoCloseable {
    private final BusinessExecutor businessExecutor;
    private final EventLoopGroup workers;
    private final Acceptor acceptor;

    public ReactorServer(ReactorConfig config) throws IOException {
        this.businessExecutor = new BusinessExecutor(
                config.businessThreads(),
                config.businessQueueCapacity()
        );
        this.workers = new EventLoopGroup(
                config.ioThreads(),
                config,
                businessExecutor
        );
        this.acceptor = new Acceptor(config, workers);
    }

    public void start() {
        workers.start();
        acceptor.start();
    }

    @Override
    public void close() {
        acceptor.close();
        workers.close();
        businessExecutor.close();
    }

    public static void main(String[] args) throws Exception {
        ReactorConfig config = ReactorConfig.defaults(9000);
        ReactorServer server = new ReactorServer(config);
        Runtime.getRuntime().addShutdownHook(new Thread(server::close));
        server.start();
    }
}

真实系统应提供阻塞等待、启动失败回滚、分阶段停机和终止超时控制,而不是仅依赖 JVM Shutdown Hook。

9. 读路径:从字节流到业务消息

9.1 推荐处理链

9.2 单轮读取预算

不能在一个连接上无限循环读取,否则高流量连接可能独占 EventLoop。可以设置:

  • 单轮最大读取字节数
  • 单轮最大帧数
  • 单轮最大处理时间

例如:

int readBudget = 256 * 1024;
int consumed = 0;
while (consumed < readBudget) {
    int n = channel.read(buffer);
    if (n <= 0) {
        break;
    }
    consumed += n;
}

9.3 业务有序性

同一连接连续发送请求 A、B 时,业务线程池可能先完成 B。是否允许乱序取决于协议:

  • 请求携带唯一 ID,响应允许乱序:可直接并发执行
  • 协议要求严格顺序:需要连接级串行执行器
  • 只有部分命令要求顺序:按会话键或业务键分片执行

连接级串行执行器可采用“每连接任务队列 + 原子运行标记”,避免为每个连接创建独立线程。

10. 写路径:半写、队列与 OP_WRITE

10.1 为什么会半写

非阻塞 SocketChannel.write() 只负责尽可能写入内核发送缓冲区。对端读取慢、网络拥塞或发送缓冲区不足时,返回值可能是:

  • 大于 0:只写入部分或全部字节
  • 等于 0:当前无法继续写
  • 抛出异常:连接失效

因此不能这样写:

channel.write(buffer);
// 错误地认为 buffer 已全部发送

10.2 正确状态机

业务产生响应
   ↓
进入连接 outboundQueue
   ↓
尝试立即写出
   ├─ 全部写完 → 队列为空 → 关闭 OP_WRITE
   └─ 仍有剩余 → 保留 ByteBuffer position → 开启 OP_WRITE
                         ↓
                 下次可写事件继续

10.3 聚合写

大量小包可以使用 gathering write:

ByteBuffer[] buffers = batch.toArray(ByteBuffer[]::new);
long written = channel.write(buffers);

需要权衡:

  • 降低系统调用次数
  • 不要无限聚合导致单连接独占 EventLoop
  • 仍要更新每个 Buffer 的 position 和队列字节数

10.4 写队列不能只限制条数

不同消息大小差异很大,限制“最多 1000 条消息”没有足够意义。应至少限制:

  • 单连接待发送字节数
  • 全局待发送字节数
  • 单消息最大字节数

11. 背压与过载保护

11.1 为什么只限制线程池队列不够

一个完整系统存在多级队列:

内核接收缓冲区
→ Java 读缓冲区
→ 解码完成消息
→ 业务线程池队列
→ 下游连接池或数据库队列
→ 响应写队列
→ 内核发送缓冲区

任何一层无限增长都可能导致内存耗尽或尾延迟雪崩。

11.2 四层背压

第一层:连接级读暂停

当某连接待发送数据超过高水位时,移除该连接的 OP_READ

第二层:业务执行器拒绝

业务队列满时,可按协议选择:

  • 暂停该连接读取
  • 返回 SERVER_BUSY
  • 丢弃低优先级消息
  • 关闭异常连接

第三层:全局内存预算

使用原子计数器维护全局待发送字节:

if (!globalBudget.tryReserve(bytes)) {
    rejectOrClose();
}

消息写完或连接关闭时必须归还预算。

第四层:入口限流

按以下维度限制:

  • 新建连接速率
  • 单 IP 连接数
  • 单租户请求速率
  • 单连接每秒消息数
  • 未认证连接存活时间

11.3 高低水位而不是单阈值

若只设置一个阈值,队列会在边界附近反复暂停和恢复。采用高低水位:

达到 High Watermark → 暂停读
下降到 Low Watermark → 恢复读

这与电路中的迟滞思想相同,可以显著减少状态抖动。

12. 缓冲区与堆外内存治理

12.1 HeapBuffer 与 DirectBuffer

类型 优点 缺点
Heap ByteBuffer 分配便宜,受 GC 管理,调试简单 Socket I/O 可能需要额外拷贝
Direct ByteBuffer 更适合底层 I/O,减少部分中间拷贝 分配释放成本高,泄漏排查更困难,受直接内存上限约束

不要简单地得出“DirectBuffer 一定更快”。对于小消息和短生命周期对象,堆内缓冲区可能更经济。必须通过目标负载压测。

12.2 连接级固定读缓冲区的内存预算

如果每连接固定分配 64 KiB:

100,000 connections × 64 KiB ≈ 6.1 GiB

这还没有计算写队列、协议对象、TLS 状态和业务缓存。因此,十万连接不能机械地给每个连接分配大块 DirectBuffer。

常见策略:

  • 小型连接级基础缓冲区,例如 4 KiB 或 8 KiB
  • 需要时扩容,设置最大上限
  • 采用分级池,如 4 KiB、16 KiB、64 KiB
  • 大帧采用分段缓冲或流式处理
  • 空闲连接回收大缓冲区,降级到基础容量

12.3 缓冲池要求

一个可用的缓冲池至少要考虑:

  • 容量分级
  • 最大缓存数量
  • 跨线程归还
  • 重复归还检测
  • 泄漏检测
  • 直接内存预算
  • 连接关闭时释放

简单的 ConcurrentLinkedQueue<ByteBuffer> 只能作为演示,不能自动解决碎片、无限缓存和泄漏问题。

12.4 零拷贝的准确理解

FileChannel.transferTo() 允许操作系统采用高效路径,将文件系统缓存中的数据传输到目标 Channel。是否完全避免复制以及具体实现路径由操作系统、文件系统、目标 Channel 和 JDK 实现共同决定。

正确的非阻塞发送还要处理返回 0:

long transferred = fileChannel.transferTo(position, remaining, socketChannel);
if (transferred == 0) {
    // 当前无法继续,等待后续 OP_WRITE,而不是死循环。
}

TLS 加密场景下,明文文件通常不能直接通过同一路径发送到网络,因为数据还需要进入加密处理流程。

13. 线程数与任务调度策略

13.1 I/O 线程数不是固定的 CPU × 2

I/O EventLoop 线程主要执行:

  • 读取和写出
  • 协议拆帧
  • 轻量编解码
  • 状态维护
  • 任务队列执行

推荐从以下基线开始:

I/O 线程数 ≈ 容器可用 CPU 核数

再根据压测调整。若 EventLoop 执行了较重的压缩、加密或序列化,可能需要增加线程或拆出专用执行器,但更重要的是不要阻塞 EventLoop。

13.2 业务线程数

CPU 密集任务可从接近 CPU 核数开始。存在阻塞 I/O 的任务可使用更多平台线程,或在 Java 21 及以上评估虚拟线程执行器:

try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) {
    executor.submit(() -> callBlockingService());
}

但虚拟线程不应被无边界地理解为“无需限流”。数据库连接、远程服务、内存和文件描述符仍然是有限资源,因此入口并发仍需通过信号量、连接池和超时控制。

13.3 EventLoop 公平性

单轮循环应在以下工作之间保持平衡:

  • 已就绪 I/O
  • 跨线程任务
  • 定时任务
  • 维护任务

常见策略:

  • 限制每轮任务数
  • 限制单连接读写预算
  • 统计 EventLoop 任务等待时间
  • 检测长任务并告警

14. 连接生命周期与异常处理

14.1 状态机

14.2 对端正常关闭

channel.read() 返回 -1 表示对端已关闭输出方向。是否立即关闭整个连接取决于协议是否支持半关闭,但多数应用协议可直接关闭。

14.3 异常分类

建议区分:

类型 示例 处理
协议错误 非法长度、错误 Magic 记录采样日志并关闭
对端网络错误 Connection reset 低级别日志或计数
业务错误 参数校验失败 返回业务错误,不一定断开
过载错误 队列满、内存预算不足 限流、忙响应或关闭
框架错误 EventLoop 未捕获异常 高优先级告警

避免每次连接重置都打印完整堆栈,否则异常流量会反过来压垮日志系统。

14.4 空闲连接与心跳

需要区分:

  • 读空闲:长时间未收到任何数据
  • 写空闲:长时间未发送数据
  • 全空闲:读写均无活动
  • 应用心跳超时:已发送 Ping,但未收到 Pong

对于大量连接,不能为每个连接创建一个独立高精度定时器。可使用:

  • 时间轮
  • 分桶扫描
  • EventLoop 内最小堆
  • 周期性批量检查

15. 网络参数与 Kubernetes 部署

15.1 应用层 Socket 参数

常见参数:

channel.setOption(StandardSocketOptions.TCP_NODELAY, true);
channel.setOption(StandardSocketOptions.SO_KEEPALIVE, true);
channel.setOption(StandardSocketOptions.SO_RCVBUF, receiveBufferBytes);
channel.setOption(StandardSocketOptions.SO_SNDBUF, sendBufferBytes);

注意:

  • TCP_NODELAY 是否开启取决于小包延迟与带宽效率的权衡
  • SO_KEEPALIVE 的探测周期仍由操作系统参数决定
  • 设置的收发缓冲区值可能被内核调整
  • 不应盲目把缓冲区设置得极大,因为连接数乘积会形成显著内存占用

15.2 Backlog

应用传入的 backlog 只是请求值,实际还受操作系统和容器环境限制。需要联合检查:

  • 应用 bind(address, backlog)
  • net.core.somaxconn
  • SYN 队列相关参数
  • 负载均衡器与 Service 的连接行为

15.3 SO_REUSEPORT

SO_REUSEPORT 的支持情况与语义依赖操作系统。使用前应检查:

if (serverChannel.supportedOptions().contains(StandardSocketOptions.SO_REUSEPORT)) {
    serverChannel.setOption(StandardSocketOptions.SO_REUSEPORT, true);
}

它适合在特定场景下让多个监听 Socket 分担接入,但会增加端口绑定、连接分配和故障排查复杂度。一个 Acceptor 通常足以接收大量已建立连接,真正瓶颈需要用指标确认。

15.4 不要机械设置 tcp_tw_reuse

TIME_WAIT 主要出现在主动关闭端。是否调整相关内核参数取决于:

  • 服务是长连接还是短连接
  • 哪一端主动关闭
  • 是否存在 NAT 或四层负载均衡
  • 内核版本与参数语义
  • 端口耗尽是否真实发生

没有证据时,不应把 net.ipv4.tcp_tw_reuse=1 当作通用优化模板。

15.5 Kubernetes 资源配置

示例:

apiVersion: apps/v1
kind: Deployment
metadata:
  name: nio-gateway
spec:
  replicas: 3
  selector:
    matchLabels:
      app: nio-gateway
  template:
    metadata:
      labels:
        app: nio-gateway
    spec:
      terminationGracePeriodSeconds: 60
      containers:
        - name: gateway
          image: example/nio-gateway:2.0.0
          ports:
            - name: tcp
              containerPort: 9000
          resources:
            requests:
              cpu: "2"
              memory: 2Gi
            limits:
              cpu: "4"
              memory: 4Gi
          readinessProbe:
            httpGet:
              path: /ready
              port: 8080
            periodSeconds: 3
            failureThreshold: 2
          livenessProbe:
            httpGet:
              path: /live
              port: 8080
            periodSeconds: 10
            failureThreshold: 3
          lifecycle:
            preStop:
              httpGet:
                path: /management/drain
                port: 8080

关键点:

  • Readiness 表示是否接收新流量
  • Liveness 只检测进程是否需要重启,不应依赖外部下游
  • preStop 先进入摘流状态
  • terminationGracePeriodSeconds 必须覆盖摘流传播、连接排空与进程关闭时间
  • CPU limit 过低或频繁节流会明显放大 EventLoop 延迟

15.6 长连接与滚动发布

长连接不会因为 Pod Readiness 变为 false 就自动迁移。滚动升级需要明确策略:

  • 停止接受新连接
  • 向客户端发送服务端下线通知
  • 客户端随机退避重连
  • 设置最大排空时间
  • 到期后关闭剩余连接
  • 避免所有客户端同时重连造成惊群

16. 可观测性设计

16.1 核心指标

连接指标

  • 当前连接数
  • 活跃连接数
  • 每秒新建连接数
  • 每秒关闭连接数
  • 认证前连接数
  • 按 EventLoop 分布的连接数

I/O 指标

  • 每秒读取字节数
  • 每秒写出字节数
  • 每秒完整帧数
  • 半写次数
  • 单轮读写字节分布
  • select() 返回事件数

队列与背压指标

  • 业务队列长度与等待时间
  • 被拒绝任务数
  • 暂停读连接数
  • 单连接待发送字节分布
  • 全局待发送字节
  • 超过高水位次数
  • 超过硬上限关闭次数

延迟指标

  • 接收完整帧到提交业务的时间
  • 业务排队时间
  • 业务执行时间
  • 响应入队到写完的时间
  • 端到端 P50、P95、P99、P99.9

EventLoop 健康指标

  • 每轮循环耗时
  • 任务队列等待时间
  • 超过阈值的长任务次数
  • EventLoop CPU 时间
  • Selector 异常或重建次数

16.2 日志字段

建议结构化输出:

{
  "event": "connection_closed",
  "connectionId": "c-123456",
  "remote": "10.0.0.8:52131",
  "eventLoop": "nio-event-loop-3",
  "reason": "OUTBOUND_HARD_LIMIT",
  "queuedBytes": 9437184,
  "lifetimeMs": 18342
}

高频事件需要采样或聚合,防止日志放大故障。

16.3 线程栈与 JFR

排查 EventLoop 延迟时重点关注:

  • EventLoop 是否执行了数据库、HTTP 或文件阻塞调用
  • 是否出现锁竞争
  • 是否进行大型对象分配
  • 是否在执行压缩、证书校验或复杂反序列化
  • GC 暂停是否与尾延迟同步

JFR 可用于观察 Socket I/O、线程停顿、对象分配、锁和 GC 等事件。

17. 优雅启动、摘流与停机

17.1 正确顺序

17.2 停机阶段

  1. Readiness 置为 false
  2. 等待负载均衡摘除传播
  3. Acceptor 停止接收新连接
  4. 现有连接进入 draining
  5. 拒绝新业务请求或提示客户端重连
  6. 等待在途任务完成
  7. 尝试写完响应
  8. 到达最大停机期限后强制关闭
  9. 关闭业务线程池、EventLoop 和 Selector

17.3 启动就绪

只有在以下条件满足后才返回 Ready:

  • 配置加载完成
  • 业务依赖初始化完成
  • EventLoop 已启动
  • 监听端口绑定成功
  • 必需的协议处理器已注册
  • 监控端点可用

18. 压测方法与容量规划

18.1 不要只报告一个 QPS 数字

“单节点 63 万 QPS”如果没有以下上下文,几乎没有可比性:

  • CPU 型号与核数
  • 内存与网络带宽
  • JDK 版本与 GC
  • 消息大小
  • 长连接还是短连接
  • 是否启用 TLS
  • 请求是否包含真实业务
  • 客户端数量与机器分布
  • 延迟分位数
  • 丢包、错误率和连接失败率
  • 预热时间与测试时长

18.2 测试矩阵

建议至少覆盖:

维度 测试值示例
连接数 1k、10k、50k、100k
消息体 64 B、1 KiB、16 KiB、256 KiB
请求模式 Ping-Pong、单向推送、突发流量
TLS 关闭、开启、会话复用
业务耗时 0 ms、1 ms、10 ms、外部调用
客户端速度 正常、慢读、慢写、断续网络
持续时间 10 分钟、1 小时、24 小时

18.3 三类压测

建连压测

关注:

  • 每秒成功建连数
  • 建连 P99
  • Accept 队列溢出
  • 文件描述符使用量
  • 客户端端口耗尽

稳态吞吐压测

关注:

  • QPS 与带宽
  • CPU 每核利用率
  • P99/P99.9
  • EventLoop 负载均衡
  • GC 与分配速率

过载与恢复压测

关注:

  • 队列是否有界
  • 背压是否生效
  • 是否出现 OOM
  • 拒绝率是否符合预期
  • 流量下降后延迟能否恢复
  • 是否出现长时间排队尾巴

18.4 容量估算

一个粗略连接内存模型:

M_connection =
    connection object
  + selection key / channel state
  + read cumulation
  + average outbound queue
  + protocol/session state
  + TLS state if enabled

实例容量:

M_total ≈ connectionCount × M_connection
        + business queues
        + direct memory pool
        + heap baseline
        + safety margin

生产上应保留安全余量,不能把压测极限值直接当作日常配额。通常还要为突发流量、GC、滚动发布期间连接迁移和节点故障预留容量。

19. 常见错误与修复方式

错误 1:先 register,再 wakeup

channel.register(selector, SelectionKey.OP_READ);
selector.wakeup();

问题:注册动作可能与正在进行的选择操作发生同步等待,且跨线程状态维护分散。

修复:任务先入目标 EventLoop 队列,调用 wakeup(),由 EventLoop 完成注册。

错误 2:把连接读缓冲区直接交给业务线程

问题:I/O 线程继续读取后会覆盖数据,ByteBuffer 状态也会发生竞态。

修复:解码为独立消息,或者使用带引用计数的池化切片。

错误 3:业务线程直接修改 interestOps

问题:虽然部分 API 本身支持并发访问,但业务状态、写队列和事件兴趣集合之间不再具备原子顺序,容易遗漏唤醒或产生竞态。

修复:所有连接状态变更返回所属 EventLoop 执行。

错误 4:始终监听 OP_WRITE

问题:Socket 大多数时间都可写,Selector 持续返回,引发 CPU 空转。

修复:写队列非空时开启,队列清空后立即关闭。

错误 5:一次 write 后丢弃 Buffer

问题:非阻塞写可能只写出部分数据。

修复:保留 Buffer 的 position,后续可写事件继续发送。

错误 6:半包时 clear()

问题:未消费字节被丢弃。

修复:使用 mark/resetcompact() 保留剩余数据。

错误 7:使用超大无界业务队列

问题:吞吐不足时,任务不断排队,最终表现为内存上涨和分钟级延迟。

修复:有界队列、拒绝策略、读暂停和入口限流组合使用。

错误 8:把 Selector 空轮询重建当作固定模板

历史上部分 JDK/内核组合出现过 Selector 异常空轮询问题,但现代 JDK 中不应仅凭 select() 返回 0 就重建 Selector。超时返回、wakeup()、任务提交、取消 Key 和信号中断都可能造成合法的零返回。

修复:

  • 先监控 EventLoop CPU、循环频率、选中 Key 数和 wakeup 次数
  • 确认存在持续高频空转
  • 再实现带时间窗口、最小持续时间和告警的降级重建
  • 优先升级到受支持的 JDK 更新版本

错误 9:认为 DirectBuffer 自动等于零拷贝

问题:DirectBuffer 只是堆外缓冲区,不代表整个数据路径没有复制。

修复:准确区分缓冲区位置、系统调用次数、内核态复制、协议编解码和 TLS 加密路径。

错误 10:把压测极限当生产容量

问题:真实生产还有日志、监控、TLS、业务调用、流量偏斜、发布和故障转移。

修复:以稳定运行区间而不是峰值极限制定容量,并保留故障冗余。

20. 自研 NIO、Netty 与虚拟线程的决策边界

20.1 选择自研 NIO

适合:

  • 协议非常简单且稳定
  • 依赖严格受限
  • 团队具备底层网络、并发和内存治理能力
  • 愿意长期维护 TLS、编解码、背压、监控与安全能力
  • 性能收益经过基准测试证明,而非凭感觉判断

20.2 选择 Netty

适合:

  • 需要快速稳定上线
  • 协议和编解码复杂
  • 需要成熟内存池
  • 需要 TLS、HTTP/2、WebSocket、MQTT 等能力
  • 需要 Linux epoll、BSD/macOS kqueue 等原生传输
  • 需要经过大量生产验证的异常处理和生态组件

截至当前 Netty 4.2 已引入新的 EventLoopGroup 配置方式,升级时应按官方迁移指南处理已弃用的传输专用 EventLoopGroup API,而不是机械复制旧版启动代码。

20.3 选择虚拟线程

适合:

  • 同步编程模型能显著降低复杂度
  • 每个请求包含多个阻塞式 I/O 调用
  • 不需要极致控制每个 Socket 的事件兴趣与缓冲区
  • 使用标准 HTTP/RPC 服务栈

20.4 混合架构

一种常见组合是:

NIO / Netty EventLoop
    ↓ 解码为业务消息
虚拟线程执行业务阻塞调用
    ↓ 生成响应
返回 EventLoop 写出

这种模式保留事件驱动通信层的连接密度和背压能力,同时让业务代码继续使用直观的同步风格。仍需保证:

  • 虚拟线程任务入口有容量控制
  • 下游连接池有上限
  • 业务完成后通过 EventLoop 写回
  • 不在 EventLoop 内等待虚拟线程结果

21. 生产上线检查清单

协议与安全

  • 定义 Magic、版本、消息类型和长度字段
  • 设置最大帧长度
  • 校验负数长度和整数溢出
  • 未认证连接有超时
  • 单 IP 和单租户有限流
  • 慢读、慢写和心跳超时有处理策略
  • 模糊测试覆盖非法报文

线程与并发

  • 每个连接固定归属一个 EventLoop
  • 跨线程注册通过任务队列执行
  • interestOps 只在所属 EventLoop 修改
  • 业务线程不持有可复用读缓冲区
  • 业务线程池或虚拟线程入口有容量限制
  • 单连接顺序语义已明确

读写正确性

  • 正确处理半包、粘包和多帧
  • 正确处理半写和返回 0
  • OP_WRITE 按需启停
  • 单轮读写有预算
  • 连接关闭逻辑幂等
  • 所有异常路径释放全局内存预算

内存与背压

  • 业务队列有界
  • 单连接写队列按字节限制
  • 全局待发送内存有限额
  • 高低水位控制读暂停与恢复
  • DirectBuffer 有容量预算与泄漏监控
  • 十万连接场景完成真实内存测量

可观测性

  • 连接数和建连速率
  • EventLoop 循环延迟
  • 业务排队与执行时间
  • 待读取和待发送字节
  • 背压、拒绝与关闭原因
  • P50/P95/P99/P99.9
  • GC、直接内存和文件描述符

部署与运维

  • Readiness 与 Liveness 语义正确
  • 支持停止 Accept
  • 支持连接排空和最大停机期限
  • 客户端具备随机退避重连
  • 压测包含滚动发布和 Pod 故障
  • 内核参数经过目标环境验证
  • 容量包含节点故障与发布冗余

22. 总结

主从 Reactor 的核心不是画出 MainReactor、SubReactor 和业务线程池三个方框,而是建立严格的状态归属和流量约束:

  1. Acceptor 只负责接入
  2. 每个连接固定归属于一个 EventLoop
  3. 连接状态在 EventLoop 内串行维护
  4. 跨线程动作通过任务队列回到 EventLoop
  5. 读取必须处理 TCP 字节流边界
  6. 写出必须处理半写并按需监听 OP_WRITE
  7. 业务队列、写队列和内存都有硬上限
  8. 背压必须贯穿内核缓冲区、通信层、业务层和下游资源
  9. 所有高性能结论必须由可复现压测和尾延迟指标证明
  10. 自研框架只在明确收益大于长期维护成本时成立

真正稳定的高并发通信系统,不是让 CPU 永远满载,而是在流量正常时保持低延迟,在流量过载时有界退化,在依赖故障时快速隔离,在发布和扩缩容时平稳迁移。希望本文对高并发系统设计的讨论,能为你的技术选型和架构演进提供一些参考。

参考资料

  1. Oracle Java SE 25 API:java.nio.channelsSelectorSocketChannelFileChannel
  2. Oracle Java SE 25 API:StandardSocketOptions
  3. OpenJDK JEP 444:Virtual Threads
  4. OpenJDK JEP 491:Synchronize Virtual Threads without Pinning
  5. Doug Lea:Scalable IO in Java
  6. Netty 官方文档:Native Transports
  7. Netty 4.2 Migration Guide
  8. Linux Kernel Documentation:IP Sysctl

说明:本文代码用于解释生产级设计原则,已修复常见的跨线程 Selector 注册、缓冲区异步复用、半写遗漏和 OP_WRITE 空转问题。正式项目仍需增加单元测试、并发测试、协议模糊测试、TLS、鉴权、指标采集、配置中心和故障注入验证。




上一篇:量子竞赛开启:比特币4700亿美元“审判日”还有多远?
下一篇:全球手机出货量下滑6%,Nothing被迫转型AI优先公司
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-8-3 06:25 , Processed in 0.804138 second(s), 41 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表