破茧成蝶:Java后端从0到资深工程师的进阶之路(五)
破茧成蝶:Java后端从0到资深工程师的进阶之路(五)并发篇——多线程与高并发实战
现代后端系统,高并发是绕不开的挑战。多线程编程就像一把双刃剑:用得好了,系统吞吐量飙升;用得不好,死锁、内存泄漏、性能骤降接踵而至。本篇将带你深入 Java 并发核心,从
CompletableFuture异步编排到线程池调优,再到并发容器的底层原理,助你成为真正的并发编程高手。
写在前面
很多开发者提到并发,第一反应是 synchronized 和 volatile,但对 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 提供了 orTimeout 和 completeOnTimeout 方法(Java 9+)。
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
// 模拟耗时操作
Thread.sleep(5000);
return "result";
}).orTimeout(3, TimeUnit.SECONDS)
.exceptionally(ex -> "timeout"); // 超时后返回默认值
1.1.3 异常处理
CompletableFuture 中的异常不会直接抛出,而是通过 exceptionally 或 handle 捕获。
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 是阿里开源的组件,专门解决线程池场景下上下文传递的问题。它通过 TtlRunnable 或 TtlExecutorService 包装任务,在任务执行前将父线程的上下文复制到当前线程,执行后清除,实现安全传递。
引入依赖:
<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 实现思路
- 封装一个
DynamicThreadPool,持有ThreadPoolExecutor实例。 - 提供方法更新核心线程数、最大线程数、队列容量等参数。
- 监听配置变化,调用
setCorePoolSize、setMaximumPoolSize等方法。 - 注意
setCorePoolSize和setMaximumPoolSize的规则:如果新值大于旧值,会立即创建新线程;如果小于旧值,会在空闲时回收。
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);
}
}
💡 资深提示:动态调整线程池参数时,需要关注
corePoolSize和maximumPoolSize的调整策略,避免频繁创建销毁线程。同时,队列容量的动态替换需要格外小心,推荐使用可扩容的队列实现(如自定义ResizableCapacityLinkedBlockingQueue)。
三、并发容器与锁优化
3.1 ConcurrentHashMap 的 size 方法原理
ConcurrentHashMap 的 size() 方法并不像普通 HashMap 那样直接返回一个字段,而是通过一种近似统计的方式实现,以平衡并发与准确性。
JDK 1.8 的实现:
- 维护一个
baseCount和CounterCell数组(类似分段计数)。 - 当多个线程同时更新时,每个线程会随机选择
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;
}
结论:ConcurrentHashMap 的 size() 方法并不是精确的,在高并发场景下可能返回一个旧值。如果需要精确计数,应使用 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 业务选择建议
| 场景 | 推荐工具 | 原因 |
|---|---|---|
| 缓存(读远多于写) | ReentrantReadWriteLock 或 StampedLock |
读并发高,写锁影响小 |
| 高频读写(如计数器) | LongAdder |
分段思想,无锁,性能极高 |
| 通用并发 Map | ConcurrentHashMap |
综合性能最好,无需额外锁 |
| 需要精确控制锁粒度 | 分段锁(自定义) | 可基于 key 分段,减少竞争 |
💡 资深提示:不要盲目使用读写锁。如果读操作非常快(例如内存读取),使用读写锁带来的上下文切换开销可能超过锁竞争的开销,此时普通锁可能更快。性能调优应基于压测数据。
总结
本篇我们从并发编程的三个关键领域深入:
-
JUC 核心类库:
CompletableFuture实现异步任务编排、超时与异常处理,提升系统吞吐量。ThreadLocal内存泄漏分析与TransmittableThreadLocal在线程池场景下的安全传递。
-
线程池避坑:
- 避免使用
Executors创建无界队列或无限线程的线程池。 - 实现动态线程池参数调整,应对流量变化。
- 避免使用
-
并发容器与锁优化:
- 理解
ConcurrentHashMap的size()方法弱一致性。 - 根据业务场景选择读写锁、分段锁或无锁结构,追求极致性能。
- 理解
并发编程是 Java 后端进阶的必经之路,理解这些原理和实战技巧,你将能在高并发场景下写出更稳定、更高效的代码。
下篇预告: 《中间件篇——消息队列与数据缓存的博弈》将带你深入 Redis 缓存三大坑、消息队列的可靠性保证以及分布式事务的最终一致性方案,敬请期待!
如果觉得本文对你有帮助,欢迎点赞、收藏、评论,你的支持是我持续创作的动力!
更多推荐


所有评论(0)