本文目录导读:

我将为您提供几个Java异步编程的实践案例,涵盖不同的实现方式。
CompletableFuture 异步编程案例
案例:异步获取用户信息和订单信息
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
public class CompletableFutureExample {
public static void main(String[] args) throws ExecutionException, InterruptedException {
// 案例1:基础异步任务
basicAsyncExample();
// 案例2:异步任务组合
combinedAsyncExample();
// 案例3:异步任务异常处理
exceptionHandlingExample();
// 案例4:并行执行多个异步任务
parallelAsyncExample();
// 案例5:超时控制
timeoutExample();
}
// 基础异步任务示例
public static void basicAsyncExample() throws ExecutionException, InterruptedException {
System.out.println("=== 基础异步任务示例 ===");
long startTime = System.currentTimeMillis();
// 创建异步任务
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return "异步任务结果";
});
// 主线程不阻塞,可以做其他事情
System.out.println("主线程继续执行其他操作...");
// 获取异步结果(阻塞等待)
String result = future.get();
System.out.println("异步任务返回: " + result);
System.out.println("总耗时: " + (System.currentTimeMillis() - startTime) + "ms");
}
// 异步任务组合示例
public static void combinedAsyncExample() throws ExecutionException, InterruptedException {
System.out.println("\n=== 异步任务组合示例 ===");
// 模拟获取用户信息
CompletableFuture<User> userFuture = CompletableFuture.supplyAsync(() -> {
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return new User(1, "张三", "zhangsan@example.com");
});
// 用户信息获取完成后,获取订单信息
CompletableFuture<Order> orderFuture = userFuture.thenCompose(user -> {
System.out.println("已获取用户: " + user.getName());
// 根据用户ID异步获取订单
return CompletableFuture.supplyAsync(() -> {
try {
Thread.sleep(1500);
} catch (InterruptedException e) {
e.printStackTrace();
}
return new Order(1001, user.getId(), 299.99);
});
});
// 等待最终结果
Order order = orderFuture.get();
System.out.println("最终订单结果: " + order);
}
// 异常处理示例
public static void exceptionHandlingExample() {
System.out.println("\n=== 异步任务异常处理示例 ===");
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
// 模拟异常
if (true) {
throw new RuntimeException("模拟业务异常");
}
return "正常结果";
}).exceptionally(ex -> {
System.out.println("捕获异常: " + ex.getMessage());
return "降级处理结果";
}).handle((result, ex) -> {
if (ex != null) {
System.out.println("handle方法捕获异常: " + ex.getMessage());
return "handle异常处理";
}
return result + " 处理完成";
});
future.thenAccept(result -> {
System.out.println("最终结果: " + result);
});
}
// 并行执行多个异步任务
public static void parallelAsyncExample() throws ExecutionException, InterruptedException {
System.out.println("\n=== 并行执行异步任务示例 ===");
long startTime = System.currentTimeMillis();
// 创建3个并行任务
CompletableFuture<String> task1 = CompletableFuture.supplyAsync(() -> {
sleep(2000);
return "任务1完成";
});
CompletableFuture<String> task2 = CompletableFuture.supplyAsync(() -> {
sleep(1500);
return "任务2完成";
});
CompletableFuture<String> task3 = CompletableFuture.supplyAsync(() -> {
sleep(3000);
return "任务3完成";
});
// 等待所有任务完成
CompletableFuture<Void> allTasks = CompletableFuture.allOf(task1, task2, task3);
allTasks.get();
System.out.println("所有任务完成,耗时: " + (System.currentTimeMillis() - startTime) + "ms");
System.out.println("任务1: " + task1.get());
System.out.println("任务2: " + task2.get());
System.out.println("任务3: " + task3.get());
// 等待任意一个任务完成
CompletableFuture<Object> anyTask = CompletableFuture.anyOf(task1, task2, task3);
System.out.println("\n最快完成的任务: " + anyTask.get());
}
// 超时控制示例
public static void timeoutExample() throws ExecutionException, InterruptedException {
System.out.println("\n=== 异步任务超时控制示例 ===");
CompletableFuture<String> slowFuture = CompletableFuture.supplyAsync(() -> {
sleep(3000);
return "慢任务返回结果";
});
try {
// 设置2秒超时
String result = slowFuture.orTimeout(2, TimeUnit.SECONDS).join();
System.out.println("成功获取结果: " + result);
} catch (Exception e) {
System.out.println("任务超时: " + e.getMessage());
}
// 另一种方式:completeOnTimeout(固定返回默认值)
CompletableFuture<String> futureWithDefault = CompletableFuture
.supplyAsync(() -> {
sleep(3000);
return "正常结果";
})
.completeOnTimeout("默认结果", 1, TimeUnit.SECONDS);
System.out.println("超时返回默认值: " + futureWithDefault.get());
}
private static void sleep(int ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// 辅助类
static class User {
private Integer id;
private String name;
private String email;
public User(Integer id, String name, String email) {
this.id = id;
this.name = name;
this.email = email;
}
public Integer getId() { return id; }
public String getName() { return name; }
@Override
public String toString() {
return "User{id=" + id + ", name='" + name + "'}";
}
}
static class Order {
private Integer orderId;
private Integer userId;
private Double amount;
public Order(Integer orderId, Integer userId, Double amount) {
this.orderId = orderId;
this.userId = userId;
this.amount = amount;
}
@Override
public String toString() {
return "Order{orderId=" + orderId + ", userId=" + userId + ", amount=" + amount + "}";
}
}
}
Spring Boot 异步编程案例
定义异步任务
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import org.springframework.beans.factory.annotation.Autowired;
import java.util.concurrent.CompletableFuture;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@Service
public class AsyncDemoService {
private static final Logger logger = LoggerFactory.getLogger(AsyncDemoService.class);
@Autowired
private EmailService emailService;
// 简单异步方法
@Async("taskExecutor")
public void asyncSendEmail(String recipient, String content) {
logger.info("开始异步发送邮件给: {}", recipient);
try {
// 模拟发送邮件
Thread.sleep(2000);
logger.info("邮件发送成功: {}", recipient);
} catch (InterruptedException e) {
logger.error("邮件发送失败", e);
Thread.currentThread().interrupt();
}
}
// 返回CompletableFuture的异步方法
@Async("taskExecutor")
public CompletableFuture<String> asyncProcessData(String data) {
logger.info("开始异步处理数据: {}", data);
try {
Thread.sleep(3000);
String result = "处理完成: " + data.toUpperCase();
logger.info("数据处理完成");
return CompletableFuture.completedFuture(result);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return CompletableFuture.failedFuture(e);
}
}
// 多个异步任务并行处理
@Async
public CompletableFuture<Order> asyncCreateOrder(Product product) {
logger.info("异步创建订单,产品: {}", product.getName());
try {
Thread.sleep(2000);
// 模拟创建订单
Order order = new Order();
order.setOrderId(generateOrderId());
order.setProductName(product.getName());
order.setPrice(product.getPrice());
order.setStatus("CREATED");
// 异步发送订单确认邮件
emailService.sendOrderConfirmation(order);
return CompletableFuture.completedFuture(order);
} catch (Exception e) {
return CompletableFuture.failedFuture(e);
}
}
private String generateOrderId() {
return "ORDER-" + System.currentTimeMillis();
}
}
自定义线程池配置
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();
// 核心线程数:CPU核心数
executor.setCorePoolSize(Runtime.getRuntime().availableProcessors());
// 最大线程数:CPU核心数的2倍
executor.setMaxPoolSize(Runtime.getRuntime().availableProcessors() * 2);
// 队列容量
executor.setQueueCapacity(100);
// 线程名前缀
executor.setThreadNamePrefix("async-task-");
// 拒绝策略:由调用线程处理
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
// 优雅关闭
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(60);
executor.initialize();
return executor;
}
}
异步事件监听
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);
// 执行后续操作:发送通知、更新统计等
notifyCustomer(event.getCustomerEmail(), event.getOrderId());
updateSalesStatistics(event.getOrderAmount());
logger.info("订单事件处理完成: {}", event.getOrderId());
} catch (Exception e) {
logger.error("订单事件处理失败: {}", event.getOrderId(), e);
}
}
private void notifyCustomer(String email, String orderId) {
logger.info("向客户 {} 发送订单确认邮件,订单号: {}", email, orderId);
}
private void updateSalesStatistics(Double amount) {
logger.info("更新销售统计,金额: {}", amount);
}
}
// 定义事件类
public class OrderCreatedEvent {
private String orderId;
private String customerEmail;
private Double orderAmount;
public OrderCreatedEvent(String orderId, String customerEmail, Double orderAmount) {
this.orderId = orderId;
this.customerEmail = customerEmail;
this.orderAmount = orderAmount;
}
public String getOrderId() { return orderId; }
public String getCustomerEmail() { return customerEmail; }
public Double getOrderAmount() { return orderAmount; }
}
使用 ExecutorService 的传统异步编程
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
public class ExecutorServiceExample {
public static void main(String[] args) throws Exception {
// 创建线程池
ExecutorService executorService = Executors.newFixedThreadPool(5);
// 案例1:单任务执行
submitSingleTask(executorService);
// 案例2:批量任务执行
invokeAllTasks(executorService);
// 案例3:定时任务
scheduleTasks();
// 关闭线程池
executorService.shutdown();
}
private static void submitSingleTask(ExecutorService executorService) throws Exception {
System.out.println("=== 单任务示例 ===");
Future<String> future = executorService.submit(() -> {
Thread.sleep(2000);
return "任务执行结果";
});
System.out.println("主线程继续工作...");
String result = future.get(3, TimeUnit.SECONDS); // 3秒超时
System.out.println("任务结果: " + result);
}
private static void invokeAllTasks(ExecutorService executorService) throws Exception {
System.out.println("\n=== 批量任务示例 ===");
List<Callable<String>> tasks = new ArrayList<>();
tasks.add(() -> {
Thread.sleep(1000);
return "任务1完成";
});
tasks.add(() -> {
Thread.sleep(2000);
return "任务2完成";
});
tasks.add(() -> {
Thread.sleep(3000);
return "任务3完成";
});
long startTime = System.currentTimeMillis();
// 执行所有任务并等待全部完成
List<Future<String>> futures = executorService.invokeAll(tasks);
for (Future<String> future : futures) {
System.out.println("结果: " + future.get());
}
System.out.println("所有任务完成,耗时: " + (System.currentTimeMillis() - startTime) + "ms");
}
private static void scheduleTasks() throws Exception {
System.out.println("\n=== 定时任务示例 ===");
ScheduledExecutorService scheduledExecutor = Executors.newScheduledThreadPool(2);
// 延迟执行
scheduledExecutor.schedule(() -> {
System.out.println("延迟2秒执行的任务");
}, 2, TimeUnit.SECONDS);
// 固定频率执行(会追任务)
scheduledExecutor.scheduleAtFixedRate(() -> {
System.out.println("每3秒执行一次");
}, 0, 3, TimeUnit.SECONDS);
// 固定延迟执行(不追任务)
scheduledExecutor.scheduleWithFixedDelay(() -> {
System.out.println("上次任务完成后延迟1秒执行");
}, 0, 1, TimeUnit.SECONDS);
// 等待一段时间后关闭
Thread.sleep(8000);
scheduledExecutor.shutdown();
}
}
生产环境最佳实践示例
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
/**
* 生产环境异步处理框架
*/
public class AsyncProcessingFramework {
// 处理器接口
public interface AsyncTask<T> {
T process() throws Exception;
}
// 异步处理器
public static class AsyncProcessor {
private final ExecutorService executorService;
private final CountDownLatch latch;
private final AtomicInteger successCount = new AtomicInteger(0);
private final AtomicInteger failureCount = new AtomicInteger(0);
public AsyncProcessor(int threadCount) {
this.executorService = Executors.newFixedThreadPool(threadCount,
new ThreadFactory() {
private final AtomicInteger counter = new AtomicInteger();
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r);
t.setName("async-processor-" + counter.incrementAndGet());
t.setDaemon(true);
return t;
}
});
this.latch = new CountDownLatch(threadCount);
}
// 执行异步任务
public CompletableFuture<ProcessResult> executeAsync(AsyncTask<?> task, String taskName) {
return CompletableFuture.supplyAsync(() -> {
long startTime = System.currentTimeMillis();
try {
Object result = task.process();
successCount.incrementAndGet();
return ProcessResult.success(taskName, result,
System.currentTimeMillis() - startTime);
} catch (Exception e) {
failureCount.incrementAndGet();
return ProcessResult.failure(taskName, e,
System.currentTimeMillis() - startTime);
}
}, executorService);
}
// 批量执行任务
public List<ProcessResult> executeBatch(List<AsyncTask<?>> tasks) {
List<CompletableFuture<ProcessResult>> futures = new ArrayList<>();
for (int i = 0; i < tasks.size(); i++) {
CompletableFuture<ProcessResult> future =
executeAsync(tasks.get(i), "Task-" + (i + 1));
futures.add(future);
}
// 等待所有任务完成
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
// 收集结果
List<ProcessResult> results = new ArrayList<>();
for (CompletableFuture<ProcessResult> future : futures) {
try {
results.add(future.get());
} catch (Exception e) {
results.add(ProcessResult.failure("Unknown Task", e, 0));
}
}
return results;
}
public void shutdown() {
executorService.shutdown();
try {
if (!executorService.awaitTermination(30, TimeUnit.SECONDS)) {
executorService.shutdownNow();
}
} catch (InterruptedException e) {
executorService.shutdownNow();
Thread.currentThread().interrupt();
}
}
public int getSuccessCount() { return successCount.get(); }
public int getFailureCount() { return failureCount.get(); }
}
// 处理结果类
public static class ProcessResult {
private final String taskName;
private final boolean success;
private final Object data;
private final long duration;
private final Throwable error;
private ProcessResult(String taskName, boolean success, Object data,
long duration, Throwable error) {
this.taskName = taskName;
this.success = success;
this.data = data;
this.duration = duration;
this.error = error;
}
public static ProcessResult success(String taskName, Object data, long duration) {
return new ProcessResult(taskName, true, data, duration, null);
}
public static ProcessResult failure(String taskName, Throwable error, long duration) {
return new ProcessResult(taskName, false, null, duration, error);
}
// getters
public String getTaskName() { return taskName; }
public boolean isSuccess() { return success; }
public Object getData() { return data; }
public long getDuration() { return duration; }
public Throwable getError() { return error; }
@Override
public String toString() {
return String.format("ProcessResult{task='%s', success=%s, duration=%dms, error=%s}",
taskName, success, duration, error != null ? error.getMessage() : "null");
}
}
// 测试代码
public static void main(String[] args) {
AsyncProcessor processor = new AsyncProcessor(4);
System.out.println("=== 测试异步处理器框架 ===");
// 模拟不同任务
List<AsyncTask<?>> tasks = new ArrayList<>();
for (int i = 1; i <= 10; i++) {
final int taskId = i;
tasks.add(() -> {
Thread.sleep(1000 + (int)(Math.random() * 2000));
if (taskId % 5 == 0) {
throw new RuntimeException("Task " + taskId + " 模拟失败");
}
return "数据-" + taskId;
});
}
List<ProcessResult> results = processor.executeBatch(tasks);
results.forEach(System.out::println);
System.out.println("\n任务统计 - 成功: " + processor.getSuccessCount() +
", 失败: " + processor.getFailureCount());
processor.shutdown();
}
}
-
选择正确的异步方式
CompletableFuture:适合复杂的异步流程编排ExecutorService:适合传统任务提交和批量执行- Spring
@Async:适合Spring Boot应用中的简单异步方法
-
线程池配置建议
- 核心线程数:根据任务类型(CPU密集型 vs IO密集型)决定
- 队列容量:防止内存溢出
- 拒绝策略:确保任务不会丢失
-
注意事项
- 处理线程池耗尽问题
- 注意线程安全
- 合理设置超时时间
- 做好异常处理和降级机制
- 监控线程池状态和任务执行情况
这些案例覆盖了Java异步编程的主要场景,在实际项目中可以根据需求选择合适的方式来异步处理任务。