适用范围:Java 17/21/25、Linux、TCP 长连接、自定义二进制协议、网关、即时通信、设备接入与 RPC 通信层。
摘要
Java NIO 的价值并不只是“一个线程管理很多连接”。真正决定系统吞吐、尾延迟与稳定性的,是围绕 Selector、SocketChannel、缓冲区、任务调度和反压机制建立的一组并发约束:
- 一个连接的 I/O 状态必须由固定 EventLoop 串行维护
- 跨线程操作必须先进入 EventLoop 任务队列,再由所属线程执行
- TCP 读取必须处理拆包与粘包,写出必须处理半写
OP_WRITE 只能在确有待发送数据时开启
- 业务线程池、连接级积压和全局内存都必须有上限
- 读写缓冲区的生命周期必须清晰,禁止在异步线程之间无约束共享
- 高并发指标必须建立在可复现压测、容量预算和可观测性之上
本文从模型选择开始,完整讲解主从 Reactor 的线程划分、连接归属、协议编解码、异步业务调度、读写反压、内存治理、优雅停机、Kubernetes 部署与压测方法,并给出一套可落地的 Java NIO 参考实现。如果你正在设计或维护一个高并发通信层,希望这里对线程模型、背压机制和内存治理的讨论能给你一些新的视角。也欢迎在云栈社区与更多开发者交流架构设计和实战经验。
目录
- 一次典型的通信层扩容事故
- 先做技术选型:BIO、虚拟线程、NIO 与 Netty
- Java NIO 的核心对象与操作系统映射
- Reactor 模型的三种形态
- 生产级主从 Reactor 总体架构
- 必须遵守的八条并发不变量
- 协议设计:先解决边界,再谈业务
- 核心源码实现
- 读路径:从字节流到业务消息
- 写路径:半写、队列与 OP_WRITE
- 背压与过载保护
- 缓冲区与堆外内存治理
- 线程数与任务调度策略
- 连接生命周期与异常处理
- 网络参数与 Kubernetes 部署
- 可观测性设计
- 优雅启动、摘流与停机
- 压测方法与容量规划
- 常见错误与修复方式
- 自研 NIO、Netty 与虚拟线程的决策边界
- 生产上线检查清单
- 总结
1. 一次典型的通信层扩容事故
某次大促开始后,核心接入网关出现以下现象:
- 新建连接速率突然上升
- 活跃连接数持续增长
- CPU 使用率达到 95%
- P99 延迟由个位数毫秒恶化到数百毫秒
- 线程数、上下文切换和堆内存同步上涨
- 部分实例触发 Full GC,随后被健康检查摘除
旧架构采用“一个连接对应一个平台线程”的阻塞式处理模型。大量连接处于空闲或慢读状态时,线程仍然占用栈空间和调度资源。请求高峰到来后,线程池、连接队列与业务队列相互放大,最终形成如下故障链:
连接突增
↓
平台线程数上升
↓
线程栈与上下文切换成本增加
↓
业务队列等待时间上升
↓
超时重试进一步放大流量
↓
CPU 饱和、P99 恶化、实例失稳
这里真正需要解决的并不是“把 BIO 改成 NIO”这么简单,而是重新设计整个通信层:
- 如何让少量 I/O 线程管理大量连接
- 如何确保慢业务不阻塞 I/O 线程
- 如何解决 TCP 字节流的拆包、粘包与半写
- 如何在业务消费能力不足时限制读入速度
- 如何控制单连接和全局待发送内存
- 如何在容器环境中完成摘流与优雅停机
- 如何用指标证明系统是否真的达到目标
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() 负责:
- 将任务放入线程安全队列
- 当调用方不是 EventLoop 自身时执行
selector.wakeup()
- 由 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 公平性
单轮循环应在以下工作之间保持平衡:
常见策略:
- 限制每轮任务数
- 限制单连接读写预算
- 统计 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 停机阶段
- Readiness 置为 false
- 等待负载均衡摘除传播
- Acceptor 停止接收新连接
- 现有连接进入 draining
- 拒绝新业务请求或提示客户端重连
- 等待在途任务完成
- 尝试写完响应
- 到达最大停机期限后强制关闭
- 关闭业务线程池、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/reset 和 compact() 保留剩余数据。
错误 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 和业务线程池三个方框,而是建立严格的状态归属和流量约束:
- Acceptor 只负责接入
- 每个连接固定归属于一个 EventLoop
- 连接状态在 EventLoop 内串行维护
- 跨线程动作通过任务队列回到 EventLoop
- 读取必须处理 TCP 字节流边界
- 写出必须处理半写并按需监听
OP_WRITE
- 业务队列、写队列和内存都有硬上限
- 背压必须贯穿内核缓冲区、通信层、业务层和下游资源
- 所有高性能结论必须由可复现压测和尾延迟指标证明
- 自研框架只在明确收益大于长期维护成本时成立
真正稳定的高并发通信系统,不是让 CPU 永远满载,而是在流量正常时保持低延迟,在流量过载时有界退化,在依赖故障时快速隔离,在发布和扩缩容时平稳迁移。希望本文对高并发系统设计的讨论,能为你的技术选型和架构演进提供一些参考。
参考资料
- Oracle Java SE 25 API:
java.nio.channels、Selector、SocketChannel、FileChannel
- Oracle Java SE 25 API:
StandardSocketOptions
- OpenJDK JEP 444:Virtual Threads
- OpenJDK JEP 491:Synchronize Virtual Threads without Pinning
- Doug Lea:Scalable IO in Java
- Netty 官方文档:Native Transports
- Netty 4.2 Migration Guide
- Linux Kernel Documentation:IP Sysctl
说明:本文代码用于解释生产级设计原则,已修复常见的跨线程 Selector 注册、缓冲区异步复用、半写遗漏和 OP_WRITE 空转问题。正式项目仍需增加单元测试、并发测试、协议模糊测试、TLS、鉴权、指标采集、配置中心和故障注入验证。