Java CompletableFuture案例

wen java案例 3

Java CompletableFuture 完整案例详解

基础概念

CompletableFuture 是 Java 8 引入的异步编程工具,实现了 Future 和 CompletionStage 接口,支持函数式编程风格。

Java CompletableFuture案例

基础用法案例

1 创建 CompletableFuture

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
public class BasicExample {
    public static void main(String[] args) throws ExecutionException, InterruptedException {
        // 1. 使用 supplyAsync 创建异步任务(有返回值)
        CompletableFuture<String> future1 = CompletableFuture.supplyAsync(() -> {
            try {
                TimeUnit.SECONDS.sleep(1);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            return "Hello from async task";
        });
        // 2. 使用 runAsync 创建异步任务(无返回值)
        CompletableFuture<Void> future2 = CompletableFuture.runAsync(() -> {
            try {
                TimeUnit.SECONDS.sleep(1);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            System.out.println("Running async task without return value");
        });
        // 获取结果
        System.out.println("Future1 result: " + future1.get());
        future2.get();
        // 3. 使用 completedFuture 创建已完成的任务
        CompletableFuture<String> completedFuture = CompletableFuture.completedFuture("Already completed");
        System.out.println("Completed future: " + completedFuture.get());
    }
}

异步任务编排案例

1 任务串联 (thenApply, thenAccept, thenRun)

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public class ChainExample {
    public static void main(String[] args) throws Exception {
        // thenApply: 对结果进行转换
        CompletableFuture<String> future = CompletableFuture
            .supplyAsync(() -> "Hello")
            .thenApply(s -> s + " World")           // 第一个转换
            .thenApply(String::toUpperCase);         // 第二个转换
        System.out.println("thenApply result: " + future.get());
        // thenAccept: 消费结果(无返回值)
        CompletableFuture.supplyAsync(() -> "Data")
            .thenAccept(data -> System.out.println("Processing: " + data));
        // thenRun: 纯粹执行任务(不关心结果)
        CompletableFuture.supplyAsync(() -> "Source")
            .thenRun(() -> System.out.println("Task completed"));
        // 等待所有任务完成
        TimeUnit.SECONDS.sleep(2);
    }
}

2 组合多个 CompletableFuture

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public class CombineExample {
    public static void main(String[] args) throws Exception {
        // thenCompose: 将两个 CompletableFuture 串联
        CompletableFuture<String> composedFuture = getUserInfo()
            .thenCompose(user -> getOrders(user));
        System.out.println("Composed result: " + composedFuture.get());
        // thenCombine: 并行执行两个任务并组合结果
        CompletableFuture<String> orderFuture = CompletableFuture
            .supplyAsync(() -> "订单信息")
            .thenCombine(
                CompletableFuture.supplyAsync(() -> "用户信息"),
                (order, user) -> order + " + " + user
            );
        System.out.println("Combined result: " + orderFuture.get());
        // thenAcceptBoth: 并行执行两个任务,消费两个结果
        CompletableFuture.supplyAsync(() -> "Order1")
            .thenAcceptBoth(
                CompletableFuture.supplyAsync(() -> "User1"),
                (order, user) -> System.out.println(order + " belongs to " + user)
            );
        TimeUnit.SECONDS.sleep(2);
    }
    private static CompletableFuture<String> getUserInfo() {
        return CompletableFuture.supplyAsync(() -> {
            sleep(1);
            return "user123";
        });
    }
    private static CompletableFuture<String> getOrders(String userId) {
        return CompletableFuture.supplyAsync(() -> {
            sleep(1);
            return "用户 " + userId + " 的订单列表";
        });
    }
    private static void sleep(int seconds) {
        try {
            TimeUnit.SECONDS.sleep(seconds);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
}

实际业务场景案例

1 异步获取多个接口数据

import java.util.Arrays;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
public class RealWorldExample {
    public static void main(String[] args) throws Exception {
        // 模拟从多个服务获取数据
        CompletableFuture<String> userInfo = getUserInfoAsync();
        CompletableFuture<List<String>> userOrders = getUserOrdersAsync();
        CompletableFuture<Double> userBalance = getUserBalanceAsync();
        // 等待所有任务完成
        CompletableFuture<Void> allOf = CompletableFuture.allOf(userInfo, userOrders, userBalance);
        allOf.get();  // 等待所有完成
        // 获取所有结果
        System.out.println("用户信息: " + userInfo.get());
        System.out.println("用户订单: " + userOrders.get());
        System.out.println("用户余额: " + userBalance.get());
        // 使用 anyOf 等待任意一个完成
        CompletableFuture<Object> anyResult = CompletableFuture.anyOf(userInfo, userOrders, userBalance);
        System.out.println("最先完成的任务: " + anyResult.get());
    }
    private static CompletableFuture<String> getUserInfoAsync() {
        return CompletableFuture.supplyAsync(() -> {
            sleep(2);
            return "张三, 30岁, 北京";
        });
    }
    private static CompletableFuture<List<String>> getUserOrdersAsync() {
        return CompletableFuture.supplyAsync(() -> {
            sleep(1);
            return Arrays.asList("订单1:手机", "订单2:电脑", "订单3:耳机");
        });
    }
    private static CompletableFuture<Double> getUserBalanceAsync() {
        return CompletableFuture.supplyAsync(() -> {
            sleep(3);
            return 9999.99;
        });
    }
    private static void sleep(int seconds) {
        try {
            TimeUnit.SECONDS.sleep(seconds);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
}

2 批量异步任务处理

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
public class BatchProcessExample {
    public static void main(String[] args) throws Exception {
        List<Integer> userIds = Arrays.asList(1, 2, 3, 4, 5, 6, 7, 8, 9, 10);
        // 并行处理所有用户
        List<CompletableFuture<UserData>> futures = userIds.stream()
            .map(id -> processUserAsync(id))
            .collect(Collectors.toList());
        // 等待所有任务完成并收集结果
        CompletableFuture<List<UserData>> allResults = CompletableFuture
            .allOf(futures.toArray(new CompletableFuture[0]))
            .thenApply(v -> futures.stream()
                .map(CompletableFuture::join)
                .collect(Collectors.toList()));
        List<UserData> results = allResults.get();
        results.forEach(System.out::println);
    }
    private static CompletableFuture<UserData> processUserAsync(Integer userId) {
        return CompletableFuture.supplyAsync(() -> {
            try {
                TimeUnit.MILLISECONDS.sleep(100);  // 模拟耗时操作
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            return new UserData(userId, "User_" + userId, Math.random() * 100);
        });
    }
    static class UserData {
        private Integer id;
        private String name;
        private double score;
        public UserData(Integer id, String name, double score) {
            this.id = id;
            this.name = name;
            this.score = score;
        }
        @Override
        public String toString() {
            return String.format("User{id=%d, name='%s', score=%.2f}", id, name, score);
        }
    }
}

异常处理案例

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public class ExceptionExample {
    public static void main(String[] args) throws Exception {
        // 1. exceptionally: 异常时提供默认值
        CompletableFuture<String> future1 = CompletableFuture
            .supplyAsync(() -> {
                if (Math.random() > 0.5) {
                    throw new RuntimeException("模拟异常1");
                }
                return "成功结果1";
            })
            .exceptionally(ex -> {
                System.out.println("处理异常1: " + ex.getMessage());
                return "默认值1";
            });
        System.out.println("Future1 result: " + future1.get());
        // 2. handle: 无论成功或失败都会执行
        CompletableFuture<String> future2 = CompletableFuture
            .supplyAsync(() -> {
                if (Math.random() > 0.5) {
                    throw new RuntimeException("模拟异常2");
                }
                return "成功结果2";
            })
            .handle((result, ex) -> {
                if (ex != null) {
                    System.out.println("捕获异常2: " + ex.getMessage());
                    return "错误恢复2";
                }
                return "成功结果2";
            });
        System.out.println("Future2 result: " + future2.get());
        // 3. 使用 whenComplete 处理完成事件(不改变结果)
        CompletableFuture.supplyAsync(() -> "测试数据")
            .whenComplete((result, ex) -> {
                if (ex == null) {
                    System.out.println("任务完成,结果: " + result);
                } else {
                    System.out.println("任务失败: " + ex.getMessage());
                }
            })
            .join();
    }
}

超时和取消控制

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
public class TimeoutExample {
    public static void main(String[] args) {
        // 1. 设置超时时间
        CompletableFuture<String> future1 = CompletableFuture
            .supplyAsync(() -> {
                try {
                    TimeUnit.SECONDS.sleep(3);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                return "长时间任务结果";
            })
            .completeOnTimeout("超时默认值", 2, TimeUnit.SECONDS);
        // 2. 或抛出 TimeoutException
        CompletableFuture<String> future2 = CompletableFuture
            .supplyAsync(() -> {
                try {
                    TimeUnit.SECONDS.sleep(3);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                return "长时间任务结果2";
            })
            .orTimeout(1, TimeUnit.SECONDS);
        try {
            // 获取结果
            System.out.println("Future1: " + future1.get(3, TimeUnit.SECONDS));
            // 这个会抛出 TimeoutException
            System.out.println("Future2: " + future2.get());
        } catch (ExecutionException | InterruptedException | TimeoutException e) {
            System.out.println("Future2 超时了: " + e.getMessage());
        }
        // 手动取消任务
        CompletableFuture<String> future3 = CompletableFuture.supplyAsync(() -> {
            try {
                Thread.sleep(5000);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            return "可取消的任务";
        });
        // 取消任务
        future3.cancel(true);
        System.out.println("Future3 已取消: " + future3.isCancelled());
    }
}

自定义线程池案例

import java.util.concurrent.*;
public class CustomExecutorExample {
    public static void main(String[] args) throws Exception {
        // 创建自定义线程池
        ThreadPoolExecutor executor = new ThreadPoolExecutor(
            4,                    // 核心线程数
            8,                    // 最大线程数
            60L,                  // 空闲线程存活时间
            TimeUnit.SECONDS,    // 时间单位
            new ArrayBlockingQueue<>(100),  // 任务队列
            new ThreadPoolExecutor.CallerRunsPolicy()  // 拒绝策略
        );
        // 1. 使用自定义线程池
        CompletableFuture<String> future = CompletableFuture
            .supplyAsync(() -> {
                System.out.println("当前线程: " + Thread.currentThread().getName());
                return "自定义线程池任务";
            }, executor);
        System.out.println("结果: " + future.get());
        // 2. 多个任务共享同一个线程池
        CompletableFuture<Integer> task1 = CompletableFuture
            .supplyAsync(() -> {
                try { Thread.sleep(1000); } catch (InterruptedException e) {}
                return 5;
            }, executor);
        CompletableFuture<Integer> task2 = CompletableFuture
            .supplyAsync(() -> {
                try { Thread.sleep(500); } catch (InterruptedException e) {}
                return 3;
            }, executor);
        // 组合两个任务
        CompletableFuture<Integer> combined = task1.thenCombine(task2, (x, y) -> x * y);
        System.out.println("组合结果: " + combined.get());
        // 关闭线程池
        executor.shutdown();
    }
}

实际综合应用示例

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
public class OrderProcessingSystem {
    public static void main(String[] args) throws Exception {
        OrderProcessingSystem system = new OrderProcessingSystem();
        // 创建订单并处理
        Order order = system.processOrder("ORD-001", "user123")
            .get();  // 同步等待最终结果
        System.out.println("订单处理完成:");
        System.out.println("订单号: " + order.getOrderId());
        System.out.println("用户: " + order.getUser());
        System.out.println("商品: " + order.getProducts());
        System.out.println("总价: " + order.getTotalPrice());
        System.out.println("状态: " + order.getStatus());
    }
    public CompletableFuture<Order> processOrder(String orderId, String userId) {
        // 1. 获取用户信息
        CompletableFuture<String> userInfoFuture = CompletableFuture
            .supplyAsync(() -> getMockUserInfo(userId));
        // 2. 获取库存信息
        CompletableFuture<Boolean> stockCheckFuture = CompletableFuture
            .supplyAsync(() -> checkStock(orderId));
        // 3. 计算价格
        CompletableFuture<Double> priceFuture = CompletableFuture
            .supplyAsync(() -> calculatePrice(orderId));
        // 并行执行所有任务,然后组合结果
        return userInfoFuture
            .thenCombine(stockCheckFuture, (user, inStock) -> {
                if (!inStock) {
                    throw new RuntimeException("库存不足");
                }
                return user;
            })
            .thenCombine(priceFuture, (user, price) -> {
                // 创建订单
                Order order = new Order(orderId, user);
                order.setTotalPrice(price);
                order.setStatus("CONFIRMED");
                return order;
            })
            .thenApply(order -> {
                // 模拟其他处理
                order.setStatus("PROCESSING");
                return order;
            })
            .exceptionally(ex -> {
                System.out.println("订单处理失败: " + ex.getMessage());
                Order failedOrder = new Order(orderId, "Unknown");
                failedOrder.setStatus("FAILED");
                return failedOrder;
            });
    }
    private String getMockUserInfo(String userId) {
        try {
            TimeUnit.MILLISECONDS.sleep(200);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return "用户: " + userId + ", 等级: VIP";
    }
    private Boolean checkStock(String orderId) {
        try {
            TimeUnit.MILLISECONDS.sleep(150);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return true;  // 简化模拟
    }
    private Double calculatePrice(String orderId) {
        try {
            TimeUnit.MILLISECONDS.sleep(100);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        return 2999.00;
    }
    static class Order {
        private String orderId;
        private String user;
        private double totalPrice;
        private String status;
        private List<String> products = new ArrayList<>();
        public Order(String orderId, String user) {
            this.orderId = orderId;
            this.user = user;
            this.products.add("iPhone 15");
            this.products.add("AirPods Pro");
        }
        // Getters and Setters
        public String getOrderId() { return orderId; }
        public String getUser() { return user; }
        public double getTotalPrice() { return totalPrice; }
        public void setTotalPrice(double totalPrice) { this.totalPrice = totalPrice; }
        public String getStatus() { return status; }
        public void setStatus(String status) { this.status = status; }
        public List<String> getProducts() { return products; }
    }
}

最佳实践建议

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class BestPractices {
    private final ExecutorService executor = Executors.newFixedThreadPool(10);
    public void bestPracticesDemo() {
        // 1. 始终使用自定义线程池,避免使用公共ForkJoinPool
        CompletableFuture<String> future = CompletableFuture
            .supplyAsync(() -> "任务", executor);
        // 2. 避免阻塞,使用异步方法
        future.thenAcceptAsync(result -> {
            // 异步处理结果
        }, executor);
        // 3. 正确处理异常
        future.exceptionally(throwable -> {
            System.err.println("处理失败: " + throwable);
            return "";
        });
        // 4. 设置合理的超时
        future.completeOnTimeout("timeout", 5, java.util.concurrent.TimeUnit.SECONDS);
        // 5. 使用 join() 而不是 get() 避免检查异常
        try {
            String result = future.join();
        } catch (Exception e) {
            e.printStackTrace();
        }
        // 6. 资源管理:确保关闭线程池
        executor.shutdown();
    }
}

这些案例覆盖了 CompletableFuture 的主要功能和应用场景,可根据实际需求选择合适的模式使用。

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