目录
正在加载目录…
专栏文章
专栏文章
Java 专栏
1. Java:跨平台语言与生态全景 2. Java 并发:JMM、锁与线程池原理 3. JVM:类加载、内存、GC 与性能诊断 4. Java 性能调优:CPU、内存、锁与 IO 排查 5. Java 新特性速查:从 Lambda 到虚拟线程 6. Java 业务开发:高频陷阱与 Code Review 清单 7. Java 核心知识:集合、并发与 JVM 面试要点 8. AQS 原理:同步队列与加锁解锁全流程 9. JVM 与 Linux 内存:分区、分配与 GC 边界 10. AtomicLong 与 LongAdder:并发计数器选型 11. Java 引用与内存泄漏:六类场景与修复 12. Java 测试实践:JUnit、Mockito 与 Testcontainers 13. Java 并发入门:线程、锁与线程池全景

Java 并发入门:线程、锁与线程池全景

发布于 2026-06-24 04:25 · 最后编辑于 2026-07-31 15:51 · 字数 2,938 👁 185 次阅读

体系化梳理 Java 并发编程的核心工具与适用场景:从线程基础、三大并发问题,到 synchronized / Lock / 原子类 / 线程池 / 并发工具类,每个工具都配有典型场景与 JDK 8 兼容代码示例,帮助快速建立并发编程全局视图。

目录

章节说明
为什么需要并发并发的价值与代价
线程基础创建线程、状态机、基本操作
三大并发问题可见性、原子性、有序性
同步工具synchronized、volatile、Lock、ReadWriteLock
并发工具类CountDownLatch、CyclicBarrier、Semaphore
线程安全集合ConcurrentHashMap、阻塞队列、CopyOnWriteArrayList
原子类AtomicXxx、LongAdder、ABA 问题
线程池ThreadPoolExecutor 核心参数与实战
快速选型指南场景 → 工具的决策树

为什么需要并发

并发的价值

  • 提升吞吐量:IO 密集型任务(读数据库、调接口)等待期间,让 CPU 处理其他请求
  • 充分利用多核:CPU 密集型任务可并行分治,用多核加速
  • 提升响应速度:后台异步任务不阻塞主流程,用户感知更快

并发的代价:引入了可见性原子性有序性三类线程安全问题(详见下文)。

选择并发的判断标准:任务之间有没有依赖关系?相互独立 → 可并发;有依赖 → 需要同步协调。

线程基础

创建线程的三种方式

import java.util.concurrent.*;

// 方式一:继承 Thread(不推荐,占用 Java 的单继承槽位)
class PrintThread extends Thread {
    @Override
    public void run() {
        System.out.println("Thread: " + Thread.currentThread().getName());
    }
}
new PrintThread().start();

// 方式二:实现 Runnable(推荐,解耦任务逻辑与线程管理)
// 匿名内部类写法(JDK 7-)
Thread t1 = new Thread(new Runnable() {
    @Override
    public void run() {
        System.out.println("Runnable task");
    }
}, "my-thread");
// Lambda 写法(JDK 8+)
Thread t2 = new Thread(() -> System.out.println("Lambda task"), "lambda-thread");
t2.start();

// 方式三:Callable + Future(有返回值,异常可向调用方传递)
ExecutorService executor = Executors.newFixedThreadPool(4);
// Lambda 写法(编译器推断为 Callable<Integer>,允许 checked exception)
Future<Integer> future = executor.submit(() -> {
    Thread.sleep(100);  // 模拟耗时计算
    return 42;
});
try {
    System.out.println(future.get());  // 阻塞等待结果
} catch (ExecutionException e) {
    System.err.println("任务失败:" + e.getCause().getMessage());
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
}
executor.shutdown();

线程状态机

flowchart LR
    NEW -->|"start()"| RUNNABLE
    RUNNABLE -->|"竞争 synchronized 锁失败"| BLOCKED
    BLOCKED -->|"获得锁"| RUNNABLE
    RUNNABLE -->|"Object.wait()"| WAITING
    WAITING -->|"notify() / notifyAll()"| RUNNABLE
    RUNNABLE -->|"sleep(n) / join(n)"| TIMED_WAITING
    TIMED_WAITING -->|"超时 / 唤醒"| RUNNABLE
    RUNNABLE -->|"run() 结束"| TERMINATED
    style NEW fill:#e8f4f8,stroke:#4a90d9
    style TERMINATED fill:#f8e8e8,stroke:#c00
    style RUNNABLE fill:#e8f8e8,stroke:#060
状态含义进入原因
NEW已创建,未启动new Thread() 之后,start() 之前
RUNNABLE运行中或就绪(等 CPU 调度)start() 之后
BLOCKED等待 synchronized 锁竞争 synchronized 代码块
WAITING无限期等待Object.wait()Thread.join()
TIMED_WAITING有限期等待Thread.sleep(n)Object.wait(n)
TERMINATED已结束run() 正常返回或抛出异常

常用线程操作

Thread t = new Thread(() -> { /* 任务逻辑 */ });
t.start();               // 启动线程,JVM 调度后执行 run()
t.join();                // 当前线程阻塞,等待 t 执行完毕
t.join(1000);            // 最多等 1 秒,不保证 t 已结束
t.interrupt();           // 发送中断信号(只设标志位,不强制停止线程)
t.isInterrupted();       // 检查中断标志(不清除)
Thread.interrupted();    // 检查并清除当前线程的中断标志
Thread.sleep(100);       // 当前线程休眠 100ms,释放 CPU 但不释放锁
Thread.currentThread();  // 获取当前线程对象

⚠️ 不要使用 Thread.stop()(已废弃):强制停止可能导致对象处于不一致状态。正确做法是用 interrupt() + 主动检查标志位来协作退出。

三大并发问题

可见性:CPU 缓存的代价

多核 CPU 每个核心有独立的 L1/L2 缓存,线程 A 修改变量后写入缓存,线程 B 从另一个核心的缓存读取时可能看到旧值。

// ❌ 可能死循环:reader 线程可能永远看不到 running = false
class StopThreadBad {
    private static boolean running = true;

    public static void main(String[] args) throws Exception {
        new Thread(new Runnable() {
            @Override
            public void run() {
                while (running) { /* 持续循环 */ }
                System.out.println("线程退出");  // 可能永远不会执行
            }
        }).start();

        Thread.sleep(1000);
        running = false;  // 对 reader 线程不可见!
    }
}

// ✅ 修复:加 volatile,强制读写直接访问主内存
class StopThreadOk {
    private static volatile boolean running = true;
    // ... 其他代码不变
}

原子性:线程切换的代价

count++ 在 JVM 层面是三条指令:读内存 → 加 1 → 写内存。OS 可以在任意两条指令之间切换线程,导致两个线程基于同一旧值计算,相互覆盖写入结果。

// ❌ 两个线程各执行 10000 次,最终结果不是 20000
class UnsafeCounter {
    long count = 0;
    void add() { count++; }  // 读-改-写三步,非原子
}

// ✅ 方案一:synchronized 加锁(保证原子性 + 可见性)
class LockedCounter {
    long count = 0;
    synchronized void add() { count++; }
    synchronized long get() { return count; }
}

// ✅ 方案二:AtomicLong(无锁,性能更好)
class AtomicCounter {
    AtomicLong count = new AtomicLong(0);
    void add() { count.incrementAndGet(); }
    long get() { return count.get(); }
}

有序性:编译优化的代价

编译器和 CPU 为提升性能会对指令重排序。最经典的问题是双重检查单例(DCL):

// ❌ 问题:instance 未加 volatile,new Singleton() 的步骤可能被重排序
// 正常步骤:① 分配内存 → ② 初始化对象 → ③ 赋值给 instance
// 重排后可能:① 分配内存 → ② 赋值给 instance(非空!)→ ③ 初始化对象
// 线程 B 看到 instance 非空就直接返回,但对象还没初始化完!
class SingletonBad {
    static SingletonBad instance;

    static SingletonBad getInstance() {
        if (instance == null) {
            synchronized (SingletonBad.class) {
                if (instance == null)
                    instance = new SingletonBad();  // 危险!可能重排序
            }
        }
        return instance;
    }
}

// ✅ 修复:volatile 禁止 new 操作的指令重排序
class SingletonOk {
    static volatile SingletonOk instance;

    static SingletonOk getInstance() {
        if (instance == null) {
            synchronized (SingletonOk.class) {
                if (instance == null)
                    instance = new SingletonOk();  // volatile 保证赋值在初始化之后
            }
        }
        return instance;
    }
}

同步工具

synchronized — 最简单的互斥锁

适用场景:保护临界区(需要原子执行的代码块),竞争不激烈时首选,代码简洁,自动释放锁。

class BankAccount {
    private long balance;

    // 场景 1:方法级锁,锁对象 = this(实例方法)
    public synchronized void deposit(long amount) {
        balance += amount;
    }

    // 场景 2:代码块锁,锁粒度更小,减少持锁时间
    public void withdraw(long amount) {
        if (amount <= 0) throw new IllegalArgumentException("金额非法");
        synchronized (this) {  // 只锁核心操作
            if (balance < amount) throw new IllegalStateException("余额不足");
            balance -= amount;
        }
    }

    // 场景 3:静态方法锁,锁对象 = BankAccount.class
    public static synchronized BankAccount createDefault() {
        return new BankAccount();
    }

    public synchronized long getBalance() { return balance; }
}

⚠️ 常见错误

// ❌ 每次 new 一个新锁对象,完全无保护效果
synchronized (new Object()) { count++; }

// ❌ add() 和 sub() 用了不同的锁,两个方法可以同时执行
class Wrong {
    private final Object lock = new Object();
    long count = 0;

    synchronized void add() { count++; }          // 锁 this
    void sub() { synchronized (lock) { count--; } }  // 锁 lock(不同对象!)
}

volatile — 轻量级可见性与有序性保证

适用场景:状态标志位、单次写-多次读的场景。不保证原子性,不能替代 synchronized。

// 场景 1:停止后台线程
class BackgroundWorker {
    private volatile boolean shutdown = false;

    public void run() {
        while (!shutdown) {
            doWork();
        }
        System.out.println("Worker 已停止");
    }

    public void stop() {
        shutdown = true;  // 立即对其他线程可见
    }
}

// 场景 2:DCL 单例(见三大并发问题示例)

// ❌ volatile 不保证原子性!
class NotSafe {
    private volatile long count = 0;
    void add() { count++; }  // count++ 仍是三步操作,线程不安全
}

Lock — 灵活的显式锁

适用场景:需要以下特性时替代 synchronized:

  • 超时获锁tryLock(timeout)):避免死等,可降级处理
  • 可中断获锁lockInterruptibly()):支持任务取消
  • 多个条件变量Condition):精确唤醒某类等待线程
import java.util.concurrent.locks.*;
import java.util.LinkedList;
import java.util.Queue;

// 场景:有界阻塞缓冲区,生产者/消费者精确唤醒
class BoundedBuffer<T> {
    private final Lock lock = new ReentrantLock();
    private final Condition notFull  = lock.newCondition();  // "未满"条件
    private final Condition notEmpty = lock.newCondition();  // "非空"条件
    private final Queue<T> queue = new LinkedList<T>();
    private final int capacity;

    public BoundedBuffer(int capacity) { this.capacity = capacity; }

    // 生产者:队列满时阻塞,等待"未满"信号
    public void put(T item) throws InterruptedException {
        lock.lock();
        try {
            while (queue.size() == capacity) {
                notFull.await();   // 释放锁并等待
            }
            queue.add(item);
            notEmpty.signal();    // 通知一个消费者
        } finally {
            lock.unlock();        // 必须在 finally 中释放,防止异常时死锁
        }
    }

    // 消费者:队列空时阻塞,等待"非空"信号
    public T take() throws InterruptedException {
        lock.lock();
        try {
            while (queue.isEmpty()) {
                notEmpty.await();
            }
            T item = queue.poll();
            notFull.signal();
            return item;
        } finally {
            lock.unlock();
        }
    }
}

// 场景:超时获锁,避免死等
class TimeoutExample {
    private final ReentrantLock lock = new ReentrantLock();

    void tryUpdate() throws InterruptedException {
        if (lock.tryLock(500, java.util.concurrent.TimeUnit.MILLISECONDS)) {
            try {
                // 成功获锁,执行操作
                doUpdate();
            } finally {
                lock.unlock();
            }
        } else {
            // 超时,执行降级逻辑
            System.out.println("获锁超时,跳过本次更新");
        }
    }

    private void doUpdate() { /* ... */ }
}

ReadWriteLock — 读多写少场景

读操作并发(多个读者同时持锁),写操作独占(写者持锁时读者阻塞)。

import java.util.concurrent.locks.*;
import java.util.HashMap;
import java.util.Map;

// 场景:本地缓存(读远多于写)
class LocalCache<K, V> {
    private final ReadWriteLock rwl = new ReentrantReadWriteLock();
    private final Map<K, V> cache = new HashMap<K, V>();

    // 多个读线程可并发执行
    public V get(K key) {
        rwl.readLock().lock();
        try {
            return cache.get(key);
        } finally {
            rwl.readLock().unlock();
        }
    }

    // 写操作独占,读写互斥
    public void put(K key, V value) {
        rwl.writeLock().lock();
        try {
            cache.put(key, value);
        } finally {
            rwl.writeLock().unlock();
        }
    }
}

三种工具对比

工具原子性可见性有序性核心特点
synchronized简单,自动释放锁
volatile✅(禁重排)轻量,只保证可见性
ReentrantLock灵活,需手动释放
ReentrantReadWriteLock读并发,写独占

并发工具类

CountDownLatch — 等待多个任务完成

场景:主线程等待 N 个子任务全部完成后再继续(一次性,不可重置)。

import java.util.concurrent.*;

// 场景:并行初始化多个服务,全部就绪后才开放流量
final int serviceCount = 3;
final CountDownLatch latch = new CountDownLatch(serviceCount);

final String[] services = {"DB", "Cache", "MQ"};
for (final String name : services) {
    new Thread(new Runnable() {
        @Override
        public void run() {
            try {
                initService(name);  // 模拟各服务初始化(耗时不同)
                System.out.println(name + " 初始化完成");
            } finally {
                latch.countDown();  // 无论成功失败,计数减 1
            }
        }
    }).start();
}

latch.await();  // 主线程阻塞,直到计数归零
System.out.println("所有服务就绪,开始接收请求");

// 也可以指定最大等待时间
boolean done = latch.await(30, TimeUnit.SECONDS);
if (!done) System.err.println("服务启动超时!");

CyclicBarrier — 多线程同步到达集合点

场景:N 个线程都到达某个同步点后,再一起继续(可重复使用,适合多阶段并行任务)。

import java.util.concurrent.*;

// 场景:分阶段并行计算,每阶段所有线程完成后汇总再进入下一阶段
final int workerCount = 3;
final CyclicBarrier barrier = new CyclicBarrier(workerCount, new Runnable() {
    @Override
    public void run() {
        // 当所有线程都到达 barrier 时,由最后到达的线程执行此回调
        System.out.println("=== 阶段完成,开始下一阶段 ===");
    }
});

for (int i = 0; i < workerCount; i++) {
    final int workerId = i;
    new Thread(new Runnable() {
        @Override
        public void run() {
            try {
                computePhase1(workerId);  // 并行执行第一阶段
                barrier.await();         // 等待所有线程完成第一阶段
                computePhase2(workerId); // 并行执行第二阶段
                barrier.await();         // 等待所有线程完成第二阶段(CyclicBarrier 可重用)
            } catch (InterruptedException | BrokenBarrierException e) {
                Thread.currentThread().interrupt();
            }
        }
    }).start();
}

Semaphore — 限制并发数量

场景:流量控制(数据库连接池、对外接口限流)。

import java.util.concurrent.*;

// 场景:最多允许 5 个线程同时查询数据库
final Semaphore semaphore = new Semaphore(5);

void queryDatabase(String sql) {
    try {
        semaphore.acquire();  // 获取 1 个许可(无许可时阻塞)
        try {
            executeQuery(sql);
        } finally {
            semaphore.release();  // 必须在 finally 中释放!
        }
    } catch (InterruptedException e) {
        Thread.currentThread().interrupt();
    }
}

// 也可以非阻塞尝试
if (semaphore.tryAcquire()) {
    try {
        executeQuery(sql);
    } finally {
        semaphore.release();
    }
} else {
    // 无许可可用,快速失败
    throw new RuntimeException("服务繁忙,请稍后重试");
}
工具核心用途是否可重置
CountDownLatch等待 N 个事件发生否(一次性)
CyclicBarrierN 个线程同步到集合点
Semaphore控制最大并发数量

线程安全集合

ConcurrentHashMap

场景:多线程共享 Map,替代线程不安全的 HashMap(易数据丢失/死循环)和性能差的 Hashtable(全表锁)。

import java.util.concurrent.*;
import java.util.function.*;

ConcurrentHashMap<String, Integer> map = new ConcurrentHashMap<String, Integer>();

// 普通读写
map.put("key", 1);
Integer val = map.get("key");

// 原子性复合操作(不能用 get + put 代替,那不是原子的!)
map.putIfAbsent("key", 1);   // 不存在才插入,返回原来的值(或 null)

// Java 8:compute 系列(lambda 在持锁期间执行,保证原子性)
map.computeIfAbsent("newKey", new Function<String, Integer>() {
    @Override
    public Integer apply(String k) {
        return k.length();   // 不存在时计算并插入
    }
});
// Lambda 简写
map.computeIfAbsent("anotherKey", k -> k.length());

// Java 8:merge(存在则合并,适合计数器场景)
map.merge("counter", 1, new BiFunction<Integer, Integer, Integer>() {
    @Override
    public Integer apply(Integer oldVal, Integer newVal) {
        return oldVal + newVal;
    }
});
// Lambda 简写
map.merge("counter", 1, Integer::sum);  // 不存在则插入 1,存在则累加

⚠️ key 和 value 都不能为 null;并发修改时 size() 是近似值。

阻塞队列(生产者-消费者)

场景:解耦生产速度与消费速度,队列满/空时自动阻塞,无需手写 wait/notify。

import java.util.concurrent.*;

// 推荐使用有界队列,防止生产过快导致 OOM
final BlockingQueue<String> queue = new ArrayBlockingQueue<String>(100);

// 生产者线程
Thread producer = new Thread(new Runnable() {
    @Override
    public void run() {
        try {
            while (true) {
                String task = buildTask();
                queue.put(task);  // 队列满时阻塞,直到有空间
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
});

// 消费者线程
Thread consumer = new Thread(new Runnable() {
    @Override
    public void run() {
        try {
            while (true) {
                String task = queue.take();  // 队列空时阻塞,直到有任务
                process(task);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
});
实现类是否有界推荐度特点
ArrayBlockingQueue⭐⭐⭐推荐,防 OOM,数组实现
LinkedBlockingQueue可选(默认无界)⭐⭐默认无界时注意 OOM
SynchronousQueue无缓冲⭐⭐生产者必须等消费者接手才返回
DelayQueue无界⭐⭐按延时时间出队(定时任务)
PriorityBlockingQueue无界⭐⭐按优先级出队(元素需实现 Comparable)

CopyOnWriteArrayList

场景:读多写极少的 List(如监听器列表、白名单配置),读操作完全无锁。

import java.util.concurrent.*;

CopyOnWriteArrayList<String> whitelist = new CopyOnWriteArrayList<String>();

// 写操作:底层数组整体复制,然后替换引用(代价高)
whitelist.add("192.168.1.1");

// 读 / 迭代:完全无锁,读的是调用时刻的快照,并发写入不影响当前迭代
for (String ip : whitelist) {
    if (isValidIp(ip)) { /* 处理 */ }
}
特性说明
适用场景写操作极少(每次写要复制整个数组)
读取语义快照语义,迭代时不会看到并发新增的元素
迭代器只读,不支持 iterator.remove()

原子类

场景:单个变量的无锁原子操作(CAS 实现)。比 synchronized 性能更好,适合高并发计数、ID 生成、状态标志原子更新。

import java.util.concurrent.atomic.*;

// ===== 基本类型原子类 =====
AtomicInteger count = new AtomicInteger(0);
count.incrementAndGet();              // 原子 ++i,返回新值
count.getAndIncrement();              // 原子 i++,返回旧值
count.addAndGet(5);                   // 原子 += 5,返回新值
count.compareAndSet(10, 0);           // CAS:当前值 == 10 才改为 0,返回是否成功

// ===== 高并发累加器(Java 8+,比 AtomicLong 竞争更小)=====
LongAdder adder = new LongAdder();
adder.increment();          // 分桶累加,各线程操作不同桶,减少 CAS 竞争
adder.add(100);
long total = adder.sum();   // 汇总所有桶的值(注意:有并发时 sum 不精确)

// ===== 对象引用原子替换(配置热更新场景)=====
AtomicReference<String> configRef = new AtomicReference<String>("v1");
String oldConfig = configRef.get();
String newConfig = "v2";
boolean replaced = configRef.compareAndSet(oldConfig, newConfig);  // CAS 原子替换

// ===== 解决 ABA 问题:带版本号的引用 =====
// ABA 问题:值从 A 改为 B 再改回 A,普通 CAS 无法感知
AtomicStampedReference<Integer> stampedRef =
        new AtomicStampedReference<Integer>(100, 0);  // (value=100, stamp=0)

int[] stampHolder = new int[1];
int currentVal = stampedRef.get(stampHolder);  // 同时读取 value 和 stamp
int currentStamp = stampHolder[0];

// value 和 stamp 都匹配才能更新(stamp 变化说明中间被修改过)
stampedRef.compareAndSet(currentVal, currentVal + 1, currentStamp, currentStamp + 1);
场景推荐工具
高并发计数,只需最终 sumLongAdder
需要精确 CAS 的计数/版本号AtomicLong
多个字段原子更新synchronizedLock
对象引用原子替换AtomicReference
需要感知 ABA 的引用替换AtomicStampedReference

线程池

线程是重量级对象(需调用 OS 内核 API,分配数 MB 的栈空间),频繁创建/销毁代价极高。线程池复用线程,采用生产者-消费者模式:调用方提交任务到队列,池中线程循环消费。

核心参数

import java.util.concurrent.*;
import java.util.concurrent.atomic.*;

ThreadPoolExecutor executor = new ThreadPoolExecutor(
    4,                                         // corePoolSize:核心线程数,常驻不回收
    8,                                         // maximumPoolSize:最大线程数上限
    60L, TimeUnit.SECONDS,                     // 非核心线程空闲多久后回收
    new ArrayBlockingQueue<Runnable>(200),      // workQueue(必须有界,防 OOM!)
    new ThreadFactory() {                       // 自定义线程名,便于 jstack 排查
        private final AtomicInteger seq = new AtomicInteger(0);
        @Override
        public Thread newThread(Runnable r) {
            return new Thread(r, "biz-worker-" + seq.getAndIncrement());
        }
    },
    new ThreadPoolExecutor.CallerRunsPolicy()   // 队列满且达上限时:让提交方自己执行(降级)
);

任务提交流转

提交任务
  ↓
当前线程数 < corePoolSize?
  → 是:新建核心线程立即执行
  → 否:任务队列未满?
        → 是:入队等待
        → 否:当前线程数 < maximumPoolSize?
              → 是:新建临时线程执行
              → 否:执行拒绝策略

四种拒绝策略

策略行为适用场景
AbortPolicy(默认)抛出 RejectedExecutionException需要感知过载
CallerRunsPolicy由提交线程自己执行任务降级处理,不丢任务
DiscardPolicy直接丢弃,不抛异常可接受丢失(如日志异步写)
DiscardOldestPolicy丢弃队列最老的任务优先保证最新任务

线程数设置

int cpuCount = Runtime.getRuntime().availableProcessors();

// CPU 密集型(大量计算、加密、图像处理):+1 防止线程偶尔阻塞时 CPU 空闲
int cpuBoundSize = cpuCount + 1;

// IO 密集型(数据库、HTTP 调用、文件读写):大量时间在等待 IO
// 经验值:核心数 × 2,根据监控(队列积压、CPU 利用率)动态调整
int ioBoundSize = cpuCount * 2;

实战注意事项

// ===== ⚠️ 禁止使用 Executors 快捷方法 =====
// ❌ 默认使用无界 LinkedBlockingQueue,高负载时 OOM
ExecutorService badPool = Executors.newFixedThreadPool(4);
// ❌ 无上限创建线程,突发流量时 OOM
ExecutorService badPool2 = Executors.newCachedThreadPool();

// ✅ 总是显式指定有界队列
ThreadPoolExecutor goodPool = new ThreadPoolExecutor(
        4, 4, 0L, TimeUnit.MILLISECONDS,
        new ArrayBlockingQueue<Runnable>(500),
        new ThreadPoolExecutor.CallerRunsPolicy());

// ===== ⚠️ submit() 的任务内部异常会被吞掉 =====
// ❌ 这里不会打印任何异常
goodPool.execute(new Runnable() {
    @Override
    public void run() {
        throw new RuntimeException("我被吞掉了");
    }
});

// ✅ 用 Future.get() 获取任务内部异常
Future<?> future = goodPool.submit(new Runnable() {
    @Override
    public void run() {
        throw new RuntimeException("任务内部异常");
    }
});
try {
    future.get();
} catch (ExecutionException e) {
    System.err.println("任务失败:" + e.getCause());  // 真正的异常在 getCause() 里
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
}

// ===== 优雅关闭 =====
goodPool.shutdown();  // 停止接受新任务,等已提交任务执行完
try {
    if (!goodPool.awaitTermination(30, TimeUnit.SECONDS)) {
        goodPool.shutdownNow();   // 超过 30 秒强制停止
    }
} catch (InterruptedException e) {
    goodPool.shutdownNow();
    Thread.currentThread().interrupt();
}

快速选型指南

flowchart TD
    A["需要线程安全"] --> B{"只操作单个变量?"}
    B -->|"是"| C{"有 check-then-act<br/>等复合操作?"}
    C -->|"否(单步 increment)"| D["原子类<br/>AtomicXxx / LongAdder"]
    C -->|"是"| E{"需要超时/中断<br/>或精确条件唤醒?"}
    B -->|"否(多变量/代码块)"| E
    E -->|"否"| F["synchronized"]
    E -->|"是"| G["ReentrantLock"]

    H["共享数据结构"] --> I{"类型"}
    I -->|"Map"| J["ConcurrentHashMap"]
    I -->|"List(读极多)"| K["CopyOnWriteArrayList"]
    I -->|"队列/生产消费"| L["ArrayBlockingQueue"]

    M["协调多线程"] --> N{"需求"}
    N -->|"等待 N 个任务完成"| O["CountDownLatch"]
    N -->|"N 线程同步集合点"| P["CyclicBarrier"]
    N -->|"限制最大并发数"| Q["Semaphore"]

    style D fill:#cfc,stroke:#060
    style F fill:#cfc,stroke:#060
    style G fill:#ff9,stroke:#960
    style J fill:#cfc,stroke:#060
    style L fill:#cfc,stroke:#060
    style O fill:#e8f4f8,stroke:#4a90d9
    style P fill:#e8f4f8,stroke:#4a90d9
    style Q fill:#e8f4f8,stroke:#4a90d9

参考资料

← 返回列表
(1 人打了分,平均分: 5.00)

评论 (0)

暂无评论,来留下第一条吧。
登录注册 后才能发表评论