Java异步编程案例

wen java案例 4

本文目录导读:

Java异步编程案例

  1. CompletableFuture 异步编程案例
  2. Spring Boot 异步编程案例
  3. 使用 ExecutorService 的传统异步编程
  4. 生产环境最佳实践示例

我将为您提供几个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();
    }
}
  1. 选择正确的异步方式

    • CompletableFuture:适合复杂的异步流程编排
    • ExecutorService:适合传统任务提交和批量执行
    • Spring @Async:适合Spring Boot应用中的简单异步方法
  2. 线程池配置建议

    • 核心线程数:根据任务类型(CPU密集型 vs IO密集型)决定
    • 队列容量:防止内存溢出
    • 拒绝策略:确保任务不会丢失
  3. 注意事项

    • 处理线程池耗尽问题
    • 注意线程安全
    • 合理设置超时时间
    • 做好异常处理和降级机制
    • 监控线程池状态和任务执行情况

这些案例覆盖了Java异步编程的主要场景,在实际项目中可以根据需求选择合适的方式来异步处理任务。

抱歉,评论功能暂时关闭!