破茧成蝶:Java后端从0到资深工程师的进阶之路(五)并发篇——多线程与高并发实战

现代后端系统,高并发是绕不开的挑战。多线程编程就像一把双刃剑:用得好了,系统吞吐量飙升;用得不好,死锁、内存泄漏、性能骤降接踵而至。本篇将带你深入 Java 并发核心,从 CompletableFuture 异步编排到线程池调优,再到并发容器的底层原理,助你成为真正的并发编程高手。


写在前面

很多开发者提到并发,第一反应是 synchronizedvolatile,但对 CompletableFuture 的异步编排、ThreadLocal 的内存泄漏、线程池的动态调优、ConcurrentHashMap 的扩容机制等缺乏深入理解。结果就是:线上突然出现 OOM、CPU 飙升、线程阻塞,却不知道从何下手。

一个资深开发者眼中的并发:

  • 异步编程:能熟练使用 CompletableFuture 编排任务依赖,合理设置超时与异常处理。
  • 线程池:能根据业务场景选择合适的线程池类型,并实现参数的动态调整。
  • 并发容器:能深入理解 ConcurrentHashMap 的扩容机制,根据场景选择读写锁或分段锁,避免锁竞争。

本篇文章,我们将从这些核心点出发,结合实战案例,帮你构建稳固的并发知识体系。


一、JUC 核心类库的深度解析

1.1 CompletableFuture 实现异步编排(任务间的依赖、超时控制、异常处理)

CompletableFuture 是 Java 8 引入的异步编程利器,它将回调、任务编排、异常处理融为一体,让异步代码不再陷入“回调地狱”。

1.1.1 任务间的依赖编排

场景:查询用户信息、订单信息、商品信息,三者并行执行,最后汇总结果。

  • 用户信息与订单信息无依赖,可并行。
  • 商品信息依赖订单信息(从订单中获取商品 ID)。

代码实现

public CompletableFuture<UserDTO> queryUser(Long userId) {
    return CompletableFuture.supplyAsync(() -> userService.getById(userId));
}

public CompletableFuture<OrderDTO> queryOrder(Long orderId) {
    return CompletableFuture.supplyAsync(() -> orderService.getById(orderId));
}

public CompletableFuture<ProductDTO> queryProductByOrder(OrderDTO order) {
    return CompletableFuture.supplyAsync(() -> productService.getById(order.getProductId()));
}

// 编排
public CompletableFuture<ResultDTO> assembleResult(Long userId, Long orderId) {
    CompletableFuture<UserDTO> userFuture = queryUser(userId);
    CompletableFuture<OrderDTO> orderFuture = queryOrder(orderId);
    
    // 等订单查询完成后,根据订单查询商品(依赖)
    CompletableFuture<ProductDTO> productFuture = orderFuture.thenCompose(this::queryProductByOrder);
    
    // 三个任务全部完成后,合并结果
    return CompletableFuture.allOf(userFuture, orderFuture, productFuture)
            .thenApply(v -> {
                ResultDTO result = new ResultDTO();
                result.setUser(userFuture.join());
                result.setOrder(orderFuture.join());
                result.setProduct(productFuture.join());
                return result;
            });
}

常用方法

  • thenApply:转换结果(同步)。
  • thenCompose:链式组合(返回新的 CompletableFuture)。
  • thenCombine:两个任务并行,合并结果。
  • allOf:等待所有任务完成。
  • anyOf:任一任务完成即返回。
1.1.2 超时控制

避免任务无限阻塞,CompletableFuture 提供了 orTimeoutcompleteOnTimeout 方法(Java 9+)。

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    // 模拟耗时操作
    Thread.sleep(5000);
    return "result";
}).orTimeout(3, TimeUnit.SECONDS)
  .exceptionally(ex -> "timeout"); // 超时后返回默认值
1.1.3 异常处理

CompletableFuture 中的异常不会直接抛出,而是通过 exceptionallyhandle 捕获。

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    if (true) throw new RuntimeException("业务异常");
    return "success";
}).exceptionally(ex -> {
    log.error("执行失败", ex);
    return "fallback";
});

💡 资深提示CompletableFuture 默认使用 ForkJoinPool.commonPool(),其线程数是 CPU 核心数 - 1。对于 IO 密集型任务,应自定义线程池,避免与计算密集型任务争抢。


1.2 ThreadLocal 的内存泄漏分析与实战(InheritableThreadLocal 的缺陷与阿里 TransmittableThreadLocal 的解决方案)

1.2.1 ThreadLocal 的内存泄漏原理

ThreadLocal 内部使用 ThreadLocalMap 存储数据,其 Key 是 ThreadLocal 对象的弱引用。当 ThreadLocal 对象不再被强引用时,下一次 GC 会被回收,但 Value 仍然是强引用,且存在于 Thread 的 ThreadLocalMap 中。如果线程一直存活(如线程池中的线程),那么这些 Value 永远不会被回收,导致内存泄漏。

典型泄漏场景:使用线程池,但在任务结束后没有调用 remove() 清除 ThreadLocal 值。

正确用法

try {
    threadLocal.set(value);
    // 执行业务
} finally {
    threadLocal.remove(); // 务必移除
}
1.2.2 InheritableThreadLocal 的缺陷

InheritableThreadLocal 允许子线程继承父线程的变量值,但在线程池场景下存在严重问题:线程池复用了线程,子线程的变量值可能被复用,导致数据错乱。

1.2.3 阿里 TransmittableThreadLocal 解决方案

TransmittableThreadLocal 是阿里开源的组件,专门解决线程池场景下上下文传递的问题。它通过 TtlRunnableTtlExecutorService 包装任务,在任务执行前将父线程的上下文复制到当前线程,执行后清除,实现安全传递。

引入依赖

<dependency>
    <groupId>com.alibaba</groupId>
    <artifactId>transmittable-thread-local</artifactId>
    <version>2.14.2</version>
</dependency>

使用示例

// 定义 TransmittableThreadLocal
TransmittableThreadLocal<String> context = new TransmittableThreadLocal<>();

// 设置值
context.set("userId");

// 包装线程池
ExecutorService executor = TtlExecutors.getTtlExecutorService(Executors.newFixedThreadPool(10));

// 提交任务时自动复制上下文
executor.submit(() -> {
    System.out.println(context.get()); // 能正确获取父线程的 userId
});

// 最后清理(可选)
context.remove();

💡 资深提示:在生产环境中,使用线程池且需要传递上下文(如链路追踪 traceId、用户信息)时,强烈推荐使用 TransmittableThreadLocal,避免手动传递的繁琐和隐患。


二、线程池的“避坑指南”

2.1 Executors 工厂类创建的常见陷阱(OOM 风险)

Executors 提供了快速创建线程池的方法,但隐藏着致命缺陷:

方法 队列 风险
newFixedThreadPool LinkedBlockingQueue (无界队列) 任务堆积可能导致 OOM
newCachedThreadPool SynchronousQueue (无容量) 线程数无上限,可能创建过多线程导致系统崩溃
newSingleThreadExecutor LinkedBlockingQueue (无界队列) 同 fixed,任务堆积 OOM
newScheduledThreadPool DelayedWorkQueue (无界队列) 同 fixed,任务堆积 OOM

案例:某电商促销期间,使用 newFixedThreadPool 处理订单异步任务,高峰期任务堆积数千万,最终导致内存溢出,整个服务宕机。

正确做法:手动创建 ThreadPoolExecutor,明确设置队列大小和拒绝策略。

int corePoolSize = Runtime.getRuntime().availableProcessors();
int maxPoolSize = corePoolSize * 2;
long keepAliveTime = 60L;
BlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<>(2000);
RejectedExecutionHandler handler = new ThreadPoolExecutor.CallerRunsPolicy();

ThreadPoolExecutor executor = new ThreadPoolExecutor(
        corePoolSize, maxPoolSize, keepAliveTime, TimeUnit.SECONDS,
        workQueue, 
        Executors.defaultThreadFactory(),
        handler
);

拒绝策略选择

  • AbortPolicy:直接抛异常(默认),适合关键任务,需要上游重试。
  • CallerRunsPolicy:由调用线程执行,适合非关键任务,可降级。
  • DiscardPolicy:直接丢弃,不通知,慎用。
  • DiscardOldestPolicy:丢弃队列中最老的任务,适合实时性要求高的场景。

2.2 如何动态调整线程池参数(美团线程池动态调优方案的简易实现)

生产环境中,业务流量是动态变化的,静态配置的线程池参数往往无法完美适应。美团技术团队提出的动态线程池思路值得借鉴:通过配置中心(如 Apollo、Nacos)动态修改核心参数,并实时生效。

2.2.1 实现思路
  1. 封装一个 DynamicThreadPool,持有 ThreadPoolExecutor 实例。
  2. 提供方法更新核心线程数、最大线程数、队列容量等参数。
  3. 监听配置变化,调用 setCorePoolSizesetMaximumPoolSize 等方法。
  4. 注意 setCorePoolSizesetMaximumPoolSize 的规则:如果新值大于旧值,会立即创建新线程;如果小于旧值,会在空闲时回收。
2.2.2 简易实现
@Component
public class DynamicThreadPool {
    private final ThreadPoolExecutor executor;
    
    public DynamicThreadPool() {
        this.executor = new ThreadPoolExecutor(
                10, 20, 60, TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(1000),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );
    }
    
    public void updateCorePoolSize(int corePoolSize) {
        executor.setCorePoolSize(corePoolSize);
    }
    
    public void updateMaxPoolSize(int maxPoolSize) {
        executor.setMaximumPoolSize(maxPoolSize);
    }
    
    public void updateQueueCapacity(int capacity) {
        BlockingQueue<Runnable> newQueue = new ArrayBlockingQueue<>(capacity);
        // 注意:直接替换队列会丢失原有任务,需要谨慎实现
        // 更安全的方式是使用 ResizableCapacityLinkedBlockingQueue(可修改容量的队列)
    }
    
    // 提交任务的方法...
}
2.2.3 配合配置中心(以 Nacos 为例)
@RefreshScope
@Component
public class ThreadPoolConfig {
    @Value("${threadpool.coreSize:10}")
    private int coreSize;
    
    @Value("${threadpool.maxSize:20}")
    private int maxSize;
    
    @Autowired
    private DynamicThreadPool dynamicThreadPool;
    
    @PostConstruct
    public void init() {
        dynamicThreadPool.updateCorePoolSize(coreSize);
        dynamicThreadPool.updateMaxPoolSize(maxSize);
    }
    
    @EventListener
    public void onRefresh(RefreshScopeRefreshedEvent event) {
        // 配置刷新后重新设置
        dynamicThreadPool.updateCorePoolSize(coreSize);
        dynamicThreadPool.updateMaxPoolSize(maxSize);
    }
}

💡 资深提示:动态调整线程池参数时,需要关注 corePoolSizemaximumPoolSize 的调整策略,避免频繁创建销毁线程。同时,队列容量的动态替换需要格外小心,推荐使用可扩容的队列实现(如自定义 ResizableCapacityLinkedBlockingQueue)。


三、并发容器与锁优化

3.1 ConcurrentHashMapsize 方法原理

ConcurrentHashMapsize() 方法并不像普通 HashMap 那样直接返回一个字段,而是通过一种近似统计的方式实现,以平衡并发与准确性。

JDK 1.8 的实现

  • 维护一个 baseCountCounterCell 数组(类似分段计数)。
  • 当多个线程同时更新时,每个线程会随机选择 CounterCell 进行累加,减少对 baseCount 的竞争。
  • size() 方法会累加 baseCount 和所有 CounterCell 的值,但这个过程不加锁,因此得到的 size 是弱一致性的(可能不是最新值)。

源码片段(JDK 1.8)

public int size() {
    long n = sumCount();
    return ((n < 0L) ? 0 : (n > (long)Integer.MAX_VALUE) ? Integer.MAX_VALUE : (int)n);
}

final long sumCount() {
    CounterCell[] as = counterCells;
    CounterCell a;
    long sum = baseCount;
    if (as != null) {
        for (int i = 0; i < as.length; ++i) {
            if ((a = as[i]) != null)
                sum += a.value;
        }
    }
    return sum;
}

结论ConcurrentHashMapsize() 方法并不是精确的,在高并发场景下可能返回一个旧值。如果需要精确计数,应使用 LongAdder 或维护独立计数器。


3.2 分段锁、读写锁在实际业务场景下的性能对比

3.2.1 分段锁(JDK 1.7 ConcurrentHashMap)

JDK 1.7 的 ConcurrentHashMap 采用分段锁(Segment),默认 16 个 Segment,每个 Segment 独立加锁,将锁粒度从整个表缩小到单个 Segment,从而提升并发度。

适用场景:读写比例均衡,但 JDK 1.8 已放弃分段锁,改用 CAS + synchronized 优化,性能更好。

3.2.2 读写锁(ReentrantReadWriteLock)

ReentrantReadWriteLock 允许多个读线程并发,写线程互斥。适合读多写少的场景。

示例:缓存实现

public class Cache<K, V> {
    private final Map<K, V> cache = new HashMap<>();
    private final ReadWriteLock lock = new ReentrantReadWriteLock();
    private final Lock readLock = lock.readLock();
    private final Lock writeLock = lock.writeLock();
    
    public V get(K key) {
        readLock.lock();
        try {
            return cache.get(key);
        } finally {
            readLock.unlock();
        }
    }
    
    public void put(K key, V value) {
        writeLock.lock();
        try {
            cache.put(key, value);
        } finally {
            writeLock.unlock();
        }
    }
}

性能对比

  • 读多写少:读写锁 >> 普通锁(synchronized) > 分段锁(如果分段数较少)。
  • 写多:分段锁或 ConcurrentHashMap 可能更优,因为读写锁的写操作会阻塞所有读操作。
  • Java 8+StampedLock 提供了乐观读,性能比 ReentrantReadWriteLock 更高,但使用更复杂。
3.2.3 业务选择建议
场景 推荐工具 原因
缓存(读远多于写) ReentrantReadWriteLockStampedLock 读并发高,写锁影响小
高频读写(如计数器) LongAdder 分段思想,无锁,性能极高
通用并发 Map ConcurrentHashMap 综合性能最好,无需额外锁
需要精确控制锁粒度 分段锁(自定义) 可基于 key 分段,减少竞争

💡 资深提示:不要盲目使用读写锁。如果读操作非常快(例如内存读取),使用读写锁带来的上下文切换开销可能超过锁竞争的开销,此时普通锁可能更快。性能调优应基于压测数据。


总结

本篇我们从并发编程的三个关键领域深入:

  1. JUC 核心类库

    • CompletableFuture 实现异步任务编排、超时与异常处理,提升系统吞吐量。
    • ThreadLocal 内存泄漏分析与 TransmittableThreadLocal 在线程池场景下的安全传递。
  2. 线程池避坑

    • 避免使用 Executors 创建无界队列或无限线程的线程池。
    • 实现动态线程池参数调整,应对流量变化。
  3. 并发容器与锁优化

    • 理解 ConcurrentHashMapsize() 方法弱一致性。
    • 根据业务场景选择读写锁、分段锁或无锁结构,追求极致性能。

并发编程是 Java 后端进阶的必经之路,理解这些原理和实战技巧,你将能在高并发场景下写出更稳定、更高效的代码。

下篇预告: 《中间件篇——消息队列与数据缓存的博弈》将带你深入 Redis 缓存三大坑、消息队列的可靠性保证以及分布式事务的最终一致性方案,敬请期待!


如果觉得本文对你有帮助,欢迎点赞、收藏、评论,你的支持是我持续创作的动力!

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐