一文吃透高并发编程,包含JDK 21 虚拟线程避坑指南
引言
在现代计算机体系中,多核处理器已成为标配,充分利用多核能力进行并行处理是提升应用性能的关键。Java 从诞生之初就内置了对多线程的支持,经过二十多年的演进,JDK 21 在多线程领域带来了革命性的变化——虚拟线程(Virtual Threads) 的正式引入,使得高并发编程的门槛大幅降低。
本文将系统介绍 Java 高并发编程的核心概念、线程创建、同步机制、并发工具以及 JDK 21 的新特性,并通过实战案例帮助你全面掌握高并发开发。特别地,我们会重点剖析虚拟线程在实际使用中容易踩的 5 个坑,让你少走弯路。
一、线程基础
1. 进程与线程
- 进程:操作系统资源分配的基本单位,拥有独立的内存空间。
- 线程:CPU 调度的基本单位,是进程内的执行单元。同一进程内的线程共享堆内存和方法区,但每个线程拥有独立的栈和程序计数器。
Java 程序至少有一个主线程(main 方法所在的线程)。多线程允许程序同时执行多个任务,提升响应速度和吞吐量。
2. 线程的生命周期
Java 线程在 java.lang.Thread.State 枚举中定义了六种状态:
- NEW:线程已创建但尚未启动。
- RUNNABLE:线程正在 JVM 中运行或等待 CPU 调度(包括就绪和运行中)。
- BLOCKED:线程被阻塞,等待获取监视器锁(
synchronized)。 - WAITING:无限期等待其他线程执行特定操作(如
wait()、join()、park())。 - TIMED_WAITING:限时等待(如
sleep()、wait(timeout)、join(timeout))。 - TERMINATED:线程执行完毕或异常退出。
状态转换图(简化):
NEW -> RUNNABLE -> TERMINATED
RUNNABLE -> BLOCKED -> RUNNABLE
RUNNABLE -> WAITING -> RUNNABLE
RUNNABLE -> TIMED_WAITING -> RUNNABLE
二、创建线程的几种方式
1. 继承 Thread 类
class MyThread extends Thread {
@Override
public void run() {
System.out.println("线程运行中: " + Thread.currentThread().getName());
}
}
public class Main {
public static void main(String[] args) {
MyThread thread = new MyThread();
thread.start(); // 启动线程,调用 run 方法
}
}
2. 实现 Runnable 接口
推荐使用,因为 Java 单继承限制,实现接口更灵活。
class MyRunnable implements Runnable {
@Override
public void run() {
System.out.println("Runnable 线程运行中");
}
}
// 使用
Thread thread = new Thread(new MyRunnable());
thread.start();
可以使用 Lambda 简化:
Thread thread = new Thread(() -> System.out.println("Lambda 线程"));
thread.start();
3. 实现 Callable 接口与 Future
Callable 可以返回结果并抛出异常,配合 Future 获取异步结果。
import java.util.concurrent.*;
public class CallableExample {
public static void main(String[] args) throws Exception {
ExecutorService executor = Executors.newSingleThreadExecutor();
Callable<Integer> task = () -> {
TimeUnit.SECONDS.sleep(1);
return 42;
};
Future<Integer> future = executor.submit(task);
System.out.println("结果: " + future.get()); // 阻塞等待结果
executor.shutdown();
}
}
4. 使用 Executor 框架(推荐)
ExecutorService 提供了线程池管理,避免了手动创建线程的开销。
ExecutorService executor = Executors.newFixedThreadPool(5);
executor.submit(() -> System.out.println("任务执行"));
executor.shutdown();
三、线程同步
多线程环境下,共享资源的并发访问可能导致数据不一致。Java 提供了多种同步机制。
1. synchronized 关键字
synchronized 可以修饰方法或代码块,确保同一时刻只有一个线程执行被保护的代码。
public class Counter {
private int count = 0;
// 同步方法
public synchronized void increment() {
count++;
}
// 同步代码块
public void incrementBlock() {
synchronized (this) {
count++;
}
}
public synchronized int getCount() {
return count;
}
}
synchronized 基于对象监视器(Monitor)实现,是可重入锁。
2. volatile 关键字
volatile 保证变量的可见性和有序性,但不保证原子性。适用于一写多读的场景(如状态标志)。
public class FlagExample {
private volatile boolean running = true;
public void stop() {
running = false; // 写操作立即对其他线程可见
}
public void doWork() {
while (running) {
// 执行任务
}
}
}
3. Lock 接口与 ReentrantLock
java.util.concurrent.locks.Lock 提供了比 synchronized 更灵活的锁控制,支持可中断、超时、公平锁等。
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
public class LockCounter {
private final Lock lock = new ReentrantLock();
private int count = 0;
public void increment() {
lock.lock();
try {
count++;
} finally {
lock.unlock(); // 必须在 finally 中释放锁
}
}
public int getCount() {
lock.lock();
try {
return count;
} finally {
lock.unlock();
}
}
}
4. 原子类(Atomic)
java.util.concurrent.atomic 包提供了一系列原子变量,利用 CAS(Compare-And-Swap)实现无锁线程安全。
import java.util.concurrent.atomic.AtomicInteger;
public class AtomicCounter {
private final AtomicInteger count = new AtomicInteger(0);
public void increment() {
count.incrementAndGet(); // 原子自增
}
public int getCount() {
return count.get();
}
}
原子类性能通常优于锁,适合高并发计数、标志位等场景。
四、线程通信
1. wait / notify / notifyAll
wait() 使当前线程等待,释放锁;notify() 唤醒一个等待线程;notifyAll() 唤醒所有等待线程。必须在同步块中调用。
public class WaitNotifyExample {
private final Object lock = new Object();
private boolean ready = false;
public void waitForReady() throws InterruptedException {
synchronized (lock) {
while (!ready) { // 使用 while 防止虚假唤醒
lock.wait();
}
System.out.println("已就绪");
}
}
public void setReady() {
synchronized (lock) {
ready = true;
lock.notifyAll();
}
}
}
2. Condition
Lock 配合 Condition 实现更灵活的等待/通知,支持多个等待队列。
Lock lock = new ReentrantLock();
Condition condition = lock.newCondition();
// 等待
lock.lock();
try {
while (!conditionMet) {
condition.await();
}
} finally {
lock.unlock();
}
// 通知
lock.lock();
try {
conditionMet = true;
condition.signalAll();
} finally {
lock.unlock();
}
3. BlockingQueue(推荐)
阻塞队列是生产者-消费者模式的最佳选择,内部实现了线程同步。
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ArrayBlockingQueue;
public class ProducerConsumer {
public static void main(String[] args) {
BlockingQueue<Integer> queue = new ArrayBlockingQueue<>(10);
// 生产者
new Thread(() -> {
try {
for (int i = 0; i < 100; i++) {
queue.put(i); // 队列满时阻塞
System.out.println("生产: " + i);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}).start();
// 消费者
new Thread(() -> {
try {
for (int i = 0; i < 100; i++) {
int value = queue.take(); // 队列空时阻塞
System.out.println("消费: " + value);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}).start();
}
}
五、并发工具类
java.util.concurrent 提供了丰富的并发工具,简化多线程编程。
1. CountDownLatch(倒计数锁存器)
允许一个或多个线程等待一组操作完成。
CountDownLatch latch = new CountDownLatch(3);
// 三个工作线程完成后调用 latch.countDown()
// 主线程等待
latch.await(); // 阻塞直到计数为0
2. CyclicBarrier(循环屏障)
让一组线程到达一个屏障点时互相等待,然后一起继续执行,可重复使用。
CyclicBarrier barrier = new CyclicBarrier(4, () -> System.out.println("所有线程到达屏障,开始下一阶段"));
// 每个线程执行到 barrier.await() 时阻塞,直到所有线程都到达
3. Semaphore(信号量)
控制同时访问特定资源的线程数量。
Semaphore semaphore = new Semaphore(3); // 最多3个线程同时访问
semaphore.acquire(); // 获取许可,若无则阻塞
try {
// 访问资源
} finally {
semaphore.release(); // 释放许可
}
4. Exchanger(交换器)
两个线程之间交换数据。
Exchanger<String> exchanger = new Exchanger<>();
String received = exchanger.exchange("线程A的数据");
六、线程池
线程池复用线程,降低创建/销毁开销,并提供任务队列管理。
1. ExecutorService 常用工厂方法
Executors.newFixedThreadPool(int n):固定大小线程池。Executors.newCachedThreadPool():可缓存线程池,空闲线程回收。Executors.newSingleThreadExecutor():单线程池。Executors.newScheduledThreadPool(int corePoolSize):定时调度线程池。
2. ThreadPoolExecutor 核心参数
ThreadPoolExecutor executor = new ThreadPoolExecutor(
corePoolSize, // 核心线程数
maximumPoolSize, // 最大线程数
keepAliveTime, // 非核心线程空闲存活时间
TimeUnit.SECONDS,
workQueue, // 任务队列
threadFactory, // 线程工厂
rejectionHandler // 拒绝策略
);
常见拒绝策略:
- AbortPolicy(默认):抛出
RejectedExecutionException。 - CallerRunsPolicy:由调用线程执行任务。
- DiscardPolicy:静默丢弃。
- DiscardOldestPolicy:丢弃队列中最旧任务。
3. 使用示例
ExecutorService executor = new ThreadPoolExecutor(
2, 4, 60, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100),
Executors.defaultThreadFactory(),
new ThreadPoolExecutor.AbortPolicy()
);
executor.submit(() -> System.out.println("任务执行"));
executor.shutdown();
七、Fork/Join 框架
Fork/Join 用于并行处理可分解的任务(如大数组求和、归并排序等),采用工作窃取算法。
import java.util.concurrent.RecursiveTask;
import java.util.concurrent.ForkJoinPool;
public class SumTask extends RecursiveTask<Long> {
private static final int THRESHOLD = 1000;
private final long[] array;
private final int start, end;
public SumTask(long[] array, int start, int end) {
this.array = array;
this.start = start;
this.end = end;
}
@Override
protected Long compute() {
if (end - start <= THRESHOLD) {
long sum = 0;
for (int i = start; i < end; i++) {
sum += array[i];
}
return sum;
} else {
int mid = start + (end - start) / 2;
SumTask left = new SumTask(array, start, mid);
SumTask right = new SumTask(array, mid, end);
left.fork(); // 异步执行左子任务
long rightResult = right.compute(); // 同步执行右子任务
long leftResult = left.join(); // 获取左子任务结果
return leftResult + rightResult;
}
}
public static void main(String[] args) {
long[] array = new long[10000];
// 初始化数组...
ForkJoinPool pool = new ForkJoinPool();
Long result = pool.invoke(new SumTask(array, 0, array.length));
System.out.println("总和: " + result);
}
}
八、CompletableFuture 异步编程
CompletableFuture 提供了强大的异步编程能力,支持链式调用、组合、异常处理。
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
public class CompletableFutureExample {
public static void main(String[] args) throws ExecutionException, InterruptedException {
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> "Hello")
.thenApply(s -> s + " World")
.thenApply(String::toUpperCase);
System.out.println(future.get()); // HELLO WORLD
// 组合两个异步任务
CompletableFuture<String> future1 = CompletableFuture.supplyAsync(() -> "Hello");
CompletableFuture<String> future2 = CompletableFuture.supplyAsync(() -> "World");
CompletableFuture<String> combined = future1.thenCombine(future2, (a, b) -> a + " " + b);
System.out.println(combined.get()); // Hello World
// 异常处理
CompletableFuture<String> withError = CompletableFuture.supplyAsync(() -> {
throw new RuntimeException("错误");
}).exceptionally(ex -> "默认值");
System.out.println(withError.get()); // 默认值
}
}
九、JDK 21 新特性:虚拟线程
1. 虚拟线程简介
虚拟线程(Virtual Threads)是 JDK 21 正式引入的轻量级线程实现(JEP 444)。它们由 JVM 管理,而非操作系统线程,创建成本极低(一个虚拟线程约几 KB 内存),可以轻松创建数十万甚至上百万个虚拟线程。虚拟线程非常适合 I/O 密集型、高并发的场景,大幅简化了传统的"每个任务一个线程"模型。
传统平台线程(Platform Thread)与虚拟线程对比:
| 特性 | 平台线程 | 虚拟线程 |
|---|---|---|
| 实现 | 操作系统线程的包装 | JVM 内部调度 |
| 创建成本 | 高(约 1MB 栈) | 极低(可动态调整栈) |
| 数量限制 | 数百到数千 | 可达百万级 |
| 适用场景 | CPU 密集型 | I/O 密集型、高并发 |
2. 创建虚拟线程
// 方式一:Thread.ofVirtual()
Thread virtualThread = Thread.ofVirtual()
.name("virtual-1")
.start(() -> System.out.println("虚拟线程运行"));
// 方式二:Thread.startVirtualThread()
Thread vThread = Thread.startVirtualThread(() -> {
System.out.println("虚拟线程执行");
});
// 方式三:ExecutorService 创建虚拟线程池
ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
executor.submit(() -> System.out.println("任务在虚拟线程中执行"));
executor.shutdown();
虚拟线程使用方式与普通线程相同,可以调用 sleep()、join()、使用 synchronized、Lock 等同步机制。但推荐使用 ReentrantLock 而非 synchronized,因为虚拟线程在 synchronized 阻塞时可能占用平台线程(不过在 JDK 21 中已有优化,但仍需注意)。
3. 实战案例:高并发网络请求
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.util.concurrent.Executors;
import java.util.concurrent.ExecutorService;
public class VirtualThreadHttpDemo {
public static void main(String[] args) {
ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor();
HttpClient client = HttpClient.newHttpClient();
for (int i = 0; i < 1000; i++) {
int id = i;
executor.submit(() -> {
try {
HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create("https://example.com/api/" + id))
.build();
HttpResponse<String> response = client.send(request, HttpResponse.BodyHandlers.ofString());
System.out.println("请求 " + id + " 完成: " + response.statusCode());
} catch (Exception e) {
e.printStackTrace();
}
});
}
executor.shutdown();
}
}
这段代码可以轻松创建 1000 个虚拟线程并发执行 HTTP 请求,而传统线程池可能需要限制在几十个线程。
4. ⚠️ 虚拟线程的坑:必须注意的 5 个问题
虚拟线程虽好,但并非万能。以下是实际开发中最容易踩的坑,务必仔细阅读。
坑 1:synchronized 会钉住平台线程(Pinning)
问题:当虚拟线程在 synchronized 块内执行阻塞操作(如 Thread.sleep()、I/O)时,会钉住(pin) 底层平台线程,导致该平台线程无法被其他虚拟线程复用,从而降低并发性能。
// ❌ 错误示范:synchronized 内做阻塞 I/O,会钉住平台线程
public class PinningDemo {
private static final Object LOCK = new Object();
public static void main(String[] args) throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 100; i++) {
int id = i;
executor.submit(() -> {
synchronized (LOCK) { // 进入同步块
try {
Thread.sleep(100); // 阻塞操作 → 钉住平台线程
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
System.out.println("任务 " + id + " 完成");
}
});
}
}
}
}
解决方案:优先使用 ReentrantLock 替代 synchronized,因为 ReentrantLock 在阻塞时不会钉住平台线程。
// ✅ 正确示范:使用 ReentrantLock,避免钉住
public class PinningFixDemo {
private static final Lock LOCK = new ReentrantLock();
public static void main(String[] args) throws Exception {
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 100; i++) {
int id = i;
executor.submit(() -> {
LOCK.lock();
try {
Thread.sleep(100); // 阻塞时不会钉住平台线程
System.out.println("任务 " + id + " 完成");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
LOCK.unlock();
}
});
}
}
}
}
注意:JDK 24 中 JEP 491 已修复
synchronized钉住问题,但在 JDK 21 中仍需规避。
坑 2:线程池混用导致虚拟线程失效
问题:将虚拟线程提交到固定大小的平台线程池中,虚拟线程的优势将完全丧失——所有任务仍被限制在有限的平台线程上执行。
// ❌ 错误示范:虚拟线程被提交到固定线程池
ExecutorService pool = Executors.newFixedThreadPool(10);
for (int i = 0; i < 1000; i++) {
pool.submit(() -> {
// 这里虽然用了虚拟线程 API,但实际跑在 10 个平台线程上
});
}
解决方案:虚拟线程必须使用 Executors.newVirtualThreadPerTaskExecutor(),或直接使用 Thread.ofVirtual().start()。
// ✅ 正确示范:使用虚拟线程专属执行器
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1000; i++) {
executor.submit(() -> {
// 每个任务都运行在独立的虚拟线程上
});
}
}
坑 3:CPU 密集型任务误用虚拟线程
问题:虚拟线程适合 I/O 密集型任务,但不适合 CPU 密集型任务。CPU 密集型任务(如复杂计算、加密解密)会持续占用平台线程,虚拟线程无法带来性能提升,反而因调度开销导致性能下降。
// ❌ 错误示范:CPU 密集型任务使用虚拟线程
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1000; i++) {
executor.submit(() -> {
// 大量 CPU 计算,虚拟线程无优势
double result = 0;
for (int j = 0; j < 1_000_000; j++) {
result += Math.sqrt(j);
}
});
}
}
解决方案:CPU 密集型任务应使用固定大小的平台线程池,线程数通常设为 CPU 核心数 + 1。
// ✅ 正确示范:CPU 密集型任务使用平台线程池
int cores = Runtime.getRuntime().availableProcessors();
ExecutorService pool = Executors.newFixedThreadPool(cores + 1);
坑 4:ThreadLocal 内存泄漏
问题:虚拟线程支持 ThreadLocal,但由于虚拟线程数量巨大,如果每个虚拟线程都设置 ThreadLocal 且不清理,会导致严重的内存泄漏。
// ❌ 错误示范:大量虚拟线程设置 ThreadLocal 不清理
ThreadLocal<byte[]> local = new ThreadLocal<>();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1_000_000; i++) {
executor.submit(() -> {
local.set(new byte[1024 * 1024]); // 每个虚拟线程持有 1MB
// 任务结束,但 ThreadLocal 未清理
});
}
}
解决方案:使用 try-finally 清理 ThreadLocal,或改用 ScopedValue(JDK 21 预览特性)。
// ✅ 正确示范:使用 try-finally 清理
ThreadLocal<byte[]> local = new ThreadLocal<>();
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1_000_000; i++) {
executor.submit(() -> {
try {
local.set(new byte[1024 * 1024]);
// 业务逻辑
} finally {
local.remove(); // 必须清理
}
});
}
}
坑 5:盲目替换所有线程池
问题:虚拟线程并非银弹,盲目将现有线程池全部替换为虚拟线程可能导致意想不到的问题,如线程池的限流作用消失、资源耗尽等。
// ❌ 错误示范:无脑替换
// 原来:限制并发 10 个
ExecutorService pool = Executors.newFixedThreadPool(10);
// 现在:无限制创建虚拟线程,可能打爆下游服务
ExecutorService pool = Executors.newVirtualThreadPerTaskExecutor();
解决方案:虚拟线程适合无共享状态、无资源限制的 I/O 密集型任务。对于需要限流的场景,应使用 Semaphore 控制并发数。
// ✅ 正确示范:虚拟线程 + Semaphore 限流
Semaphore semaphore = new Semaphore(10); // 最多 10 个并发
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
for (int i = 0; i < 1000; i++) {
executor.submit(() -> {
try {
semaphore.acquire();
// 调用下游服务
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
semaphore.release();
}
});
}
}
小结
| 坑 | 核心问题 | 解决方案 |
|---|---|---|
| 坑 1:synchronized 钉住 | 阻塞时占用平台线程 | 改用 ReentrantLock |
| 坑 2:线程池混用 | 虚拟线程优势丧失 | 使用专属执行器 |
| 坑 3:CPU 密集型误用 | 性能不升反降 | 使用平台线程池 |
| 坑 4:ThreadLocal 泄漏 | 内存被大量占用 | 及时清理或改用 ScopedValue |
| 坑 5:盲目替换 | 限流失效、资源耗尽 | 结合 Semaphore 限流 |
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐


所有评论(0)