线程池优化网络服务端 | JavaSE

线程池优化网络服务端

一、学习目标

学完本章,你应该能够:

  1. 能够解释“一客户端一线程”模型为什么存在资源扩展问题。
  2. 能够说明线程池优化网络服务端的核心思想。
  3. 能够把一个客户端 Socket 封装成 Runnable 任务。
  4. 能够使用 ExecutorService.execute() 将网络任务提交给线程池。
  5. 能够使用 ThreadPoolExecutor 配置核心线程数、最大线程数、任务队列、线程工厂和拒绝策略。
  6. 能够解释线程池“核心线程 → 队列 → 最大线程 → 拒绝”的任务处理流程。
  7. 能够理解线程池可以限制线程数量,但并不会让阻塞式 Socket IO 自动变成非阻塞 IO。
  8. 能够发现线程池拒绝 Socket 任务时潜在的连接资源泄漏问题。

二、核心知识

2.1 上一章还存在什么问题

上一章实现了:

一个客户端连接
      ↓
创建一个 Thread
      ↓
这个 Thread 持续处理 Socket

代码类似:

while (true) {
    Socket socket = serverSocket.accept();

    new ServerReader(socket).start();
}

这个方案已经可以:

Client A → Thread A
Client B → Thread B
Client C → Thread C

实现多客户端并发通信。

但是它存在一个明显问题:

客户端数量决定线程数量。

假设:

10 个客户端
→ 大约 10 个工作线程

1000 个客户端
→ 大约 1000 个工作线程

10000 个客户端
→ 可能尝试创建 10000 个工作线程

线程并不是免费的。

每个线程都需要:

  • JVM / 操作系统管理资源
  • 线程栈空间
  • CPU 调度
  • 上下文切换

因此不能简单认为:

new Thread(...).start();

可以无限执行。


2.2 线程池解决什么问题

线程池(Thread Pool)的核心思想是:

提前维护有限数量的工作线程,让任务交给这些线程执行,而不是每来一个任务就永久创建一个新线程。

模型:

                     Thread Pool

Socket A ──→ Task A ─┐
Socket B ──→ Task B ─┼──→ 工作线程
Socket C ──→ Task C ─┤
Socket D ──→ Task D ─┘

线程可以被:

任务 A
   ↓
执行完成
   ↓
线程保留
   ↓
继续执行任务 B

从而实现:

线程复用。


2.3 网络服务端中的“任务”到底是什么

上一章:

class ServerReader extends Thread

意味着:

业务逻辑
+
线程本身

绑定在一起。

使用线程池以后,更合理的设计是:

class ServerReaderRunnable implements Runnable

此时这个类只描述:

需要完成什么工作。

例如:

拿到 Socket
    ↓
读取客户端消息
    ↓
处理消息
    ↓
客户端断开
    ↓
关闭 Socket

而:

到底由哪个线程执行?

交给:

ExecutorService

决定。

这体现了一个很重要的设计思想:

任务
与
线程执行机制
解耦

2.4 ExecutorService

ExecutorService 位于:

java.util.concurrent

它表示:

可以接收和管理任务执行的 Executor 服务。

处理普通 Runnable 任务最常用的方法:

void execute(Runnable command)

例如:

pool.execute(
        new ServerReaderRunnable(socket)
);

含义:

Socket
   ↓
封装 Runnable
   ↓
提交线程池
   ↓
由线程池安排线程执行

而不是:

new Thread(task).start();

2.5 ThreadPoolExecutor

本课程已经在多线程章节学习过:

ThreadPoolExecutor

网络编程阶段不再重新完整讲一遍线程池,而是重点理解:

之前学过的线程池到底如何应用到真正的 TCP 服务端。

课程原始示例采用:

ExecutorService pool =
        new ThreadPoolExecutor(
                3,
                10,
                10,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                Executors.defaultThreadFactory(),
                new ThreadPoolExecutor.AbortPolicy()
        );

参数分别表示:

corePoolSize
=
3

核心线程数。

maximumPoolSize
=
10

最大线程数。

keepAliveTime
=
10 秒

非核心空闲线程允许存活的时间。

workQueue
=
ArrayBlockingQueue<>(100)

最多暂存 100 个等待执行的任务。

ThreadFactory
=
defaultThreadFactory()

负责创建线程。

RejectedExecutionHandler
=
AbortPolicy

线程池完全无法接受新任务时:

抛出 RejectedExecutionException

2.6 ThreadPoolExecutor 接收任务的真正顺序

这是本章非常重要的地方。

很多初学者会错误理解为:

核心线程 3 个用完
↓
马上扩容到 10 个
↓
然后才排队

实际上不是。

对于当前:

core = 3
max = 10
queue = 100

任务处理过程可以理解为:

提交任务
   ↓
当前线程数 < 3?
   │
  是
   ↓
创建核心工作线程

核心线程已经达到 3 个后:

继续提交任务
   ↓
尝试进入任务队列

只要队列还没有满:

任务进入队列等待

只有:

核心线程已满
+
任务队列也满

才会尝试:

继续创建线程

一直扩展到:

maximumPoolSize = 10

如果:

线程数已经 10
+
队列也已经满

新任务最终:

触发拒绝策略

因此大致顺序是:

1. 核心线程

2. 工作队列

3. 最大线程扩容

4. 拒绝策略

这个执行顺序必须真正理解。


三、使用方法

3.1 将客户端处理逻辑改造成 Runnable

首先把上一章:

extends Thread

改为:

implements Runnable

完整示例:

import java.io.DataInputStream;
import java.io.EOFException;
import java.io.IOException;
import java.net.Socket;

public class ServerReaderRunnable
        implements Runnable {

    private final Socket socket;

    public ServerReaderRunnable(Socket socket) {
        this.socket = socket;
    }

    @Override
    public void run() {

        String client =
                socket.getInetAddress()
                        .getHostAddress()
                +
                ":"
                +
                socket.getPort();

        try (
                Socket clientSocket = socket;

                DataInputStream input =
                        new DataInputStream(
                                clientSocket
                                        .getInputStream()
                        )
        ) {

            while (true) {

                String message =
                        input.readUTF();

                System.out.println(
                        "[" + client + "] "
                                + message
                );
            }

        } catch (EOFException e) {

            System.out.println(
                    "客户端正常下线:"
                            + client
            );

        } catch (IOException e) {

            System.out.println(
                    "客户端连接异常:"
                            + client
                            + ","
                            + e.getMessage()
            );
        }
    }
}

注意这个类现在:

不是线程

而是:

Runnable 网络处理任务。


3.2 创建线程池版服务端

import java.net.ServerSocket;
import java.net.Socket;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class TCPPoolServer {

    public static void main(String[] args)
            throws Exception {

        System.out.println(
                "=== TCP 线程池服务端启动 ==="
        );

        ExecutorService pool =
                new ThreadPoolExecutor(
                        3,
                        10,
                        10,
                        TimeUnit.SECONDS,
                        new ArrayBlockingQueue<>(100),
                        Executors.defaultThreadFactory(),
                        new ThreadPoolExecutor.AbortPolicy()
                );

        try (
                ServerSocket serverSocket =
                        new ServerSocket(9999)
        ) {

            while (true) {

                Socket socket =
                        serverSocket.accept();

                String client =
                        socket.getInetAddress()
                                .getHostAddress()
                        +
                        ":"
                        +
                        socket.getPort();

                System.out.println(
                        "客户端连接:"
                                + client
                );

                try {

                    pool.execute(
                            new ServerReaderRunnable(
                                    socket
                            )
                    );

                } catch (
                        RejectedExecutionException e
                ) {

                    System.out.println(
                            "服务器繁忙,拒绝客户端:"
                                    + client
                    );

                    socket.close();
                }
            }

        } finally {

            pool.shutdown();
        }
    }
}

核心代码已经从:

new ServerReader(socket).start();

升级为:

pool.execute(
        new ServerReaderRunnable(socket)
);

3.3 主线程现在负责什么

主线程依然只负责:

ServerSocket
    ↓
accept()
    ↓
得到 Socket
    ↓
封装任务
    ↓
提交线程池
    ↓
重新 accept()

对应:

while (true) {

    Socket socket =
            serverSocket.accept();

    pool.execute(
            new ServerReaderRunnable(socket)
    );
}

因此:

主线程不应该自己进入客户端的 readUTF() 循环。


3.4 工作线程负责什么

线程池中的线程负责运行:

ServerReaderRunnable.run()

也就是:

拿到客户端 Socket
        ↓
获取 InputStream
        ↓
循环 readUTF
        ↓
处理消息
        ↓
客户端下线
        ↓
释放 Socket

最终模型:

                  Main Thread

                 ServerSocket
                      │
                    accept
                      │
           ┌──────────┴──────────┐
           ▼                     ▼
        Socket A              Socket B
           │                     │
           ▼                     ▼
        Task A                 Task B
           │                     │
           └────────┬────────────┘
                    ▼
                ThreadPool
              ┌─────┼─────┐
              ▼     ▼     ▼
            Worker Worker Worker

3.5 为什么拒绝任务时还要关闭 Socket

假设线程池已经:

工作线程全部占用
+
队列完全满
+
线程数达到最大值

此时:

pool.execute(task);

AbortPolicy 下会:

throw new RejectedExecutionException();

但注意:

Socket socket =
        serverSocket.accept();

已经执行成功。

也就是说:

TCP 连接已经被服务端接受。

如果我们只是:

catch (RejectedExecutionException e) {
}

却不:

socket.close();

这个 Socket 可能成为:

已经建立
但没人处理
也没有关闭

的无效连接资源。

因此:

catch (RejectedExecutionException e) {
    socket.close();
}

是很重要的资源管理意识。


四、原理与进阶

4.1 线程池真的可以同时处理 10 个长连接吗

先看原课程参数:

core = 3
max = 10
queue = 100

假设每个 Socket 任务都是:

while (true) {
    input.readUTF();
}

也就是一个长期运行、长期阻塞的任务。

前三个客户端:

Client 1
Client 2
Client 3

会分别占据:

3 个核心线程

第 4 个任务到来时:

核心线程已达到 3

线程池不会立即创建第 4 个线程。

它首先:

放入队列

因此:

Client 4

虽然 TCP 连接可能已经被 accept() 接受,

但是它对应的:

ServerReaderRunnable

可能仍然:

排队等待执行

直到:

某个核心工作线程空闲

或者:

队列最终被填满

后线程池才继续扩展线程。


4.2 为什么这个参数对长连接尤其值得注意

对于普通短任务:

任务 A:20 ms
任务 B:50 ms
任务 C:10 ms

排队通常只是:

短暂等待

但 Socket 长连接任务可能:

运行几分钟
几小时
甚至长期不结束

所以:

大队列 + 少量核心线程

可能导致许多已经建立的 Socket:

长时间排在任务队列里,没有线程真正读取它们的数据。

因此原课程:

3 / 10 / queue 100

非常适合帮助我们学习:

ThreadPoolExecutor 参数
+
Socket 任务提交

但不能直接推导:

“这就是生产环境最合理的长连接服务器参数。”

真实参数必须根据:

  • 请求生命周期
  • 是否长期阻塞
  • CPU 核心数
  • 内存
  • 最大连接数
  • 延迟要求
  • 业务吞吐量

进行设计和压测。


4.3 线程池解决的不是“无限并发”

线程池真正做的是:

限制资源
+
复用线程
+
管理任务

不是:

让一台机器无限处理连接

假设:

maximumPoolSize = 10

对于每个都持续阻塞的 Socket 任务:

最多也只有有限数量工作线程
真正同时执行这些任务

更多任务:

排队
或
被拒绝

所以线程池本质上是在建立:

容量边界。


4.4 阻塞式 Socket + 线程池仍然是阻塞 IO

我们现在仍然使用:

input.readUTF();

没有数据:

工作线程阻塞

因此:

线程池

并不会神奇地把:

Blocking IO

变成:

Non-blocking IO

它只是:

使用有限、可管理的线程执行这些阻塞任务。

后续更加高级的服务器体系还可能涉及:

  • Java NIO
  • Selector
  • Netty
  • Virtual Thread
  • 异步 IO

但这些已经超出当前 JavaSE 冻结章节的主线。

本章只要求掌握:

BIO Socket
+
ThreadPool

服务器模型。


4.5 为什么 Runnable 比 extends Thread 更适合线程池

线程池需要的是:

Runnable

任务。

这意味着:

ServerReaderRunnable
=
业务任务

而:

ThreadPoolExecutor
=
负责如何执行这些任务

形成:

What to do
和
How to run
分离

这就是:

任务与执行策略解耦。


五、实践应用

5.1 网络聊天室服务端

前面:

Socket
+
Thread

现在升级:

Socket
+
Runnable
+
ThreadPoolExecutor

以后局域网即时通信系统可以进一步加入:

在线用户集合
+
消息协议
+
线程池
+
消息转发

形成:

                Chat Server

              Thread Pool
                  │
        ┌─────────┼─────────┐
        ▼         ▼         ▼
    Client A   Client B   Client C

5.2 Web 服务器

浏览器连接 Web Server 时,本质上也是:

Browser
   ↓
TCP Connection
   ↓
Server

服务器收到请求后:

封装任务
   ↓
线程池
   ↓
生成 HTTP Response

这正是下一章:

B/S + HTTP

的直接基础。


六、常见问题

6.1 线程池优化以后,一个客户端还对应一个线程吗

对于当前:

阻塞式长连接模型

一个正在执行的客户端任务,通常仍然会占用:

一个工作线程

区别在于:

线程由线程池统一管理

而不是:

每个连接都自行无限 new Thread

6.2 线程池为什么叫“复用线程”

当一个任务执行结束:

Worker Thread

通常不会马上销毁。

它可以继续获取:

下一个 Runnable

执行。


6.3 maximumPoolSize=10 为什么第四个任务不一定创建第四个线程

因为:

核心线程满

以后通常先尝试:

加入 workQueue

只有队列不能继续接受任务时,才会继续创建非核心线程,直到 maximumPoolSize。


6.4 队列越大越好吗

不是。

大队列:

可以缓冲更多任务

但也可能:

增加排队延迟

特别对于:

长期不结束的 Socket 任务

必须更加谨慎。


6.5 maximumPoolSize 越大越好吗

也不是。

线程越多:

  • 内存占用越大
  • 调度成本越高
  • 上下文切换可能越频繁

必须结合实际业务。


6.6 AbortPolicy 做什么

当线程池无法继续接受新任务:

new ThreadPoolExecutor.AbortPolicy()

会:

抛出 RejectedExecutionException

6.7 为什么任务被拒绝后还需要考虑 Socket

因为:

任务被拒绝
≠
之前 accept 得到的 Socket 自动消失

应该根据业务:

关闭连接
或者
返回服务器繁忙响应

而不是泄漏资源。


6.8 为什么不直接用 Executors.newFixedThreadPool

当然可以创建简单线程池。

但是在本课程中:

ThreadPoolExecutor

可以明确展示:

  • 核心线程
  • 最大线程
  • 队列
  • KeepAlive
  • ThreadFactory
  • 拒绝策略

更利于理解服务端资源边界。


七、练习与验收

7.1 知识问答

  1. 一连接一线程模型有什么问题?
  2. 什么叫线程池?
  3. 为什么线程池能够复用线程?
  4. 网络连接应该如何封装成 Runnable
  5. ExecutorService.execute() 有什么作用?
  6. corePoolSizemaximumPoolSize 分别是什么?
  7. workQueue 解决什么问题?
  8. AbortPolicy 做什么?
  9. 为什么 maximumPoolSize=10 不意味着第四个任务一定创建第四个线程?
  10. 为什么长连接任务对任务队列大小更加敏感?
  11. 线程池是否会把阻塞 IO 自动变成非阻塞 IO?
  12. Socket 任务被拒绝时为什么要考虑关闭连接?

7.2 代码阅读

阅读:

ExecutorService pool =
        new ThreadPoolExecutor(
                3,
                10,
                10,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                Executors.defaultThreadFactory(),
                new ThreadPoolExecutor.AbortPolicy()
        );

回答:

  1. 核心线程数量是多少?
  2. 最大线程数量是多少?
  3. 队列最多存多少个任务?
  4. 第 4 个长期运行任务到达时一定会立即创建新线程吗?
  5. 在什么条件下线程池才会从 3 个线程继续扩大?
  6. 在什么情况下执行拒绝策略?

7.3 手写代码

关闭 AI 自动补全。

将上一章:

ServerReader extends Thread

修改为:

ServerReaderRunnable
        implements Runnable

然后实现:

TCPPoolServer

要求:

  • ServerSocket 监听 9999
  • 创建 ThreadPoolExecutor
  • 主线程循环 accept()
  • 每个 Socket 封装为 Runnable
  • 使用 pool.execute()
  • 客户端断开后关闭资源
  • 线程池拒绝任务时关闭 Socket

7.4 Debug

观察:

while (true) {

    Socket socket =
            serverSocket.accept();

    try {

        pool.execute(
                new ServerReaderRunnable(socket)
        );

    } catch (
            RejectedExecutionException e
    ) {

        System.out.println("任务拒绝");
    }
}

请回答:

  1. 线程池拒绝任务以后 socket 当前是什么状态?
  2. 这个程序有什么资源泄漏风险?
  3. 应该增加什么处理?

7.5 综合训练

配置:

core = 2
max = 4
queue = 2

假设每个任务都长期不结束。

依次提交:

Task A
Task B
Task C
Task D
Task E
Task F
Task G

不要运行程序。

先手工推导:

哪些任务创建核心线程?
哪些任务进入队列?
什么时候创建非核心线程?
哪个任务开始可能被拒绝?

然后编写程序验证自己的推理。

7.6 本章验收

如果你能够闭卷解释:

Socket
  ↓
Runnable
  ↓
ExecutorService
  ↓
ThreadPoolExecutor
  ↓
Worker Thread

并能够准确画出:

提交任务
   ↓
核心线程未满?
   ↓ no
任务队列能放?
   ↓ no
最大线程未满?
   ↓ no
拒绝策略

同时可以独立把:

new ServerReader(socket).start();

重构成:

pool.execute(
        new ServerReaderRunnable(socket)
);

就说明已经真正掌握线程池优化网络服务端的核心思想。