引言

在现代计算机体系中,多核处理器已成为标配,充分利用多核能力进行并行处理是提升应用性能的关键。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()、使用 synchronizedLock 等同步机制。但推荐使用 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 限流
Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐