Java 增强流解析
一、增强流的基本概念
1.1 什么是增强流
增强流(Stream)是 Java 8 中引入的一个全新的抽象概念,它代表了一组支持串行和并行聚合操作的元素序列。与传统的集合(Collection)不同,流并不实际存储元素,而是通过管道从数据源(如集合、数组、I/O通道等)中获取元素,并对其进行高效处理。
流的核心特性
-
管道化处理:流操作可以像Unix管道一样串联起来,形成处理流水线。例如:
List<String> names = Arrays.asList("John", "Alice", "Bob"); long count = names.stream() .filter(s -> s.startsWith("A")) .map(String::toUpperCase) .count(); -
操作类型:
- 中间操作(Intermediate Operations):返回新流,可以链式调用,如
filter(),map(),sorted() - 终端操作(Terminal Operations):产生结果或副作用,如
forEach(),collect(),reduce()
- 中间操作(Intermediate Operations):返回新流,可以链式调用,如
-
数据源多样性:可以从多种数据源创建流:
- 集合:
Collection.stream() - 数组:
Arrays.stream() - 文件:
Files.lines() - 生成器:
Stream.iterate(),Stream.generate()
- 集合:
1.2 流的特点
-
无存储特性:
- 流本身不是数据结构,不存储数据元素
- 数据元素存储在底层数据源中(如集合、数组等)
- 示例:
Stream.of(1,2,3)并不实际存储这些数字
-
函数式编程特性:
- 流操作不会修改源数据,而是产生新结果
- 例如
filter()操作会返回包含满足条件元素的新流,原集合不受影响 - 这一特性使得流操作更安全,适合并行处理
-
惰性执行机制:
- 中间操作不会立即执行,只有在终端操作触发时才会真正处理
- 示例:
Stream<String> stream = list.stream().filter(s -> { System.out.println("Filtering: " + s); return s.length() > 3; }); // 此时不会输出任何内容 stream.count(); // 此时才会执行过滤操作并输出
-
一次性消费特性:
- 流只能被消费一次,类似迭代器
- 尝试重复使用已关闭的流会抛出
IllegalStateException - 解决方案:每次需要时重新创建流
-
并行处理能力:
- 通过
parallelStream()可轻松实现并行处理 - 示例:
List<Integer> numbers = Arrays.asList(1,2,3,4,5); int sum = numbers.parallelStream() .filter(n -> n % 2 == 0) .mapToInt(Integer::intValue) .sum(); - 底层使用Fork/Join框架,自动利用多核处理器
- 通过
1.3流的典型应用场景
- 数据过滤和转换
- 集合元素的聚合计算
- I/O操作的高效处理
- 大数据量的并行处理
- 函数式编程风格的实现
流的这些特性使其成为Java 8中处理集合数据的强大工具,相比传统的循环迭代方式,流提供了更简洁、更高效的编程模式。
二、流的创建方式
流(Stream)是Java 8引入的一个强大的数据处理抽象,它允许我们以声明式方式处理数据集合。以下是Java中创建流的多种方式及其详细说明:
2.1 从集合创建流
集合框架是Java中最常用的数据结构,Java 8为Collection接口新增了两个方法来创建流:
List<String> list = Arrays.asList("a", "b", "c");
// 创建串行流 - 元素按顺序处理
Stream<String> stream = list.stream();
// 创建并行流 - 元素可以并行处理(适用于大数据集)
Stream<String> parallelStream = list.parallelStream();
应用场景:
- 串行流适合小数据集或需要保证处理顺序的情况
- 并行流适合大数据集且处理顺序不重要的场景,可以利用多核处理器提高性能
2.2 从数组创建流
数组是Java中的基本数据结构,Arrays工具类提供了创建流的方法:
String[] array = {"a", "b", "c"};
// 创建整个数组的流
Stream<String> fullArrayStream = Arrays.stream(array);
// 创建数组部分范围的流(从索引1开始,包含2个元素)
Stream<String> partialArrayStream = Arrays.stream(array, 1, 3);
特点:
- 可以创建整个数组或部分数组的流
- 数组流保持了元素的原始顺序
2.3 使用Stream静态方法创建流
Stream类提供了多种静态工厂方法来创建流:
of()方法
// 创建包含固定元素的流
Stream<String> fixedStream = Stream.of("a", "b", "c");
// 也可以创建包含数组元素的流
Stream<String> arrayStream = Stream.of(new String[]{"a", "b", "c"});
empty()方法
// 创建空流(常用于方法返回值的默认值)
Stream<String> emptyStream = Stream.empty();
generate()方法
// 创建无限流(通过Supplier函数生成元素)
Stream<Double> randomStream = Stream.generate(Math::random);
// 限制无限流的大小
randomStream.limit(10).forEach(System.out::println);
iterate()方法
// 创建无限流(通过种子和迭代函数)
Stream<Integer> evenNumbers = Stream.iterate(0, n -> n + 2);
// 限制流的大小并收集结果
List<Integer> first5Evens = evenNumbers.limit(5).collect(Collectors.toList());
2.4 从文件创建流
Java NIO的Files类提供了便捷的方法来处理文件内容:
// 使用try-with-resources确保流自动关闭
try (Stream<String> lines = Files.lines(Paths.get("data.txt"))) {
// 处理文件每行内容
lines.filter(line -> !line.isEmpty())
.forEach(System.out::println);
} catch (IOException e) {
e.printStackTrace();
}
特点:
- 自动处理文件编码(默认UTF-8)
- 可以指定字符集:
Files.lines(path, Charset.forName("GBK")) - 自动关闭流资源
2.5 从其他数据源创建流
从缓冲读取器创建流
try (BufferedReader reader = new BufferedReader(new FileReader("data.txt"))) {
Stream<String> lines = reader.lines();
long emptyLines = lines.filter(String::isEmpty).count();
System.out.println("空行数量: " + emptyLines);
} catch (IOException e) {
e.printStackTrace();
}
从正则表达式创建流
// 使用正则表达式分割字符串为流
Pattern pattern = Pattern.compile("\\s+"); // 按空白字符分割
Stream<String> words = pattern.splitAsStream("Java 8 Stream API");
// 使用正则表达式匹配创建流
Pattern digitPattern = Pattern.compile("\\d");
Stream<String> digits = digitPattern.matcher("a1b2c3").results()
.map(MatchResult::group);
其他特殊流
// IntStream、LongStream、DoubleStream等基本类型流
IntStream intStream = IntStream.range(1, 100); // 1-99
LongStream longStream = LongStream.rangeClosed(1, 100); // 1-100
// 随机数流
Random random = new Random();
DoubleStream randomDoubles = random.doubles(5); // 5个随机double
通过这些多样化的创建方式,Java流可以方便地处理各种数据源,为函数式编程提供了强大的支持。
三、流的操作分类
流的操作可以分为中间操作(Intermediate Operations)和终端操作(Terminal Operations)两大类,它们在流处理过程中扮演不同角色。
3.1 中间操作
中间操作是构建流处理管道的基础,它们会返回一个新的流,允许链式调用。中间操作具有惰性执行(Lazy Evaluation)特性,只有当终端操作被调用时才会真正执行。
3.1.1 筛选与切片
这些操作用于从流中筛选或截取特定元素。
filter(Predicate<? super T> predicate):返回一个包含满足Predicate条件的元素的流。Predicate是一个函数式接口,接收一个参数并返回布尔值。
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
Stream<Integer> evenNumbers = numbers.stream().filter(n -> n % 2 == 0); // 包含2,4
distinct():返回一个去除重复元素的流,基于元素的equals()方法判断是否重复。对于自定义对象,需要正确实现equals()和hashCode()方法。
List<Integer> numbers = Arrays.asList(1, 2, 2, 3, 3, 3);
Stream<Integer> distinctNumbers = numbers.stream().distinct(); // 包含1,2,3
limit(long maxSize):返回一个只包含前maxSize个元素的流,常用于限制处理元素数量。
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
Stream<Integer> limitedStream = numbers.stream().limit(3); // 包含1,2,3
skip(long n):返回一个跳过前n个元素的流,常用于分页处理。
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
Stream<Integer> skippedStream = numbers.stream().skip(2); // 包含3,4,5
3.1.2 映射操作
映射操作将流中的元素转换为其他形式。
map(Function<? super T, ? extends R> mapper):将每个元素通过mapper函数转换后返回一个新的流。
List<String> words = Arrays.asList("a", "bb", "ccc");
Stream<Integer> wordLengths = words.stream().map(String::length); // 包含1,2,3
flatMap(Function<? super T, ? extends Stream<? extends R>> mapper):将每个元素转换为一个流,然后将所有流合并为一个扁平化的流。常用于处理嵌套集合。
List<String> words = Arrays.asList("hello", "world");
Stream<Character> characters = words.stream()
.flatMap(word -> word.chars().mapToObj(c -> (char)c)); // 包含h,e,l,l,o,w,o,r,l,d
基本类型特化映射:
- mapToInt(ToIntFunction<? super T> mapper):转换为IntStream
- mapToLong(ToLongFunction<? super T> mapper):转换为LongStream
- mapToDouble(ToDoubleFunction<? super T> mapper):转换为DoubleStream
List<String> numbers = Arrays.asList("1", "2", "3");
IntStream intStream = numbers.stream().mapToInt(Integer::parseInt);
3.1.3 排序操作
sorted():返回一个按自然顺序排序的流,元素需实现Comparable接口。
List<Integer> numbers = Arrays.asList(3, 1, 2);
Stream<Integer> sortedStream = numbers.stream().sorted(); // 1,2,3
sorted(Comparator<? super T> comparator):返回一个按指定比较器排序的流。
List<String> words = Arrays.asList("b", "a", "c");
Stream<String> sortedStream = words.stream()
.sorted(Comparator.reverseOrder()); // c,b,a
3.1.4 元素消费(peek)
peek(Consumer<? super T> action):对每个元素执行action操作,主要用于调试,不影响流内容。
List<Integer> numbers = Arrays.asList(1, 2, 3);
numbers.stream()
.peek(n -> System.out.println("Before filter: " + n)) // 输出1,2,3
.filter(n -> n % 2 == 0) // 筛选偶数
.peek(n -> System.out.println("After filter: " + n)) // 输出2
.forEach(System.out::println); // 最终输出2
注意:peek不应用于修改流内容,它的主要用途是观察流处理过程中的中间状态。
3.2 终端操作
终端操作是流处理的最后一步,它会产生一个结果或副作用。需要注意的是,一旦执行终端操作后,流就会被消耗掉,无法再被使用。终端操作可以分为以下几类:
3.2.1 遍历元素
forEach(Consumer<? super T> action)
对每个元素执行给定的操作,操作顺序在并行流中是不确定的。
List<String> words = Arrays.asList("a", "b", "c");
words.stream().forEach(System.out::println); // 输出顺序可能是a,b,c,也可能是其他顺序
forEachOrdered(Consumer<? super T> action)
即使在并行流中,也会按原始顺序对每个元素执行给定的操作。
List<String> words = Arrays.asList("a", "b", "c");
words.parallelStream().forEachOrdered(System.out::println); // 始终按a,b,c顺序输出
3.2.2 收集结果
collect(Collector<? super T, A, R> collector)
将流中的元素收集到一个结果容器中,Collectors工具类提供了许多预定义的收集器。
List<String> words = Arrays.asList("a", "b", "c");
// 收集到List
List<String> collectedList = words.stream().collect(Collectors.toList());
// 收集到Set
Set<String> collectedSet = words.stream().collect(Collectors.toSet());
// 收集到Map
Map<String, Integer> lengthMap = words.stream()
.collect(Collectors.toMap(s -> s, String::length));
toArray()
将流中的元素转换为一个Object数组。
List<String> words = Arrays.asList("a", "b", "c");
Object[] array = words.stream().toArray();
toArray(IntFunction<A[]> generator)
将流中的元素转换为指定类型的数组。
List<String> words = Arrays.asList("a", "b", "c");
String[] array = words.stream().toArray(String[]::new); // 使用数组构造器引用
3.2.3 查找与匹配
allMatch(Predicate<? super T> predicate)
检查是否所有元素都满足给定条件。
List<Integer> numbers = Arrays.asList(2, 4, 6);
boolean allEven = numbers.stream().allMatch(n -> n % 2 == 0); // true
List<Integer> mixed = Arrays.asList(2, 4, 5);
boolean allEven2 = mixed.stream().allMatch(n -> n % 2 == 0); // false
anyMatch(Predicate<? super T> predicate)
检查是否有任何元素满足给定条件。
List<Integer> numbers = Arrays.asList(1, 3, 5);
boolean hasEven = numbers.stream().anyMatch(n -> n % 2 == 0); // false
List<Integer> mixed = Arrays.asList(1, 3, 6);
boolean hasEven2 = mixed.stream().anyMatch(n -> n % 2 == 0); // true
noneMatch(Predicate<? super T> predicate)
检查是否没有元素满足给定条件。
List<Integer> numbers = Arrays.asList(1, 3, 5);
boolean noEven = numbers.stream().noneMatch(n -> n % 2 == 0); // true
List<Integer> mixed = Arrays.asList(1, 3, 6);
boolean noEven2 = mixed.stream().noneMatch(n -> n % 2 == 0); // false
findFirst()
返回流中的第一个元素(Optional),在并行流中也保证返回第一个元素。
List<Integer> numbers = Arrays.asList(1, 2, 3);
Optional<Integer> first = numbers.stream().findFirst(); // Optional[1]
findAny()
返回流中的任意一个元素(Optional),在并行流中效率更高。
List<Integer> numbers = Arrays.asList(1, 2, 3);
Optional<Integer> any = numbers.parallelStream().findAny(); // 可能是1,2或3
3.2.4 归约
reduce(T identity, BinaryOperator<T> accumulator)
从初始值identity开始,通过累加器将元素累积起来。
List<Integer> numbers = Arrays.asList(1, 2, 3);
int sum = numbers.stream().reduce(0, Integer::sum); // 6
// 计算字符串连接
List<String> strings = Arrays.asList("a", "b", "c");
String concatenated = strings.stream().reduce("", String::concat); // "abc"
reduce(BinaryOperator<T> accumulator)
不提供初始值的归约,返回Optional(因为流可能为空)。
List<Integer> numbers = Arrays.asList(1, 2, 3);
Optional<Integer> sum = numbers.stream().reduce(Integer::sum); // Optional[6]
List<Integer> empty = Collections.emptyList();
Optional<Integer> emptySum = empty.stream().reduce(Integer::sum); // Optional.empty
reduce(U identity, BiFunction<U, ? super T, U> accumulator, BinaryOperator<U> combiner)
用于并行流的归约,combiner用于合并多个线程的结果。
List<Integer> numbers = Arrays.asList(1, 2, 3);
int sum = numbers.parallelStream()
.reduce(0, (partialSum, num) -> partialSum + num, Integer::sum); // 6
3.2.5 统计
count()
返回流中元素的数量。
List<String> words = Arrays.asList("a", "b", "c");
long count = words.stream().count(); // 3
long emptyCount = Stream.empty().count(); // 0
min(Comparator<? super T> comparator)
返回流中最小的元素(Optional)。
List<Integer> numbers = Arrays.asList(1, 2, 3);
Optional<Integer> min = numbers.stream().min(Integer::compare); // Optional[1]
Optional<Integer> emptyMin = Stream.<Integer>empty().min(Integer::compare); // Optional.empty
max(Comparator<? super T> comparator)
返回流中最大的元素(Optional)。
List<Integer> numbers = Arrays.asList(1, 2, 3);
Optional<Integer> max = numbers.stream().max(Integer::compare); // Optional[3]
Optional<Integer> emptyMax = Stream.<Integer>empty().max(Integer::compare); // Optional.empty
四、并行流
并行流是Java 8 Stream API提供的一个重要特性,它能够将流操作自动并行化,充分利用现代多核处理器的计算能力,在处理大规模数据集时显著提升性能。下面我们将深入探讨并行流的各个方面。
4.1 并行流的创建方式
4.1.1 通过集合的parallelStream()方法创建
这是最直接的并行流创建方式,示例代码如下:
List<String> list = Arrays.asList("a", "b", "c", "d", "e", "f", "g", "h");
// 直接创建并行流
Stream<String> parallelStream = list.parallelStream();
// 并行处理示例
parallelStream.forEach(System.out::println);
4.1.2 通过流的parallel()方法转换
对于已经存在的串行流,可以通过parallel()方法将其转换为并行流:
Stream<String> stream = Stream.of("a", "b", "c", "d", "e");
// 将串行流转为并行流
Stream<String> parallelStream = stream.parallel();
// 执行并行操作
long count = parallelStream.filter(s -> s.length() > 1).count();
4.2 并行流的执行原理
4.2.1 Fork/Join框架工作机制
并行流底层使用Java 7引入的Fork/Join框架实现,其工作流程如下:
- 任务分解(Fork):将大数据集分割成多个较小的子数据集
- 并行处理:每个子任务在不同的工作线程上独立执行
- 结果合并(Join):将各个子任务的处理结果汇总为最终结果
4.2.2 并行流的分割策略
并行流采用以下策略分割任务:
- 范围分割:对于ArrayList等可预测大小的集合,采用等分策略
- 迭代器分割:对于LinkedList等不可预测大小的集合,采用动态分割
- 深度限制:防止过度分割导致性能下降,默认递归深度限制为384
4.3 并行流的适用场景
4.3.1 适合使用并行流的场景
- 大数据集处理:当数据集规模超过10万条时,并行效果明显
- 计算密集型操作:如复杂数学运算、加密解密等耗时操作
- 无状态操作:如filter、map等不依赖前序结果的操作
// 适合并行处理的示例:大规模数据计算
List<Double> numbers = ... // 假设有100万个数字
double sum = numbers.parallelStream()
.mapToDouble(d -> complexCalculation(d))
.sum();
4.3.2 不建议使用并行流的场景
- 小数据集:数据量小于1万条时,并行开销可能超过收益
- 顺序依赖操作:如limit、findFirst等需要确定顺序的操作
- 有状态操作:如sorted、distinct等需要全局状态的操作
- 共享可变状态:操作中访问共享可变变量会导致线程安全问题
4.3.3 性能考量因素
- 数据规模:NQ模型(数据量N×单个任务处理量Q)决定并行效果
- 任务平衡性:各子任务耗时是否均衡
- 合并成本:结果合并操作的复杂度
- 硬件资源:可用CPU核心数和内存带宽
4.4 并行流使用建议
- 基准测试:使用JMH等工具进行性能测试对比
- 线程池调整:可通过
ForkJoinPool.commonPool()配置公共线程池 - 避免阻塞:不要在并行流中执行阻塞IO操作
- 有序性处理:必要时使用
forEachOrdered保持顺序
// 有序处理的并行流示例
List<String> results = data.parallelStream()
.map(processingFunction)
.collect(Collectors.toList());
通过合理使用并行流,可以在适当的场景下获得显著的性能提升,但需要根据具体业务场景和硬件环境进行细致的调优和测试。
五、收集器(Collector)
5.1 收集到集合
Collectors 类提供了多种将流元素收集到不同集合类型的方法:
toList()
将流中的元素收集到一个新的 ArrayList 中。
List<String> words = Arrays.asList("a", "b", "c");
List<String> list = words.stream().collect(Collectors.toList());
// 结果: ["a", "b", "c"]
toSet()
将流中的元素收集到一个新的 HashSet 中,自动去除重复元素。
List<String> words = Arrays.asList("a", "b", "b", "c");
Set<String> set = words.stream().collect(Collectors.toSet());
// 结果: ["a", "b", "c"] (去除了重复的"b")
toCollection(Supplier<C> collectionFactory)
允许指定具体的集合实现类型。
List<String> words = Arrays.asList("a", "b", "c");
// 收集到 LinkedList
LinkedList<String> linkedList = words.stream()
.collect(Collectors.toCollection(LinkedList::new));
// 收集到 TreeSet
TreeSet<String> treeSet = words.stream()
.collect(Collectors.toCollection(TreeSet::new));
5.2 聚合操作
Collectors 提供了多种数值聚合操作:
计数操作
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
long count = numbers.stream().collect(Collectors.counting());
// 结果: 5
求和操作
// 整数求和
int sumInt = numbers.stream()
.collect(Collectors.summingInt(n -> n));
// 结果: 15
// 长整数求和
List<Long> longNumbers = Arrays.asList(1L, 2L, 3L);
long sumLong = longNumbers.stream()
.collect(Collectors.summingLong(n -> n));
// 结果: 6
// 浮点数求和
List<Double> doubles = Arrays.asList(1.1, 2.2, 3.3);
double sumDouble = doubles.stream()
.collect(Collectors.summingDouble(n -> n));
// 结果: 6.6
平均值计算
double average = numbers.stream()
.collect(Collectors.averagingInt(n -> n));
// 结果: 3.0
统计汇总
IntSummaryStatistics stats = numbers.stream()
.collect(Collectors.summarizingInt(n -> n));
// 可以获取以下统计信息:
// stats.getCount() - 元素数量
// stats.getSum() - 总和
// stats.getAverage() - 平均值
// stats.getMin() - 最小值
// stats.getMax() - 最大值
5.3 分组与分区
groupingBy 分组
List<String> words = Arrays.asList("a", "bb", "ccc", "dddd");
// 简单分组
Map<Integer, List<String>> groupByLength = words.stream()
.collect(Collectors.groupingBy(String::length));
// 结果: {1=["a"], 2=["bb"], 3=["ccc"], 4=["dddd"]}
// 分组后应用下游收集器
Map<Integer, Long> groupByLengthCount = words.stream()
.collect(Collectors.groupingBy(String::length, Collectors.counting()));
// 结果: {1=1, 2=1, 3=1, 4=1}
// 多级分组
Map<Integer, Map<Character, List<String>>> multiLevelGroup = words.stream()
.collect(Collectors.groupingBy(String::length,
Collectors.groupingBy(s -> s.charAt(0))));
partitioningBy 分区
// 按长度是否大于2分区
Map<Boolean, List<String>> partitionByLength = words.stream()
.collect(Collectors.partitioningBy(s -> s.length() > 2));
// 结果: {false=["a", "bb"], true=["ccc", "dddd"]}
// 分区后应用下游收集器
Map<Boolean, Long> partitionCount = words.stream()
.collect(Collectors.partitioningBy(s -> s.length() > 2, Collectors.counting()));
// 结果: {false=2, true=2}
5.4 连接字符串
简单连接
List<String> words = Arrays.asList("a", "b", "c");
String joined1 = words.stream().collect(Collectors.joining());
// 结果: "abc"
使用分隔符连接
String joined2 = words.stream().collect(Collectors.joining(","));
// 结果: "a,b,c"
String joined3 = words.stream().collect(Collectors.joining(", ", "开始: ", " 结束"));
// 结果: "开始: a, b, c 结束"
复杂示例
List<Person> people = Arrays.asList(
new Person("John", 25),
new Person("Alice", 30),
new Person("Bob", 20)
);
// 连接所有人员的姓名,用逗号分隔
String names = people.stream()
.map(Person::getName)
.collect(Collectors.joining(", "));
// 结果: "John, Alice, Bob"
// 使用前缀和后缀
String formattedNames = people.stream()
.map(Person::getName)
.collect(Collectors.joining(", ", "[", "]"));
// 结果: "[John, Alice, Bob]"
5.5 自定义收集器
自定义收集器需要实现Collector接口,该接口包含5个核心方法,每个方法都有特定的职责:
-
supplier()
返回一个Supplier函数式接口实现,用于创建新的结果容器。例如在字符串拼接场景中,这个方法会返回一个新的StringJoiner实例。 -
accumulator()
返回一个BiConsumer,定义如何将流中的元素累积到结果容器中。对于字符串收集器,这是简单的StringJoiner::add方法引用。 -
combiner()
返回一个BinaryOperator,用于并行流处理时合并两个部分结果。字符串收集器使用StringJoiner::merge来合并两个StringJoiner。 -
finisher()
返回一个Function,将中间累积类型转换为最终结果类型。对于字符串收集器,这是调用StringJoiner::toString完成最终拼接。 -
characteristics()
返回一个不可变的Set<Characteristics>,描述收集器的特性,影响流处理的优化方式。
Characteristics枚举值
Characteristics枚举定义了收集器的三种重要特性:
-
CONCURRENT
表示收集器支持并发累积,多个线程可以同时操作同一个结果容器。使用时必须确保结果容器是线程安全的。 -
UNORDERED
表示收集器不保留元素的原始顺序。这在处理无序集合(如HashSet)或并行流时可以提高性能。 -
IDENTITY_FINISH
表示finisher()方法是恒等转换,可以跳过。如果设置此标志,中间累积类型必须与最终结果类型相同。
高级实现示例
扩展字符串连接收集器实现,增加更多功能:
public class AdvancedStringJoiningCollector implements Collector<String, StringJoiner, String> {
private final String delimiter;
private final String prefix;
private final String suffix;
private final String emptyValue;
private final Predicate<String> filter;
public AdvancedStringJoiningCollector(String delimiter, String prefix,
String suffix, String emptyValue,
Predicate<String> filter) {
this.delimiter = delimiter;
this.prefix = prefix;
this.suffix = suffix;
this.emptyValue = emptyValue;
this.filter = filter;
}
@Override
public Supplier<StringJoiner> supplier() {
return () -> new StringJoiner(delimiter, prefix, suffix);
}
@Override
public BiConsumer<StringJoiner, String> accumulator() {
return (joiner, str) -> {
if (filter.test(str)) {
joiner.add(str);
}
};
}
@Override
public BinaryOperator<StringJoiner> combiner() {
return StringJoiner::merge;
}
@Override
public Function<StringJoiner, String> finisher() {
return joiner -> {
String result = joiner.toString();
return result.equals(prefix + suffix) ? emptyValue : result;
};
}
@Override
public Set<Characteristics> characteristics() {
return Collections.unmodifiableSet(EnumSet.of(Characteristics.IDENTITY_FINISH));
}
}
使用场景示例
1.处理空流情况
List<String> emptyList = Collections.emptyList();
String result = emptyList.stream()
.collect(new AdvancedStringJoiningCollector(",", "[", "]", "Empty List", s -> true));
System.out.println(result); // 输出: Empty List
2.带过滤条件的拼接
List<String> words = Arrays.asList("apple", "banana", "cherry", "date");
String filtered = words.stream()
.collect(new AdvancedStringJoiningCollector("-", "", "", "No matches",
s -> s.length() > 5));
System.out.println(filtered); // 输出: banana-cherry
3.并行流处理
List<String> bigList = IntStream.range(0, 10000)
.mapToObj(i -> "item" + i)
.collect(Collectors.toList());
String parallelResult = bigList.parallelStream()
.collect(new AdvancedStringJoiningCollector("|", "{", "}", "Empty",
s -> s.contains("99")));
System.out.println(parallelResult); // 输出所有包含"99"的项
六、流的注意事项
6.1 流的一次性使用
Java 8 中的流(Stream)是一种一次性消费的数据结构,类似于迭代器。一旦执行了终端操作(如 forEach、collect、count 等),流就会被关闭,再次使用会抛出 IllegalStateException。这是因为流设计为延迟计算(lazy evaluation)模式,只有在终端操作时才会真正执行数据处理。
Stream<String> stream = Stream.of("a", "b", "c");
// 第一次消费流
stream.forEach(System.out::println); // 输出: a b c
// 再次尝试消费已关闭的流
stream.forEach(System.out::println); // 抛出IllegalStateException: stream has already been operated upon or closed
如果需要多次操作相同数据,可以将流转换为集合后重新创建流:
List<String> list = Stream.of("a", "b", "c").collect(Collectors.toList());
Stream<String> stream1 = list.stream();
Stream<String> stream2 = list.stream(); // 可以创建多个新流
6.2 避免副作用
流的函数式操作应该是无副作用的(pure function),即不应该修改外部状态。虽然在 forEach 等操作中可以修改外部状态,但这违背了函数式编程的原则,会导致代码难以理解和维护,在并行流中还可能引发线程安全问题。
不推荐的做法(有副作用)
List<Integer> numbers = Arrays.asList(1, 2, 3, 4, 5);
List<Integer> result = new ArrayList<>();
// 在forEach中修改外部集合(副作用)
numbers.stream()
.filter(n -> n % 2 == 0)
.forEach(result::add); // 不推荐,尤其当改为parallelStream时会有线程安全问题
推荐的做法(无副作用)
List<Integer> result = numbers.stream()
.filter(n -> n % 2 == 0)
.collect(Collectors.toList()); // 使用collect替代副作用操作
6.3 注意并行流的线程安全
并行流(parallelStream)会利用多核处理器并行处理数据,但同时也带来了线程安全问题。如果在流的操作中修改共享变量,必须使用同步机制或线程安全的数据结构。
线程不安全的示例
List<Integer> numbers = IntStream.range(0, 10000).boxed().collect(Collectors.toList());
List<Integer> result = new ArrayList<>(); // 非线程安全集合
numbers.parallelStream()
.map(n -> n * 2)
.forEach(result::add); // 可能导致元素丢失、重复或ArrayIndexOutOfBoundsException
线程安全的解决方案
1.使用线程安全集合:
List<Integer> result = Collections.synchronizedList(new ArrayList<>());
numbers.parallelStream().forEach(result::add);
2.使用并发集合:
List<Integer> result = new CopyOnWriteArrayList<>();
numbers.parallelStream().forEach(result::add);
3.使用collect方法(推荐):
List<Integer> result = numbers.parallelStream()
.collect(Collectors.toList()); // 内部处理了线程安全问题
6.4 无限流的处理
Stream API 提供了生成无限流的方法,如 Stream.generate() 和 Stream.iterate()。如果不加限制地使用这些流,会导致程序陷入无限循环甚至内存溢出。
生成随机数的无限流
// 生成10个随机数
Stream.generate(Math::random)
.limit(10) // 必须限制数量
.forEach(System.out::println);
数字序列的无限流
// 生成从1开始的奇数序列,取前5个
Stream.iterate(1, n -> n + 2)
.limit(5)
.forEach(System.out::println); // 输出: 1 3 5 7 9
Java 9 还提供了带条件的 iterate 方法,可以更安全地创建有限流:
// 生成小于100的斐波那契数列
Stream.iterate(new int[]{0, 1}, t -> new int[]{t[1], t[0] + t[1]})
.takeWhile(t -> t[0] < 100)
.map(t -> t[0])
.forEach(System.out::println);
6.5 避免过度使用并行流
虽然并行流可以利用多核处理器的优势,但并非所有场景都适合使用并行流。需要考虑以下因素:
1.数据量大小:对于小数据集(通常少于1000个元素),并行流的线程创建和上下文切换开销可能超过性能收益。
// 不合适的并行流使用(数据量太小)
List<Integer> smallList = IntStream.range(0, 100).boxed().collect(Collectors.toList());
smallList.parallelStream().map(n -> n * 2).count(); // 可能比串行流更慢
2.操作复杂性:如果流中的操作非常简单(如简单的数学运算),并行化的收益可能不明显。
// 简单操作可能不适合并行化
List<Integer> numbers = IntStream.range(0, 1000000).boxed().collect(Collectors.toList());
long count = numbers.parallelStream().filter(n -> n % 2 == 0).count(); // 收益有限
3.可拆分性:某些数据源(如IO流、LinkedList)难以高效拆分,会影响并行效果。
最佳实践是先用串行流实现功能,在性能测试确认瓶颈后再考虑并行化,并实际测量性能提升。
6.6 注意流的顺序
流操作中的顺序保证是重要考虑因素:
1.串行流保持数据源顺序:
List<Integer> numbers = Arrays.asList(3, 1, 4, 1, 5, 9);
// 串行流保持顺序
numbers.stream()
.map(n -> n * 2)
.forEach(System.out::print); // 输出顺序保证: 6 2 8 2 10 18
2.并行流默认不保证顺序:
// 并行流顺序不确定
numbers.parallelStream()
.map(n -> n * 2)
.forEach(System.out::print); // 输出顺序可能变化,如: 2 6 18 10 8 2
3.强制顺序的方法:
// 使用forEachOrdered保证顺序
numbers.parallelStream()
.map(n -> n * 2)
.forEachOrdered(System.out::print); // 输出: 6 2 8 2 10 18
// 或在collect时保持顺序
List<Integer> result = numbers.parallelStream()
.map(n -> n * 2)
.collect(Collectors.toList()); // 保持原始顺序
注意:某些中间操作如 sorted() 会引入顺序保证,而 unordered() 可以显式放弃顺序约束以提高性能。
6.7 避免在流操作中使用阻塞操作
流的操作应该是非阻塞且高效的,阻塞操作会严重影响流处理的性能,特别是在并行流中可能导致线程耗尽。
不推荐的做法(包含阻塞操作)
List<String> urls = Arrays.asList("url1", "url2", "url3");
// 在流中直接进行网络IO(阻塞操作)
List<String> contents = urls.stream()
.map(url -> {
try {
// 模拟网络请求
Thread.sleep(1000);
return fetchContent(url);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
})
.collect(Collectors.toList());
推荐的替代方案
1.使用CompletableFuture实现异步:
List<CompletableFuture<String>> futures = urls.stream()
.map(url -> CompletableFuture.supplyAsync(() -> fetchContent(url)))
.collect(Collectors.toList());
List<String> contents = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
2.使用专门的反应式库(如RxJava、Reactor)处理IO密集型操作。
3.将阻塞操作与流处理分离:
// 先获取所有数据
List<String> contents = new ArrayList<>();
for (String url : urls) {
contents.add(fetchContent(url));
}
// 然后进行流处理
contents.stream().filter(...).map(...)...
对于必须包含IO操作的流处理,可以考虑使用 Stream.iterator() 或第三方库提供的特殊流实现。
七、流的实际应用示例
7.1 过滤并收集元素
List<String> words = Arrays.asList("apple", "banana", "cherry", "date");
// 使用Stream API过滤出长度大于5的单词
// filter()方法接收一个Predicate函数式接口,用于判断元素是否满足条件
// collect()方法将流中的元素收集到指定容器中,这里使用Collectors.toList()收集到List
List<String> longWords = words.stream()
.filter(word -> word.length() > 5)
.collect(Collectors.toList());
System.out.println(longWords); // 输出结果:[banana, cherry]
// 其中"apple"长度为5,"date"长度为4,均不满足条件,被过滤掉
7.2 映射并计算总和
// 定义Product类,包含名称、数量和单价属性
List<Product> products = Arrays.asList(
new Product("apple", 10, 2.5),
new Product("banana", 20, 1.5),
new Product("cherry", 15, 3.0)
);
// 使用Stream API计算所有产品的总金额
// mapToDouble()将每个产品映射为其金额值(数量*单价)
// sum()方法对流中的Double值求和
double totalAmount = products.stream()
.mapToDouble(product -> product.getQuantity() * product.getPrice())
.sum();
System.out.println("总金额: " + totalAmount);
// 详细计算过程:10*2.5=25 + 20*1.5=30 + 15*3.0=45 → 总金额=100.0
7.3 分组并统计
// 定义Student类,包含姓名、科目和分数属性
List<Student> students = Arrays.asList(
new Student("Alice", "Math", 90),
new Student("Bob", "Math", 85),
new Student("Charlie", "English", 95),
new Student("David", "English", 80)
);
// 使用Stream API按科目分组并计算平均分
// groupingBy()方法接收分类函数和下游收集器
// averagingInt()计算整数值的平均数
Map<String, Double> averageScoreBySubject = students.stream()
.collect(Collectors.groupingBy(
Student::getSubject,
Collectors.averagingInt(Student::getScore)
));
System.out.println("各科目平均分: " + averageScoreBySubject);
// 计算结果:
// Math科目:(90+85)/2=87.5
// English科目:(95+80)/2=87.5
7.4 并行流处理大量数据
// 生成100万个0-1000之间的随机整数
// Random.ints()方法生成IntStream,boxed()转换为Stream<Integer>
List<Integer> numbers = new Random()
.ints(1_000_000, 0, 1000)
.boxed()
.collect(Collectors.toList());
// 使用并行流处理大数据集
// parallelStream()创建并行流,自动利用多核处理器
// filter()筛选偶数,mapToLong()转换为long类型,sum()求和
long evenSum = numbers.parallelStream()
.filter(n -> n % 2 == 0)
.mapToLong(Long::valueOf)
.sum();
System.out.println("偶数总和: " + evenSum);
// 并行流特别适合处理大数据集,可以显著提高处理速度
// 注意:并行流不保证元素的处理顺序,且对状态共享操作需谨慎
更多推荐


所有评论(0)