Java 并发编程:CompletableFuture 与并行流

1. 并行流(Parallel Stream)
  • 基础概念
    并行流基于 Java 8 的 Stream API,利用多核处理器自动拆分任务并行执行。其核心是 ForkJoinPool,默认线程数为处理器核心数。
    示例代码:计算整数平方和

    long sum = IntStream.range(1, 1000)
                       .parallel()  // 启用并行
                       .map(x -> x * x)
                       .sum();
    

  • 适用场景

    • 数据密集型计算(如批量处理集合)
    • 无状态操作(操作间无依赖)
    • 简单任务拆分(如 $ \sum_{i=1}^{n} f(i) $ 类问题)
  • 局限性

    • 阻塞操作会降低性能(如 I/O 等待)
    • 任务拆分不均可能导致负载失衡
    • 调试复杂(异常堆栈不直观)

2. CompletableFuture
  • 基础概念
    用于异步编程的类,支持任务组合、回调和非阻塞操作。可自定义线程池,避免共享池阻塞。
    示例代码:异步获取数据并处理

    CompletableFuture.supplyAsync(() -> fetchData())  // 异步执行
                     .thenApply(data -> process(data)) // 链式处理
                     .exceptionally(ex -> handleError(ex)); // 异常处理
    

  • 核心优势

    • 任务组合:支持 thenCombine(), allOf() 等组合操作
      例如:$$ \text{Task}_C = \text{Task}_A \oplus \text{Task}_B $$
    • 灵活线程控制:可指定专用线程池
    • 异常隔离:每个任务独立处理异常
  • 适用场景

    • I/O 密集型任务(如网络请求)
    • 多任务依赖(如 B 需 A 的结果)
    • 超时控制(通过 orTimeout() 实现)

3. 对比与选型
维度 并行流 CompletableFuture
任务类型 计算密集型 I/O 密集型或混合型
线程控制 共享池,不可定制 可自定义线程池
错误处理 全流终止 单任务隔离处理
复杂度 低(声明式) 中高(需编排逻辑)
性能关键点 数据分片效率 异步回调开销

4. 实践建议
  1. 优先并行流:当满足 $ \text{数据量} \gg \text{单任务耗时} $ 且无阻塞时
  2. 选择 CompletableFuture 当:
    • 需组合多个异步源(如数据库 + API)
    • 需精细控制超时/回退(如 completeOnTimeout()
    • 避免共享池阻塞(例如:newFixedThreadPool()
  3. 混合使用
    List<CompletableFuture<Void>> tasks = dataList.stream()
        .map(item -> CompletableFuture.runAsync(() -> process(item), customPool))
        .collect(Collectors.toList());
    CompletableFuture.allOf(tasks.toArray(new CompletableFuture[0])).join();
    

关键总结

  • 并行流 ≈ 自动并行化的「批量计算流水线」
  • CompletableFuture ≈ 可编排的「异步任务工作流」
  • 在 $ \text{CPU核心数} \times \text{任务粒度} \approx \text{吞吐量} $ 的约束下优化选择
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐