Java中@Async注解的线程池隔离陷阱
作者:linmoo2006
作为多年的Java开发经验,在开发过程中经常会踩一些坑,本系列想通过一些案例分享,帮助其他开发者避免这些问题。
注意:由于框架不同版本改造会有些使用的不同,因此本次系列中使用JDK版本使用的是open-jdk21。
1. 事情起因
在一次电商系统的订单处理系统中,用户反馈系统响应越来越慢,最终导致服务不可用。经过排查发现,是因为大量使用了@Async注解进行异步处理,但没有配置自定义线程池,导致系统创建了大量线程,最终内存溢出(OOM)。
问题代码如下:
参考代码 lesson16-async-threadpool 中的AsyncThreadPoolDemo.java
package com.architect.pitfalls.async.cause;
import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.Async;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import java.lang.reflect.Method;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executor;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.atomic.AtomicInteger;
@Configuration
@EnableAsync
@ComponentScan
public class AsyncThreadPoolDemo {
public static void main(String[] args) throws InterruptedException {
System.out.println("=== @Async注解线程池隔离陷阱演示 ===\n");
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(AsyncThreadPoolDemo.class);
System.out.println("========================================");
System.out.println("场景1: 默认SimpleAsyncTaskExecutor问题");
System.out.println("========================================");
demonstrateDefaultExecutor(context);
context.close();
System.out.println("\n========================================");
System.out.println("场景2: 共享线程池导致的问题");
System.out.println("========================================");
demonstrateSharedThreadPool();
System.out.println("\n========================================");
System.out.println("场景3: 线程池耗尽导致服务阻塞");
System.out.println("========================================");
demonstrateThreadPoolExhaustion();
}
private static void demonstrateDefaultExecutor(AnnotationConfigApplicationContext context) throws InterruptedException {
System.out.println();
System.out.println("问题说明:");
System.out.println(" Spring Boot默认使用SimpleAsyncTaskExecutor");
System.out.println(" 这个类不是真正的线程池,每次调用都创建新线程");
System.out.println();
DefaultAsyncService service = context.getBean(DefaultAsyncService.class);
AtomicInteger threadCount = new AtomicInteger(0);
int taskCount = 20;
CountDownLatch latch = new CountDownLatch(taskCount);
System.out.println("步骤1: 并发调用" + taskCount + "个异步任务");
long startTime = System.currentTimeMillis();
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
service.executeAsync(() -> {
threadCount.incrementAndGet();
String threadName = Thread.currentThread().getName();
System.out.println(" 任务" + taskId + " 执行线程: " + threadName);
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
latch.countDown();
});
}
latch.await();
long duration = System.currentTimeMillis() - startTime;
System.out.println();
System.out.println("步骤2: 分析结果");
System.out.println(" 执行时间: " + duration + "ms");
System.out.println(" 创建的线程数: " + threadCount.get());
System.out.println();
System.out.println("问题分析:");
System.out.println(" ⚠️ 每个任务都创建了新线程,线程名不同");
System.out.println(" ⚠️ 高并发时可能创建大量线程,导致OOM");
System.out.println(" ⚠️ 线程创建销毁开销大,性能差");
System.out.println();
System.out.println("严重后果:");
System.out.println(" 1. 内存溢出(OOM)");
System.out.println(" 2. 线程创建开销导致性能下降");
System.out.println(" 3. 系统资源耗尽");
}
private static void demonstrateSharedThreadPool() throws InterruptedException {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(SharedPoolConfig.class);
System.out.println();
System.out.println("问题说明:");
System.out.println(" 多个业务共用一个线程池");
System.out.println(" 一个业务的慢任务会阻塞其他业务");
System.out.println();
FastService fastService = context.getBean(FastService.class);
SlowService slowService = context.getBean(SlowService.class);
int fastTaskCount = 10;
int slowTaskCount = 5;
CountDownLatch latch = new CountDownLatch(fastTaskCount + slowTaskCount);
System.out.println("步骤1: 同时提交快速任务和慢任务");
System.out.println(" 快速任务数: " + fastTaskCount);
System.out.println(" 慢任务数: " + slowTaskCount);
System.out.println(" 线程池大小: 5");
System.out.println();
long startTime = System.currentTimeMillis();
for (int i = 0; i < slowTaskCount; i++) {
slowService.executeSlowTask(() -> {
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
latch.countDown();
});
}
for (int i = 0; i < fastTaskCount; i++) {
final int taskId = i;
fastService.executeFastTask(() -> {
long waitTime = System.currentTimeMillis() - startTime;
System.out.println(" 快速任务" + taskId + " 开始执行,等待时间: " + waitTime + "ms");
latch.countDown();
});
}
latch.await();
long duration = System.currentTimeMillis() - startTime;
System.out.println();
System.out.println("步骤2: 分析结果");
System.out.println(" 总执行时间: " + duration + "ms");
System.out.println();
System.out.println("问题分析:");
System.out.println(" ⚠️ 慢任务占满了线程池");
System.out.println(" ⚠️ 快速任务被阻塞,无法及时执行");
System.out.println(" ⚠️ 业务之间相互影响");
System.out.println();
System.out.println("严重后果:");
System.out.println(" 1. 核心业务被非核心业务阻塞");
System.out.println(" 2. 系统响应时间不可控");
System.out.println(" 3. 故障隔离能力缺失");
context.close();
}
private static void demonstrateThreadPoolExhaustion() throws InterruptedException {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(ExhaustionConfig.class);
System.out.println();
System.out.println("问题说明:");
System.out.println(" 线程池配置不当,队列满后任务被拒绝");
System.out.println(" 或者线程池耗尽导致服务不可用");
System.out.println();
ExhaustionService service = context.getBean(ExhaustionService.class);
int taskCount = 30;
AtomicInteger successCount = new AtomicInteger(0);
AtomicInteger rejectCount = new AtomicInteger(0);
CountDownLatch latch = new CountDownLatch(taskCount);
System.out.println("步骤1: 提交" + taskCount + "个任务到小容量线程池");
System.out.println(" 线程池核心大小: 2");
System.out.println(" 线程池最大大小: 3");
System.out.println(" 队列容量: 5");
System.out.println();
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
try {
service.executeTask(() -> {
successCount.incrementAndGet();
try {
Thread.sleep(500);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
latch.countDown();
});
} catch (Exception e) {
rejectCount.incrementAndGet();
System.out.println(" 任务" + taskId + " 被拒绝: " + e.getClass().getSimpleName());
latch.countDown();
}
}
latch.await();
System.out.println();
System.out.println("步骤2: 分析结果");
System.out.println(" 成功执行: " + successCount.get());
System.out.println(" 被拒绝: " + rejectCount.get());
System.out.println();
System.out.println("问题分析:");
System.out.println(" ⚠️ 线程池容量不足,任务被拒绝");
System.out.println(" ⚠️ 业务请求失败,用户体验差");
System.out.println(" ⚠️ 没有合理的拒绝策略");
System.out.println();
System.out.println("严重后果:");
System.out.println(" 1. 业务请求失败");
System.out.println(" 2. 数据丢失");
System.out.println(" 3. 用户投诉");
context.close();
}
}
@Service
class DefaultAsyncService {
@Async
public void executeAsync(Runnable task) {
task.run();
}
}
@Configuration
@EnableAsync
class SharedPoolConfig {
@Bean(name = "sharedExecutor")
public Executor sharedExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(5);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("shared-");
executor.initialize();
return executor;
}
@Bean
public FastService fastService() {
return new FastService();
}
@Bean
public SlowService slowService() {
return new SlowService();
}
}
@Service
class FastService {
@Async("sharedExecutor")
public void executeFastTask(Runnable task) {
task.run();
}
}
@Service
class SlowService {
@Async("sharedExecutor")
public void executeSlowTask(Runnable task) {
task.run();
}
}
@Configuration
@EnableAsync
class ExhaustionConfig implements AsyncUncaughtExceptionHandler {
@Bean(name = "limitedExecutor")
public Executor limitedExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(2);
executor.setMaxPoolSize(3);
executor.setQueueCapacity(5);
executor.setThreadNamePrefix("limited-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());
executor.initialize();
return executor;
}
@Bean
public ExhaustionService exhaustionService() {
return new ExhaustionService();
}
@Override
public void handleUncaughtException(Throwable ex, Method method, Object... params) {
System.out.println(" 异步任务异常: " + ex.getMessage());
}
}
@Service
class ExhaustionService {
@Async("limitedExecutor")
public void executeTask(Runnable task) {
task.run();
}
}
运行结果:
=== @Async注解线程池隔离陷阱演示 === ======================================== 场景1: 默认SimpleAsyncTaskExecutor问题 ======================================== 问题说明: Spring Boot默认使用SimpleAsyncTaskExecutor 这个类不是真正的线程池,每次调用都创建新线程 步骤1: 并发调用20个异步任务 任务0 执行线程: SimpleAsyncTaskExecutor-1 任务1 执行线程: SimpleAsyncTaskExecutor-2 任务2 执行线程: SimpleAsyncTaskExecutor-3 ... 任务19 执行线程: SimpleAsyncTaskExecutor-20 步骤2: 分析结果 执行时间: 150ms 创建的线程数: 20 问题分析: ⚠️ 每个任务都创建了新线程,线程名不同 ⚠️ 高并发时可能创建大量线程,导致OOM ⚠️ 线程创建销毁开销大,性能差 严重后果: 1. 内存溢出(OOM) 2. 线程创建开销导致性能下降 3. 系统资源耗尽 ======================================== 场景2: 共享线程池导致的问题 ======================================== 问题说明: 多个业务共用一个线程池 一个业务的慢任务会阻塞其他业务 步骤1: 同时提交快速任务和慢任务 快速任务数: 10 慢任务数: 5 线程池大小: 5 快速任务0 开始执行,等待时间: 5ms 快速任务1 开始执行,等待时间: 6ms ... 快速任务4 开始执行,等待时间: 8ms 快速任务5 开始执行,等待时间: 2010ms ← 被阻塞了2秒 快速任务6 开始执行,等待时间: 2011ms ... 步骤2: 分析结果 总执行时间: 4015ms 问题分析: ⚠️ 慢任务占满了线程池 ⚠️ 快速任务被阻塞,无法及时执行 ⚠️ 业务之间相互影响 严重后果: 1. 核心业务被非核心业务阻塞 2. 系统响应时间不可控 3. 故障隔离能力缺失 ======================================== 场景3: 线程池耗尽导致服务阻塞 ======================================== 问题说明: 线程池配置不当,队列满后任务被拒绝 或者线程池耗尽导致服务不可用 步骤1: 提交30个任务到小容量线程池 线程池核心大小: 2 线程池最大大小: 3 队列容量: 5 任务8 被拒绝: TaskRejectedException 任务9 被拒绝: TaskRejectedException ... 步骤2: 分析结果 成功执行: 8 被拒绝: 22 问题分析: ⚠️ 线程池容量不足,任务被拒绝 ⚠️ 业务请求失败,用户体验差 ⚠️ 没有合理的拒绝策略 严重后果: 1. 业务请求失败 2. 数据丢失 3. 用户投诉
2. 原因分析
2.1 Spring @Async的线程池查找机制
Spring在执行@Async注解的方法时,会按照以下顺序查找线程池:
1. 首先查找名为"taskExecutor"的Bean
2. 如果没找到,查找实现了TaskExecutor接口的Bean
3. 如果都没找到,使用默认的SimpleAsyncTaskExecutor
关键源码分析(Spring Framework 6.0.x):
// AsyncExecutionInterceptor.java
protected Executor determineExecutor(Method method) {
Executor executor = this.executors.get(method);
if (executor == null) {
Executor targetExecutor = findExecutor(method);
this.executors.putIfAbsent(method, targetExecutor);
executor = targetExecutor;
}
return executor;
}
// AsyncExecutionAspectSupport.java
protected Executor getDefaultExecutor(@Nullable BeanFactory beanFactory) {
if (beanFactory != null) {
// 1. 查找名为"taskExecutor"的Bean
try {
return beanFactory.getBean(TaskExecutor.class);
} catch (NoUniqueBeanDefinitionException ex) {
// 多个TaskExecutor时,查找名为taskExecutor的
} catch (NoSuchBeanDefinitionException ex) {
// 没找到,继续
}
// 2. 查找名为"taskExecutor"的Bean
try {
return beanFactory.getBean("taskExecutor", Executor.class);
} catch (NoSuchBeanDefinitionException ex) {
// 3. 使用默认的SimpleAsyncTaskExecutor
return new SimpleAsyncTaskExecutor();
}
}
return new SimpleAsyncTaskExecutor();
}
2.2 SimpleAsyncTaskExecutor的问题
SimpleAsyncTaskExecutor不是真正的线程池,它有以下问题:
// SimpleAsyncTaskExecutor.java
public void execute(Runnable task) {
Thread thread = new Thread(getThreadGroup(), task, nextThreadName());
thread.setPriority(this.threadPriority);
thread.setDaemon(this.daemon);
thread.start(); // 每次都创建新线程!
}
问题本质:
- 每次执行任务都创建新线程
- 没有线程复用机制
- 没有队列缓冲
- 没有线程数量限制
- 高并发时会导致OOM
2.3 三种常见陷阱场景
场景1:默认线程池陷阱
- 没有配置自定义线程池
- 使用默认的SimpleAsyncTaskExecutor
- 高并发时创建大量线程
场景2:共享线程池陷阱
- 多个业务共用一个线程池
- 一个业务的慢任务阻塞其他业务
- 故障无法隔离
场景3:线程池配置不当陷阱
- 线程池参数配置不合理
- 拒绝策略选择不当
- 没有监控和告警
3. 解决方案
3.1 方案一:配置自定义线程池
手动的创建线程池
// 自定义关键代码
@Configuration
@EnableAsync
class CustomPoolConfig {
@Bean(name = "customExecutor")
public Executor customExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("custom-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(60);
executor.initialize();
return executor;
}
@Bean
public CustomAsyncService customAsyncService() {
return new CustomAsyncService();
}
}
优点: 线程复用、资源可控、可监控
缺点: 需要手动配置参数
适用场景: 所有使用@Async的场景
3.2 方案二:线程池隔离
为不同业务配置独立的线程池,实现业务隔离:
@Configuration
@EnableAsync
class IsolatedPoolConfig {
// 核心业务线程池
@Bean(name = "coreTaskExecutor")
public Executor coreTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(20);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("core-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
// 非核心业务线程池
@Bean(name = "nonCoreTaskExecutor")
public Executor nonCoreTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(3);
executor.setMaxPoolSize(5);
executor.setQueueCapacity(50);
executor.setThreadNamePrefix("non-core-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.DiscardOldestPolicy());
executor.initialize();
return executor;
}
}
优点: 业务隔离、故障隔离
缺点: 配置复杂、资源占用增加
适用场景: 多业务系统,核心业务需要保障
3.3 方案三:完善的线程池配置
@Bean(name = "taskExecutor")
public Executor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
// 核心线程数:CPU密集型 = CPU核心数,IO密集型 = CPU核心数 * 2
executor.setCorePoolSize(Runtime.getRuntime().availableProcessors());
// 最大线程数
executor.setMaxPoolSize(Runtime.getRuntime().availableProcessors() * 2);
// 队列容量
executor.setQueueCapacity(200);
// 线程名前缀
executor.setThreadNamePrefix("async-");
// 拒绝策略:调用者运行
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
// 优雅关闭
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(60);
// 允许核心线程超时
executor.setAllowCoreThreadTimeOut(true);
executor.setKeepAliveSeconds(60);
executor.initialize();
return executor;
}
优点: 配置完善、可监控、优雅关闭
缺点: 需要根据业务调整参数
适用场景: 生产环境推荐
3.4 最终代码演示
参考代码 lesson16-async-threadpool 中的CustomThreadPoolSolution.java
package com.architect.pitfalls.async.solution;
import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.Async;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import java.lang.reflect.Method;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executor;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.atomic.AtomicInteger;
@Configuration
@EnableAsync
@ComponentScan
public class CustomThreadPoolSolution {
public static void main(String[] args) throws InterruptedException {
System.out.println("=== 方案一:自定义线程池解决方案 ===\n");
demonstrateCustomThreadPool();
System.out.println("\n=== 方案二:线程池隔离解决方案 ===\n");
demonstrateIsolatedThreadPool();
System.out.println("\n=== 方案三:完善的线程池配置 ===\n");
demonstrateCompleteThreadPoolConfig();
}
private static void demonstrateCustomThreadPool() throws InterruptedException {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(CustomPoolConfig.class);
System.out.println("========================================");
System.out.println("自定义线程池演示");
System.out.println("========================================");
System.out.println();
System.out.println("方案说明:");
System.out.println(" 1. 配置自定义ThreadPoolTaskExecutor");
System.out.println(" 2. 设置合理的核心线程数、最大线程数、队列容量");
System.out.println(" 3. 配置合理的拒绝策略");
System.out.println();
CustomAsyncService service = context.getBean(CustomAsyncService.class);
int taskCount = 20;
AtomicInteger reusedThreadCount = new AtomicInteger(0);
CountDownLatch latch = new CountDownLatch(taskCount);
AtomicInteger uniqueThreads = new AtomicInteger(0);
java.util.Set<String> threadNames = new java.util.HashSet<>();
System.out.println("步骤1: 并发调用" + taskCount + "个异步任务");
long startTime = System.currentTimeMillis();
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
service.executeAsync(() -> {
String threadName = Thread.currentThread().getName();
synchronized (threadNames) {
if (threadNames.add(threadName)) {
uniqueThreads.incrementAndGet();
}
}
System.out.println(" 任务" + taskId + " 执行线程: " + threadName);
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
latch.countDown();
});
}
latch.await();
long duration = System.currentTimeMillis() - startTime;
System.out.println();
System.out.println("步骤2: 分析结果");
System.out.println(" 执行时间: " + duration + "ms");
System.out.println(" 创建的不同线程数: " + uniqueThreads.get());
System.out.println(" 总任务数: " + taskCount);
System.out.println();
System.out.println("分析:");
System.out.println(" ✅ 线程被复用,不是每次创建新线程");
System.out.println(" ✅ 线程数量可控,不会无限增长");
System.out.println(" ✅ 性能更好,资源使用更合理");
System.out.println();
System.out.println("优点:");
System.out.println(" - 线程复用,减少创建销毁开销");
System.out.println(" - 线程数量可控,避免OOM");
System.out.println(" - 可以监控线程池状态");
System.out.println();
System.out.println("缺点:");
System.out.println(" - 需要手动配置线程池参数");
System.out.println(" - 需要根据业务场景调整配置");
context.close();
}
private static void demonstrateIsolatedThreadPool() throws InterruptedException {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(IsolatedPoolConfig.class);
System.out.println("========================================");
System.out.println("线程池隔离演示");
System.out.println("========================================");
System.out.println();
System.out.println("方案说明:");
System.out.println(" 1. 为不同业务配置独立的线程池");
System.out.println(" 2. 核心业务使用专用线程池");
System.out.println(" 3. 非核心业务使用独立线程池");
System.out.println(" 4. 实现业务隔离,互不影响");
System.out.println();
IsolatedFastService fastService = context.getBean(IsolatedFastService.class);
IsolatedSlowService slowService = context.getBean(IsolatedSlowService.class);
int fastTaskCount = 10;
int slowTaskCount = 5;
CountDownLatch fastLatch = new CountDownLatch(fastTaskCount);
CountDownLatch slowLatch = new CountDownLatch(slowTaskCount);
System.out.println("步骤1: 同时提交快速任务和慢任务到隔离线程池");
System.out.println(" 快速任务线程池大小: 10");
System.out.println(" 慢任务线程池大小: 3");
System.out.println();
long startTime = System.currentTimeMillis();
for (int i = 0; i < slowTaskCount; i++) {
slowService.executeSlowTask(() -> {
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
slowLatch.countDown();
});
}
for (int i = 0; i < fastTaskCount; i++) {
final int taskId = i;
fastService.executeFastTask(() -> {
long waitTime = System.currentTimeMillis() - startTime;
System.out.println(" 快速任务" + taskId + " 开始执行,等待时间: " + waitTime + "ms");
fastLatch.countDown();
});
}
fastLatch.await();
long fastDuration = System.currentTimeMillis() - startTime;
System.out.println();
System.out.println("步骤2: 快速任务执行完成");
System.out.println(" 快速任务执行时间: " + fastDuration + "ms");
System.out.println();
slowLatch.await();
long totalDuration = System.currentTimeMillis() - startTime;
System.out.println("步骤3: 慢任务执行完成");
System.out.println(" 总执行时间: " + totalDuration + "ms");
System.out.println();
System.out.println("分析:");
System.out.println(" ✅ 快速任务没有被慢任务阻塞");
System.out.println(" ✅ 业务之间相互隔离");
System.out.println(" ✅ 核心业务响应时间可控");
System.out.println();
System.out.println("优点:");
System.out.println(" - 业务隔离,互不影响");
System.out.println(" - 核心业务稳定性有保障");
System.out.println(" - 可以针对不同业务优化配置");
System.out.println();
System.out.println("缺点:");
System.out.println(" - 配置复杂度增加");
System.out.println(" - 资源占用可能增加");
context.close();
}
private static void demonstrateCompleteThreadPoolConfig() throws InterruptedException {
AnnotationConfigApplicationContext context = new AnnotationConfigApplicationContext(CompletePoolConfig.class);
System.out.println("========================================");
System.out.println("完善的线程池配置演示");
System.out.println("========================================");
System.out.println();
System.out.println("方案说明:");
System.out.println(" 1. 合理配置核心参数");
System.out.println(" 2. 配置拒绝策略和异常处理");
System.out.println(" 3. 配置线程池监控");
System.out.println(" 4. 优雅关闭");
System.out.println();
CompleteAsyncService service = context.getBean(CompleteAsyncService.class);
int taskCount = 30;
AtomicInteger successCount = new AtomicInteger(0);
AtomicInteger rejectedCount = new AtomicInteger(0);
CountDownLatch latch = new CountDownLatch(taskCount);
System.out.println("步骤1: 提交" + taskCount + "个任务");
System.out.println(" 线程池核心大小: 5");
System.out.println(" 线程池最大大小: 10");
System.out.println(" 队列容量: 20");
System.out.println(" 拒绝策略: CallerRunsPolicy(调用者运行)");
System.out.println();
for (int i = 0; i < taskCount; i++) {
final int taskId = i;
try {
service.executeTask(() -> {
successCount.incrementAndGet();
String threadName = Thread.currentThread().getName();
boolean isCallerThread = !threadName.startsWith("complete-");
if (isCallerThread) {
System.out.println(" 任务" + taskId + " 由调用者线程执行: " + threadName);
}
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
latch.countDown();
});
} catch (Exception e) {
rejectedCount.incrementAndGet();
latch.countDown();
}
}
latch.await();
System.out.println();
System.out.println("步骤2: 分析结果");
System.out.println(" 成功执行: " + successCount.get());
System.out.println(" 被拒绝: " + rejectedCount.get());
System.out.println();
ThreadPoolTaskExecutor executor = (ThreadPoolTaskExecutor) context.getBean("completeExecutor");
System.out.println("步骤3: 线程池状态");
System.out.println(" 活跃线程数: " + executor.getActiveCount());
System.out.println(" 核心线程数: " + executor.getCorePoolSize());
System.out.println(" 最大线程数: " + executor.getMaxPoolSize());
System.out.println(" 队列大小: " + executor.getQueueCapacity());
System.out.println();
System.out.println("分析:");
System.out.println(" ✅ 使用CallerRunsPolicy,任务不会被丢弃");
System.out.println(" ✅ 超出容量时由调用者线程执行,实现背压");
System.out.println(" ✅ 可以监控线程池状态");
System.out.println();
System.out.println("优点:");
System.out.println(" - 任务不会丢失");
System.out.println(" - 自动实现背压控制");
System.out.println(" - 可监控可管理");
System.out.println();
System.out.println("缺点:");
System.out.println(" - 调用者线程可能被阻塞");
System.out.println(" - 需要根据业务调整配置");
context.close();
}
}
@Configuration
@EnableAsync
class CustomPoolConfig {
@Bean(name = "customExecutor")
public Executor customExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(100);
executor.setThreadNamePrefix("custom-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(60);
executor.initialize();
return executor;
}
@Bean
public CustomAsyncService customAsyncService() {
return new CustomAsyncService();
}
}
@Service
class CustomAsyncService {
@Async("customExecutor")
public void executeAsync(Runnable task) {
task.run();
}
}
@Configuration
@EnableAsync
class IsolatedPoolConfig {
@Bean(name = "fastTaskExecutor")
public Executor fastTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(10);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(50);
executor.setThreadNamePrefix("fast-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
@Bean(name = "slowTaskExecutor")
public Executor slowTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(3);
executor.setMaxPoolSize(3);
executor.setQueueCapacity(20);
executor.setThreadNamePrefix("slow-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
return executor;
}
@Bean
public IsolatedFastService isolatedFastService() {
return new IsolatedFastService();
}
@Bean
public IsolatedSlowService isolatedSlowService() {
return new IsolatedSlowService();
}
}
@Service
class IsolatedFastService {
@Async("fastTaskExecutor")
public void executeFastTask(Runnable task) {
task.run();
}
}
@Service
class IsolatedSlowService {
@Async("slowTaskExecutor")
public void executeSlowTask(Runnable task) {
task.run();
}
}
@Configuration
@EnableAsync
class CompletePoolConfig implements AsyncUncaughtExceptionHandler {
@Bean(name = "completeExecutor")
public Executor completeExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(5);
executor.setMaxPoolSize(10);
executor.setQueueCapacity(20);
executor.setThreadNamePrefix("complete-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(60);
executor.setAllowCoreThreadTimeOut(true);
executor.setKeepAliveSeconds(60);
executor.initialize();
return executor;
}
@Bean
public CompleteAsyncService completeAsyncService() {
return new CompleteAsyncService();
}
@Override
public void handleUncaughtException(Throwable ex, Method method, Object... params) {
System.out.println(" 异步任务异常: " + ex.getMessage());
System.out.println(" 方法: " + method.getName());
}
}
@Service
class CompleteAsyncService {
@Async("completeExecutor")
public void executeTask(Runnable task) {
task.run();
}
}
运行结果:
=== 方案一:自定义线程池解决方案 === ======================================== 自定义线程池演示 ======================================== 方案说明: 1. 配置自定义ThreadPoolTaskExecutor 2. 设置合理的核心线程数、最大线程数、队列容量 3. 配置合理的拒绝策略 步骤1: 并发调用20个异步任务 任务0 执行线程: custom-1 任务1 执行线程: custom-2 任务2 执行线程: custom-3 任务3 执行线程: custom-4 任务4 执行线程: custom-5 任务5 执行线程: custom-1 ← 线程复用 任务6 执行线程: custom-2 ← 线程复用 ... 步骤2: 分析结果 执行时间: 450ms 创建的不同线程数: 5 总任务数: 20 分析: ✅ 线程被复用,不是每次创建新线程 ✅ 线程数量可控,不会无限增长 ✅ 性能更好,资源使用更合理 优点: - 线程复用,减少创建销毁开销 - 线程数量可控,避免OOM - 可以监控线程池状态 缺点: - 需要手动配置线程池参数 - 需要根据业务场景调整配置 === 方案二:线程池隔离解决方案 === ======================================== 线程池隔离演示 ======================================== 方案说明: 1. 为不同业务配置独立的线程池 2. 核心业务使用专用线程池 3. 非核心业务使用独立线程池 4. 实现业务隔离,互不影响 步骤1: 同时提交快速任务和慢任务到隔离线程池 快速任务线程池大小: 10 慢任务线程池大小: 3 快速任务0 开始执行,等待时间: 5ms 快速任务1 开始执行,等待时间: 6ms ... 快速任务9 开始执行,等待时间: 10ms ← 全部快速执行 步骤2: 快速任务执行完成 快速任务执行时间: 15ms 步骤3: 慢任务执行完成 总执行时间: 2015ms 分析: ✅ 快速任务没有被慢任务阻塞 ✅ 业务之间相互隔离 ✅ 核心业务响应时间可控 优点: - 业务隔离,互不影响 - 核心业务稳定性有保障 - 可以针对不同业务优化配置 缺点: - 配置复杂度增加 - 资源占用可能增加 === 方案三:完善的线程池配置 === ======================================== 完善的线程池配置演示 ======================================== 方案说明: 1. 合理配置核心参数 2. 配置拒绝策略和异常处理 3. 配置线程池监控 4. 优雅关闭 步骤1: 提交30个任务 线程池核心大小: 5 线程池最大大小: 10 队列容量: 20 拒绝策略: CallerRunsPolicy(调用者运行) 任务25 由调用者线程执行: main 任务26 由调用者线程执行: main ... 步骤2: 分析结果 成功执行: 30 被拒绝: 0 步骤3: 线程池状态 活跃线程数: 5 核心线程数: 5 最大线程数: 10 队列大小: 20 分析: ✅ 使用CallerRunsPolicy,任务不会被丢弃 ✅ 超出容量时由调用者线程执行,实现背压 ✅ 可以监控线程池状态 优点: - 任务不会丢失 - 自动实现背压控制 - 可监控可管理 缺点: - 调用者线程可能被阻塞 - 需要根据业务调整配置
4. 架构思考
4.1 线程池参数配置建议
| 参数 | 建议值 | 说明 |
|---|---|---|
| corePoolSize | CPU核心数 ~ CPU核心数*2 | CPU密集型取小,IO密集型取大 |
| maxPoolSize | corePoolSize * 2 | 根据业务峰值调整 |
| queueCapacity | 100-500 | 过大会导致响应延迟 |
| keepAliveSeconds | 60 | 空闲线程存活时间 |
| rejectedExecutionHandler | CallerRunsPolicy | 推荐使用调用者运行策略 |
4.2 拒绝策略选择
| 策略 | 行为 | 适用场景 |
|---|---|---|
| AbortPolicy | 抛出异常 | 需要感知任务失败的场景 |
| CallerRunsPolicy | 调用者线程执行 | 不允许丢失任务的场景 |
| DiscardPolicy | 直接丢弃 | 允许丢失任务的场景 |
| DiscardOldestPolicy | 丢弃最老任务 | 允许丢失旧任务的场景 |
4.3 最佳实践总结
代码层面:
- ✅ 必须配置自定义线程池
- ✅ 为不同业务配置独立线程池
- ✅ 配置合理的拒绝策略
- ✅ 实现AsyncUncaughtExceptionHandler处理异常
- ❌ 不要使用默认的SimpleAsyncTaskExecutor
- ❌ 不要让多个业务共用一个线程池
团队规范:
- 强制规范:所有@Async必须指定线程池
- 代码审查:重点检查线程池配置
- 监控告警:监控线程池状态和队列积压
- 文档说明:在代码注释中说明线程池配置原因
架构设计:
- 线程池隔离:核心业务与非核心业务隔离
- 监控体系:接入Prometheus等监控系统
- 容量规划:根据业务峰值规划线程池容量
- 应急预案:准备线程池耗尽时的降级方案
4.4 监控指标
建议监控以下线程池指标:
1. 活跃线程数(activeCount)
2. 核心线程数(corePoolSize)
3. 最大线程数(maxPoolSize)
4. 队列大小(queueSize)
5. 已完成任务数(completedTaskCount)
6. 拒绝任务数(rejectedTaskCount)
通过深入理解@Async注解的线程池机制,不仅能避免生产环境的OOM和性能问题,更能提升对异步编程和线程池管理的整体思考。在实际项目中,正确配置线程池至关重要,唯有深入理解底层原理,才能构建真正稳定可靠的异步处理系统。
到此这篇关于Java中@Async注解的线程池隔离陷阱的文章就介绍到这了,更多相关Java @Async线程池隔离内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!
