Stream API 从 Java 8 开始提供,用声明式管道处理集合、数组、文件和生成序列。Stream 表示一段一次性的数据计算过程,不是保存元素的数据结构。

一条 Stream 管道通常由“数据源 → 零个或多个中间操作 → 一个终止操作”构成。中间操作通常是惰性的,真正遍历数据一般从终止操作开始。

一、Stream 的核心特征

  • 声明式:描述要完成的过滤、转换和归约,而不是手写遍历过程。
  • 内部迭代:由 Stream 管道控制元素遍历。
  • 惰性执行:filter、map 等中间操作在终止操作之前通常不会真正处理数据。
  • 一次性消费:一个 Stream 不能执行两条管道,也不能在终止后再次使用。
  • 可顺序或并行执行:并行不等于一定更快,必须根据任务与数据验证。

二、创建 Stream

List<String> languages = List.of("Java", "Python", "Go");

Stream<String> fromCollection = languages.stream();
Stream<Integer> fromValues = Stream.of(1, 2, 3, 4);
IntStream fromRange = IntStream.rangeClosed(1, 10);
Stream<Double> generated = Stream.generate(Math::random).limit(5);
Stream<Integer> iterated = Stream.iterate(1, n -> n + 1).limit(10);

数组可以通过 Arrays.stream(array) 创建流。Files.lines(path) 等 I/O 来源的 Stream 持有外部资源,应使用 try-with-resources 及时关闭。

try (Stream<String> lines = Files.lines(path)) {
    long nonBlankCount = lines
        .filter(line -> !line.isBlank())
        .count();
}

三、常用中间操作

过滤与截取

List<String> result = languages.stream()
    .filter(name -> name.length() > 2)
    .distinct()
    .skip(1)
    .limit(10)
    .toList();

map:一对一转换

List<String> upperNames = languages.stream()
    .map(String::toUpperCase)
    .toList();

flatMap:一对多展开

List<String> words = List.of("java stream", "lambda api")
    .stream()
    .flatMap(line -> Arrays.stream(line.split("\s+")))
    .distinct()
    .toList();

排序

List<User> sortedUsers = users.stream()
    .sorted(
        Comparator.comparingInt(User::age)
                  .thenComparing(User::name)
    )
    .toList();
distinct 和 sorted 是有状态中间操作,可能需要保存大量元素;filter 和 map 通常是无状态操作,更容易单遍流水处理。

四、常用终止操作

long count = users.stream().count();

boolean hasAdult = users.stream().anyMatch(user -> user.age() >= 18);
boolean allActive = users.stream().allMatch(User::active);
boolean noneBlocked = users.stream().noneMatch(User::blocked);

Optional<User> oldest = users.stream()
    .max(Comparator.comparingInt(User::age));

users.forEach(System.out::println);
  • findFirst、findAny、anyMatch、allMatch 等操作可能短路,不必处理完整输入。
  • forEach 不保证并行流的遇见顺序;需要顺序时使用 forEachOrdered。
  • min、max、findFirst 等可能没有结果,因此返回 Optional。

五、收集结果

toList 与 Collectors.toList

List<String> unmodifiable = languages.stream()
    .map(String::toLowerCase)
    .toList();

List<String> mutable = languages.stream()
    .map(String::toLowerCase)
    .collect(Collectors.toCollection(ArrayList::new));
Stream.toList() 返回不可修改的 List。Collectors.toList() 不保证具体实现或可修改性;如果明确需要可变列表,应使用 toCollection(ArrayList::new)。

分组与下游归约

Map<Integer, List<User>> usersByAge = users.stream()
    .collect(Collectors.groupingBy(User::age));

Map<Integer, Long> countByAge = users.stream()
    .collect(Collectors.groupingBy(
        User::age,
        Collectors.counting()
    ));

分区、拼接与映射到 Map

Map<Boolean, List<User>> adults = users.stream()
    .collect(Collectors.partitioningBy(user -> user.age() >= 18));

String names = users.stream()
    .map(User::name)
    .collect(Collectors.joining(", ", "[", "]"));

Map<Long, User> usersById = users.stream()
    .collect(Collectors.toMap(
        User::id,
        Function.identity(),
        (existing, replacement) -> existing
    ));
Collectors.toMap 遇到重复键时默认抛出异常。业务上可能重复时,应显式提供合并函数。

六、数值流与统计

int totalAge = users.stream()
    .mapToInt(User::age)
    .sum();

OptionalDouble averageAge = users.stream()
    .mapToInt(User::age)
    .average();

IntSummaryStatistics stats = users.stream()
    .mapToInt(User::age)
    .summaryStatistics();

System.out.println(stats.getCount());
System.out.println(stats.getMin());
System.out.println(stats.getMax());
System.out.println(stats.getAverage());

IntStream、LongStream 和 DoubleStream 可以减少基本类型的装箱开销,并直接提供 sum、average、summaryStatistics 等操作。


七、reduce 归约

int sum = Stream.of(1, 2, 3, 4)
    .reduce(0, Integer::sum);

Optional<Integer> max = Stream.of(1, 5, 3)
    .reduce(Integer::max);
  • 带 identity 的 reduce 在空流时返回 identity。
  • 不带 identity 时返回 Optional,因为流可能为空。
  • 并行归约要求操作满足结合律,并正确设计 identity 与 combiner。

八、惰性执行与短路

Optional<String> first = languages.stream()
    .filter(name -> {
        System.out.println("filter: " + name);
        return name.length() > 2;
    })
    .map(name -> {
        System.out.println("map: " + name);
        return name.toUpperCase();
    })
    .findFirst();

findFirst 找到第一个满足条件的元素后即可结束,因此后续元素可能完全不会执行 filter 或 map。不要依赖中间操作中的副作用一定发生。


九、并行流的边界

long count = users.parallelStream()
    .filter(User::active)
    .map(User::id)
    .distinct()
    .count();
  • 适合数据量足够大、计算开销较高、元素可独立处理的 CPU 密集型任务。
  • 数据拆分、线程调度、合并结果都有成本,小集合通常不会更快。
  • 阻塞 I/O、共享可变状态、严格顺序要求通常会削弱并行收益。
  • parallelStream 默认使用公共 ForkJoinPool,可能与应用其他任务竞争资源。
  • 是否使用并行流应通过基准测试和生产负载验证,而不是凭感觉决定。

十、无干扰、无状态与副作用

// 不推荐:在管道中修改外部可变集合
List<String> target = new ArrayList<>();
languages.parallelStream().forEach(target::add);

// 推荐:使用归约或收集器返回结果
List<String> safeResult = languages.parallelStream()
    .map(String::toUpperCase)
    .toList();
  • 不要在管道执行期间修改非并发数据源。
  • Lambda 应尽量无状态,结果只依赖当前输入。
  • 需要聚合时优先使用 reduce 或 collect,不要手动修改共享容器。
  • peek 主要用于调试,不应承载关键业务副作用。

十一、常见误区

  • Stream 不会自动修改原集合,除非代码显式产生副作用。
  • Stream 只能消费一次;需要再次处理时重新从数据源创建。
  • Stream.toList() 的结果不可修改。
  • 并行流不保证更快,也不自动让不安全代码变得线程安全。
  • 无限流必须搭配 limit、findFirst 等短路机制,否则终止操作可能永不结束。
  • 可读性较差的超长管道应拆分为命名方法或中间变量。

十二、快速记忆

source.stream()
    .filter(...)      // 过滤
    .map(...)         // 转换
    .flatMap(...)     // 展开
    .distinct()       // 去重
    .sorted(...)      // 排序
    .limit(...)       // 截取
    .toList();        // 终止并收集
先保证管道正确、无副作用且易读,再考虑性能;只有基准数据证明顺序流是瓶颈时,才评估并行化。