
在 Java Stream API 中,实现数据的并行处理非常简单,核心是通过 parallelStream() 方法获取并行流,而非默认的串行流(stream())。并行流会自动利用多核 CPU 的优势,将数据分成多个子任务并行执行,从而提升大数据量处理的效率。
Fork/Join 框架实现,自动将流中的元素分割成多个子流,由多个线程并行处理,最后合并结果。parallelStream() 方法(或流的 parallel() 方法将串行流转为并行流)。import java.util.Arrays;
import java.util.List;
public class ParallelStreamDemo {
public static void main(String[] args) {
// 准备一个大数据量的集合(1000万个整数)
List<Integer> numbers = Arrays.asList(new Integer[10_000_000]);
for (int i = 0; i < numbers.size(); i++) {
numbers.set(i, i);
}
// 串行流处理:计算偶数之和
long start = System.currentTimeMillis();
long serialSum = numbers.stream()
.filter(n -> n % 2 == 0)
.mapToLong(n -> n)
.sum();
long serialTime = System.currentTimeMillis() - start;
System.out.println("串行处理结果:" + serialSum + ",耗时:" + serialTime + "ms");
// 并行流处理:同样计算偶数之和
start = System.currentTimeMillis();
long parallelSum = numbers.parallelStream() // 关键:使用parallelStream()
.filter(n -> n % 2 == 0)
.mapToLong(n -> n)
.sum();
long parallelTime = System.currentTimeMillis() - start;
System.out.println("并行处理结果:" + parallelSum + ",耗时:" + parallelTime + "ms");
}
}输出(示例):
串行处理结果:24999995000000,耗时:120ms
并行处理结果:24999995000000,耗时:35ms // 并行效率更高(依赖CPU核心数)parallel() 方法)除了直接使用 parallelStream(),还可以通过 parallel() 方法将串行流转换为并行流:
List<String> words = Arrays.asList("apple", "banana", "cherry", "date");
// 串行流 → 转为并行流
long count = words.stream()
.parallel() // 切换为并行处理
.filter(word -> word.length() > 5)
.count();
System.out.println("长度大于5的单词数:" + count); // 输出:2(banana、cherry)forEach 累加全局变量),可能导致线程安全问题。
❌ 错误示例(共享变量不安全):int[] sum = {0}; // 共享数组
numbers.parallelStream()
.forEach(n -> sum[0] += n); // 多线程修改sum[0],结果可能不正确✅ 正确方式(使用线程安全的聚合操作):
long sum = numbers.parallelStream()
.mapToLong(n -> n)
.sum(); // sum() 内部线程安全filter 简单判断)的并行优势不明显,复杂操作(如大量计算)更适合并行。Ordered):并行流为提升效率可能打破元素顺序(如 forEach 输出顺序不确定),若需保持顺序,可用 forEachOrdered(但会损失部分并行性能)。Fork/Join 框架的公共线程池(ForkJoinPool.commonPool()),若需自定义线程池,可通过 ForkJoinPool 包装:import java.util.concurrent.ForkJoinPool;
ForkJoinPool pool = new ForkJoinPool(4); // 自定义4个核心线程的线程池
long sum = pool.submit(() ->
numbers.parallelStream()
.filter(n -> n % 2 == 0)
.mapToLong(n -> n)
.sum()
).get(); // 阻塞获取结果
pool.shutdown(); // 关闭线程池parallelStream() 或 stream().parallel() 获取并行流,后续操作与串行流一致。合理使用并行流能显著优化数据处理性能,但需根据具体场景评估是否适用。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。