本文目录导读:

我来给你一个完整的Spring Boot异步任务案例,包括基础的异步执行、异步回调处理,以及最佳实践。
项目依赖配置
pom.xml
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>
配置异步任务
异步配置类
package com.example.async.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import java.util.concurrent.Executor;
import java.util.concurrent.ThreadPoolExecutor;
@Configuration
public class AsyncConfig {
@Bean("taskExecutor")
public Executor taskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
// 核心线程数
executor.setCorePoolSize(5);
// 最大线程数
executor.setMaxPoolSize(10);
// 队列容量
executor.setQueueCapacity(100);
// 线程名称前缀
executor.setThreadNamePrefix("async-task-");
// 拒绝策略:调用者执行
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
// 初始化
executor.initialize();
return executor;
}
// 可以配置多个线程池处理不同业务
@Bean("emailExecutor")
public Executor emailExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(3);
executor.setMaxPoolSize(6);
executor.setQueueCapacity(50);
executor.setThreadNamePrefix("email-task-");
executor.initialize();
return executor;
}
}
或者使用配置文件方式:
# application.yml
spring:
task:
execution:
pool:
core-size: 5
max-size: 10
queue-capacity: 100
keep-alive: 60s
thread-name-prefix: async-task-
shutdown:
await-termination: true
await-termination-period: 60s
异步任务实现
基础异步任务类
package com.example.async.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
@Service
public class AsyncTaskService {
private static final Logger logger = LoggerFactory.getLogger(AsyncTaskService.class);
/**
* 无返回值异步任务
*/
@Async("taskExecutor")
public void sendEmail(String to, String content) {
logger.info("开始发送邮件到: {}", to);
try {
// 模拟耗时操作
Thread.sleep(2000);
logger.info("邮件发送成功: {}", to);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
logger.error("邮件发送失败: {}", to, e);
}
}
/**
* 带返回值的异步任务
*/
@Async("taskExecutor")
public CompletableFuture<String> generateReport(String reportType) {
logger.info("开始生成报表: {}", reportType);
try {
Thread.sleep(3000);
String result = "报表-" + reportType + "-" + System.currentTimeMillis();
logger.info("报表生成完成: {}", result);
return CompletableFuture.completedFuture(result);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return CompletableFuture.failedFuture(e);
}
}
/**
* 处理订单异步任务
*/
@Async("taskExecutor")
public void processOrder(Long orderId) {
logger.info("开始处理订单: {}", orderId);
// 订单处理逻辑
try {
Thread.sleep(1500);
// 发送通知
sendNotification(orderId);
// 更新库存
updateStock(orderId);
logger.info("订单处理完成: {}", orderId);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
logger.error("订单处理失败: {}", orderId, e);
}
}
private void sendNotification(Long orderId) {
logger.info("发送订单通知: {}", orderId);
}
private void updateStock(Long orderId) {
logger.info("更新库存: {}", orderId);
}
}
异步事件处理
package com.example.async.event;
import org.springframework.context.event.EventListener;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
@Component
public class OrderEventListener {
private static final Logger logger = LoggerFactory.getLogger(OrderEventListener.class);
@Async("taskExecutor")
@EventListener
public void handleOrderCreatedEvent(OrderCreatedEvent event) {
logger.info("处理订单创建事件: {}", event.getOrderId());
// 发送订单创建通知邮件
try {
Thread.sleep(2000);
logger.info("订单通知已发送: {}", event.getOrderId());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
@Async("taskExecutor")
@EventListener
public void handlePaymentSuccessEvent(PaymentSuccessEvent event) {
logger.info("处理支付成功事件: 订单{} 金额{}", event.getOrderId(), event.getAmount());
// 更新订单状态等
}
}
// 事件类
public class OrderCreatedEvent {
private Long orderId;
// getter/setter...
}
public class PaymentSuccessEvent {
private Long orderId;
private BigDecimal amount;
// getter/setter...
}
控制器层
package com.example.async.controller;
import com.example.async.service.AsyncTaskService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.multipart.MultipartFile;
import java.util.concurrent.CompletableFuture;
@RestController
@RequestMapping("/api/async")
public class AsyncController {
@Autowired
private AsyncTaskService asyncTaskService;
/**
* 触发异步任务
*/
@PostMapping("/email")
public ApiResponse sendEmail(@RequestParam String to, @RequestParam String content) {
asyncTaskService.sendEmail(to, content);
return ApiResponse.success("邮件发送任务已提交");
}
/**
* 获取异步任务结果
*/
@GetMapping("/report/{type}")
public CompletableFuture<ApiResponse> generateReport(@PathVariable String type) {
return asyncTaskService.generateReport(type)
.thenApply(result -> ApiResponse.success(result));
}
/**
* 批量异步处理
*/
@PostMapping("/batch-process")
public ApiResponse batchProcess(@RequestParam("files") MultipartFile[] files) {
List<CompletableFuture<String>> futures = new ArrayList<>();
for (MultipartFile file : files) {
CompletableFuture<String> future = asyncTaskService.generateReport(file.getOriginalFilename());
futures.add(future);
}
// 等待所有任务完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
.thenRun(() -> System.out.println("所有文件处理完成"));
return ApiResponse.success("批量处理任务已提交,共" + files.length + "个文件");
}
}
异步任务管理器(进阶用法)
package com.example.async.manager;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CompletableFuture;
import java.util.function.Function;
@Component
public class AsyncTaskManager {
private final Map<String, CompletableFuture<?>> tasks = new ConcurrentHashMap<>();
public void submitTask(String taskId, Runnable task) {
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
task.run();
});
tasks.put(taskId, future);
}
public <T> CompletableFuture<T> submitTask(String taskId, Supplier<T> task) {
CompletableFuture<T> future = CompletableFuture.supplyAsync(task);
tasks.put(taskId, future);
return future;
}
public CompletableFuture<?> getTaskResult(String taskId) {
return tasks.get(taskId);
}
public boolean isCompleted(String taskId) {
CompletableFuture<?> future = tasks.get(taskId);
return future != null && future.isDone();
}
public void cancelTask(String taskId) {
CompletableFuture<?> future = tasks.get(taskId);
if (future != null) {
future.cancel(true);
}
}
}
异步回调处理
package com.example.async.callback;
import org.springframework.stereotype.Component;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@Component
public class AsyncCallbackHandler {
/**
* 演示异步回调处理
*/
public CompletableFuture<User> getUserInfoAsync(Long userId) {
return CompletableFuture.supplyAsync(() -> {
// 模拟异步获取用户信息
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
User user = new User();
user.setId(userId);
user.setName("User" + userId);
return user;
}).thenApply(user -> {
// 异步结果处理
user.setStatus("PROCESSED");
return user;
}).exceptionally(ex -> {
System.err.println("获取用户信息失败: " + ex.getMessage());
return null;
});
}
/**
* 多个异步任务的组合
*/
public CompletableFuture<OrderInfo> getOrderDetails(Long orderId) {
CompletableFuture<User> userFuture = getUserInfoAsync(1L);
CompletableFuture<Order> orderFuture = CompletableFuture.supplyAsync(() -> {
// 模拟异步获取订单
return new Order(orderId);
});
// 组合多个异步任务的结果
return userFuture.thenCombine(orderFuture, (user, order) -> {
OrderInfo info = new OrderInfo();
info.setUser(user);
info.setOrder(order);
return info;
});
}
}
测试代码
package com.example.async;
import com.example.async.service.AsyncTaskService;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@SpringBootTest
public class AsyncTaskTest {
@Autowired
private AsyncTaskService asyncTaskService;
@Test
public void testAsyncTask() throws Exception {
// 测试无返回值异步任务
asyncTaskService.sendEmail("test@example.com", "测试邮件");
System.out.println("主线程继续执行...");
// 等待异步任务完成
Thread.sleep(5000);
// 测试带返回值异步任务
CompletableFuture<String> future = asyncTaskService.generateReport("SALES");
String result = future.get(10, TimeUnit.SECONDS);
System.out.println("报表结果: " + result);
}
@Test
public void testConcurrentAsyncTasks() throws Exception {
List<CompletableFuture<String>> futures = new ArrayList<>();
// 同时触发多个异步任务
for (int i = 0; i < 5; i++) {
futures.add(asyncTaskService.generateReport("报表" + i));
}
// 等待所有任务完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
.join();
// 获取所有结果
for (CompletableFuture<String> future : futures) {
System.out.println("结果: " + future.get());
}
}
}
注意事项
常见问题1:同类调用失效
@Service
public class AsyncService {
public void methodA() {
// 这样调用不会触发异步,因为代理不会处理同类调用
this.methodB();
// 正确方式:注入自身或使用ApplicationContext获取代理
}
@Async
public void methodB() {
// 异步逻辑
}
// 正确方式示例
@Autowired
private ApplicationContext context;
public void methodC() {
AsyncService proxy = context.getBean(AsyncService.class);
proxy.methodB(); // 这样才会触发异步
}
}
自定义异常处理
package com.example.async.config;
import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.AsyncConfigurer;
import java.lang.reflect.Method;
import java.util.Arrays;
@Configuration
public class AsyncExceptionConfig implements AsyncConfigurer {
@Override
public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
return new AsyncUncaughtExceptionHandler() {
@Override
public void handleUncaughtException(Throwable ex, Method method, Object... params) {
System.err.println("异步任务异常:");
System.err.println("方法: " + method.getName());
System.err.println("参数: " + Arrays.toString(params));
System.err.println("异常信息: " + ex.getMessage());
// 可以在这里发送告警通知等
}
};
}
}
这个案例涵盖了Spring Boot异步任务的核心用法,包括:
- 基础异步执行:@Async注解的基本使用
- 带返回值的异步任务:使用CompletableFuture
- 异步事件处理:结合Spring事件机制
- 异步回调处理:链式调用和组合
- 最佳实践:线程池配置、异常处理、性能优化
根据你的具体业务需求,可以选择合适的方式进行实现。