一、增强流的基本概念

1.1 什么是增强流

增强流(Stream)是 Java 8 中引入的一个全新的抽象概念,它代表了一组支持串行和并行聚合操作的元素序列。与传统的集合(Collection)不同,流并不实际存储元素,而是通过管道从数据源(如集合、数组、I/O通道等)中获取元素,并对其进行高效处理。

流的核心特性

  1. 管道化处理:流操作可以像Unix管道一样串联起来,形成处理流水线。例如:

    List<String> names = Arrays.asList("John", "Alice", "Bob");
    long count = names.stream()
                    .filter(s -> s.startsWith("A"))
                    .map(String::toUpperCase)
                    .count();
    

  2. 操作类型

    • 中间操作(Intermediate Operations):返回新流,可以链式调用,如filter(), map(), sorted()
    • 终端操作(Terminal Operations):产生结果或副作用,如forEach(), collect(), reduce()
  3. 数据源多样性:可以从多种数据源创建流:

    • 集合:Collection.stream()
    • 数组:Arrays.stream()
    • 文件:Files.lines()
    • 生成器:Stream.iterate(), Stream.generate()

1.2 流的特点

  1. 无存储特性

    • 流本身不是数据结构,不存储数据元素
    • 数据元素存储在底层数据源中(如集合、数组等)
    • 示例:Stream.of(1,2,3)并不实际存储这些数字
  2. 函数式编程特性

    • 流操作不会修改源数据,而是产生新结果
    • 例如filter()操作会返回包含满足条件元素的新流,原集合不受影响
    • 这一特性使得流操作更安全,适合并行处理
  3. 惰性执行机制

    • 中间操作不会立即执行,只有在终端操作触发时才会真正处理
    • 示例:
      Stream<String> stream = list.stream().filter(s -> {
          System.out.println("Filtering: " + s);
          return s.length() > 3;
      }); // 此时不会输出任何内容
      
      stream.count(); // 此时才会执行过滤操作并输出
      

  4. 一次性消费特性

    • 流只能被消费一次,类似迭代器
    • 尝试重复使用已关闭的流会抛出IllegalStateException
    • 解决方案:每次需要时重新创建流
  5. 并行处理能力

    • 通过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流的典型应用场景

  1. 数据过滤和转换
  2. 集合元素的聚合计算
  3. I/O操作的高效处理
  4. 大数据量的并行处理
  5. 函数式编程风格的实现

流的这些特性使其成为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框架实现,其工作流程如下:

  1. 任务分解(Fork):将大数据集分割成多个较小的子数据集
  2. 并行处理:每个子任务在不同的工作线程上独立执行
  3. 结果合并(Join):将各个子任务的处理结果汇总为最终结果

4.2.2 并行流的分割策略

并行流采用以下策略分割任务:

  1. 范围分割:对于ArrayList等可预测大小的集合,采用等分策略
  2. 迭代器分割:对于LinkedList等不可预测大小的集合,采用动态分割
  3. 深度限制:防止过度分割导致性能下降,默认递归深度限制为384

4.3 并行流的适用场景

4.3.1 适合使用并行流的场景

  1. 大数据集处理:当数据集规模超过10万条时,并行效果明显
  2. 计算密集型操作:如复杂数学运算、加密解密等耗时操作
  3. 无状态操作:如filter、map等不依赖前序结果的操作
// 适合并行处理的示例:大规模数据计算
List<Double> numbers = ... // 假设有100万个数字
double sum = numbers.parallelStream()
                   .mapToDouble(d -> complexCalculation(d))
                   .sum();

4.3.2 不建议使用并行流的场景

  1. 小数据集:数据量小于1万条时,并行开销可能超过收益
  2. 顺序依赖操作:如limit、findFirst等需要确定顺序的操作
  3. 有状态操作:如sorted、distinct等需要全局状态的操作
  4. 共享可变状态:操作中访问共享可变变量会导致线程安全问题

4.3.3 性能考量因素

  1. 数据规模:NQ模型(数据量N×单个任务处理量Q)决定并行效果
  2. 任务平衡性:各子任务耗时是否均衡
  3. 合并成本:结果合并操作的复杂度
  4. 硬件资源:可用CPU核心数和内存带宽

4.4 并行流使用建议

  1. 基准测试:使用JMH等工具进行性能测试对比
  2. 线程池调整:可通过ForkJoinPool.commonPool()配置公共线程池
  3. 避免阻塞:不要在并行流中执行阻塞IO操作
  4. 有序性处理:必要时使用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个核心方法,每个方法都有特定的职责:

  1. supplier()
    返回一个Supplier函数式接口实现,用于创建新的结果容器。例如在字符串拼接场景中,这个方法会返回一个新的StringJoiner实例。

  2. accumulator()
    返回一个BiConsumer,定义如何将流中的元素累积到结果容器中。对于字符串收集器,这是简单的StringJoiner::add方法引用。

  3. combiner()
    返回一个BinaryOperator,用于并行流处理时合并两个部分结果。字符串收集器使用StringJoiner::merge来合并两个StringJoiner

  4. finisher()
    返回一个Function,将中间累积类型转换为最终结果类型。对于字符串收集器,这是调用StringJoiner::toString完成最终拼接。

  5. characteristics()
    返回一个不可变的Set<Characteristics>,描述收集器的特性,影响流处理的优化方式。

Characteristics枚举值

Characteristics枚举定义了收集器的三种重要特性:

  1. CONCURRENT
    表示收集器支持并发累积,多个线程可以同时操作同一个结果容器。使用时必须确保结果容器是线程安全的。

  2. UNORDERED
    表示收集器不保留元素的原始顺序。这在处理无序集合(如HashSet)或并行流时可以提高性能。

  3. 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);
// 并行流特别适合处理大数据集,可以显著提高处理速度
// 注意:并行流不保证元素的处理顺序,且对状态共享操作需谨慎

Logo

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

更多推荐