Stream API 从 Java 8 引入,声明式的写法很容易让人把它当成「换一套语法写循环」。
它自带一套求值模型(惰性)、一套并发模型(ForkJoinPool)和一组明确契约(identity、结合律、有序性、一次性消费),违反契约的代码常常在单测里跑得好好的,换个数据量或换成并行流才炸。

验证环境

文中所有代码都在下面的环境编译并运行,javac -Xlint:all 无告警:

1
2
3
java version "25" 2025-09-16 LTS
Java(TM) SE Runtime Environment (build 25+37-LTS-3491)
javac 25

涉及版本能力差异的结论一律标注 JDK 版本与 JEP 编号;纯 API 增补(没有对应 JEP 的)以 javadoc 里的 @since 为准。文中出现的实测数字都是本机(Apple M1 Pro)单次运行结果,用途是展示量级差异,不是基准测试。

本文使用的数据模型

1
2
3
4
5
6
7
8
9
record OrderLine(String sku, int qty) {}

record Order(String id, String customer, String city, int amount, List<OrderLine> lines) {}

record Reading(long millis, int kelvins) {
boolean jumpOver30(Reading next) {
return next.kelvins() > kelvins + 30 || next.kelvins() < kelvins - 30;
}
}

Order.customer 允许为 null,本文关于 null 键、null 值的例子都落在它上面;其余字段都是非空的基本类型或集合。

误用一:一个 Stream 只能消费一次

错误写法

1
2
3
Stream<String> upper = List.of("a", "bb", "ccc").stream().map(String::toUpperCase);
System.out.println("[1] count = " + upper.count());
System.out.println("[2] first = " + upper.findFirst().orElseThrow());
1
2
[1] count = 3
[2] java.lang.IllegalStateException: stream has already been operated upon or closed

为什么错

Stream 不是集合,它是「一次性管道」。JDK 的实现里 AbstractPipeline 有一个 linkedOrConsumed 标志位,第一次调用终结操作(或 iterator()spliterator()onClose())就置位,之后任何操作都直接抛异常:

1
private static final String MSG_STREAM_LINKED = "stream has already been operated upon or closed";

第一次 count() 成功并不矛盾——管道被消费了,只是一次消费是合法的。这类代码通常来自「想把管道当查询对象存起来复用」的直觉,但流不支持这种用法。

正确写法:把「创建管道」包成 Supplier

1
2
3
Supplier<Stream<String>> upper = () -> List.of("a", "bb", "ccc").stream().map(String::toUpperCase);
System.out.println("[3] count = " + upper.get().count());
System.out.println("[4] first = " + upper.get().findFirst().orElseThrow());
1
2
[3] count = 3
[4] first = A

只有「管道构造逻辑本身需要复用」时才值得用 Supplier<Stream<T>>,而它每次 get() 都是重新建管道;如果数据源昂贵(磁盘、网络),缓存 List 再反复开流更合适。

同一个坑的几个变体:forEach 之后又 count();把 Stream 存成字段期望多次读取;方法返回 Stream<T> 却没说明「只能消费一次」,调用方先打印日志再返回就会炸。

误用二:惰性求值——没有终结操作就什么都不发生

中间操作不会自己执行

1
2
3
4
5
6
Stream<String> pipeline = List.of("a", "b", "c").stream()
.peek(s -> System.out.println(" peek: " + s))
.map(String::toUpperCase);
System.out.println("[1] 管道构建完成,上面的 peek 一行都不会打印");
List<String> result = pipeline.collect(Collectors.toList());
System.out.println("[2] 终结操作触发求值 = " + result);
1
2
3
4
5
[1] 管道构建完成,上面的 peek 一行都不会打印
peek: a
peek: b
peek: c
[2] 终结操作触发求值 = [A, B, C]

中间操作只是往管道上挂了一个 Sink,只有终结操作(collectforEachcountreducefindFirst……)才会驱动整条链真的跑一遍。所以「构造了流却没消费」等于完全没执行——peek 里写日志、map 里调接口,都会静默消失。

peek 会被 count() 优化掉

1
2
3
4
long n = List.of("a", "b", "c").stream()
.peek(s -> System.out.println(" peek(被跳过): " + s))
.count();
System.out.println("[3] count = " + n);
1
[3] count = 3

三行 peek 一行都没打印。这不是 bug:Stream#count 的 apiNote 明确允许这种跳过:

An implementation may choose to not execute the stream pipeline … if it is capable of computing the count directly from the stream source. In such cases no source elements will be traversed and no intermediate operations will be evaluated.

源是 ListSIZED),中间只挂了 peek(不改变元素个数),实现就可以跳过整条管道。一旦链上出现会改变数量的操作,就必须真的跑:

1
2
3
4
5
long n = List.of("a", "bbb", "c").stream()
.filter(s -> s.length() > 1)
.peek(s -> System.out.println(" peek(filter 之后): " + s))
.count();
System.out.println("[4] count = " + n);
1
2
  peek(filter 之后): bbb
[4] count = 1

结论

peek 的定位是调试,不是业务副作用。任何「靠 peek 干活」的代码都在赌实现不会优化它,业务副作用应当用 forEach 或普通循环。

误用三:Collectors.toMap 的三种翻车方式

1
2
3
4
5
6
7
8
record OrderLine(String sku, int qty) {}
record Order(String id, String customer, String city, int amount, List<OrderLine> lines) {}

static final List<Order> ORDERS = List.of(
new Order("A-1", "alice", "Beijing", 120, List.of()),
new Order("A-2", "alice", "Beijing", 80, List.of()),
new Order("B-1", "bob", "Shanghai", 200, List.of()),
new Order("C-1", null, "Beijing", 60, List.of()));

key 重复 → IllegalStateException

1
2
Map<String, Integer> byCustomer = ORDERS.stream()
.collect(Collectors.toMap(Order::customer, Order::amount));
1
java.lang.IllegalStateException: Duplicate key alice (attempted merging values 120 and 80)

两参数的 toMap 内部是 uniqKeysMapAccumulator,对每个元素做 map.putIfAbsent(key, value),key 已存在就抛异常。这是好设计,但很多人默认它像 put

value 为 null → NPE,加 merge 函数也救不了

1
2
Map<String, String> byId = List.of(new Order("A-1", null, "Beijing", 120, List.of())).stream()
.collect(Collectors.toMap(Order::id, Order::customer));
1
java.lang.NullPointerException(valueMapper 返回 null)

两参数版本的累加器第一句就是 V v = Objects.requireNonNull(valueMapper.apply(element));

改用三参数版本(带 merge 函数):

1
2
Map<String, String> byId = List.of(new Order("A-1", null, "Beijing", 120, List.of())).stream()
.collect(Collectors.toMap(Order::id, Order::customer, (a, b) -> a));
1
三参数版本仍然 java.lang.NullPointerException(HashMap.merge 不接受 null value)

三参数版本内部走 map.merge(key, value, mergeFunction),而 HashMap.merge 要求 value 非 null,所以照样 NPE。merge 函数只解决 key 冲突,不解决 null 值

key 为 null → 取决于容器

默认的 HashMap::new 允许 null key,所以 toMap 遇到 null key 不会抛;但 Collectors.groupingBy 会(见误用十),Collectors.toConcurrentMapConcurrentHashMap 承载,也会抛。同一个 null 在不同收集器上表现不同,是线上最容易漏测的一类。

正确写法:merge 函数 + mapFactory

1
2
3
4
5
6
7
8
Map<String, Integer> merged = ORDERS.stream()
.filter(o -> o.customer() != null)
.collect(Collectors.toMap(
Order::customer,
Order::amount,
Integer::sum,
LinkedHashMap::new));
System.out.println("[4] " + merged + " -> " + merged.getClass().getSimpleName());
1
[4] {alice=200, bob=200} -> LinkedHashMap

四个要点:

  • 第三个参数决定冲突策略,Integer::sum(a, b) -> a(a, b) -> b 都是常见选择,语义要写清楚;
  • 第四个参数换成 LinkedHashMap::new 能保住遇到顺序,默认的 HashMap 不保证迭代顺序:
1
2
3
Map<String, Integer> byId = ORDERS.stream()
.collect(Collectors.toMap(Order::id, Order::amount, Integer::sum));
System.out.println("[5] " + byId);
1
[5] {A-1=120, C-1=60, B-1=200, A-2=80}
  • null 值必须在进 map 之前处理:filter(o -> o.customer() != null) 或者先映射成默认值;
  • 想要「保留最后一个值」这种覆盖语义,就显式写 (a, b) -> b

误用四:map 与 flatMap 的边界

map 会保留嵌套

1
2
List<List<OrderLine>> nested = ORDERS.stream().map(Order::lines).toList();
System.out.println("[1] 元素个数 = " + nested.size() + ",每个元素仍是 List<OrderLine>");
1
[1] 元素个数 = 2,每个元素仍是 List<OrderLine>

map 是一进一出,输入 Order 输出 List<OrderLine>,结果自然是 Stream<List<OrderLine>>。接着对它 count()distinct() 都是在 List 层面操作,跟里面的商品无关。

flatMap 才展平

1
2
3
4
List<OrderLine> lines = ORDERS.stream()
.flatMap(o -> o.lines().stream())
.toList();
System.out.println("[2] 元素个数 = " + lines.size());
1
[2] 元素个数 = 3

判断标准很简单:mapper 返回的是「一个元素」还是「一堆元素」。想把 List/数组/Optional 里的内容合并进主流,就用 flatMap

真正抛 NPE 的地方:嵌套集合本身是 null

1
2
3
4
5
6
List<Order> orders = List.of(
new Order("A-1", "alice", "Beijing", 120, List.of(new OrderLine("apple", 2))),
new Order("A-2", "bob", "Beijing", 80, null));
List<OrderLine> lines = orders.stream()
.flatMap(o -> o.lines().stream()) // 对 null 调 .stream()
.toList();
1
java.lang.NullPointerException(A-2 的 lines 是 null,直接解引用)

抛 NPE 的原因是 o.lines() 为 null 之后又被解引用,跟 flatMap 本身无关。

flatMap 的 mapper 返回 null 反而不会抛

1
2
3
4
List<String> result = Stream.of("a", "b", "c")
.flatMap(s -> s.equals("b") ? null : Stream.of(s))
.toList();
System.out.println("[5] mapper 返回 null 不抛异常,元素被静默丢弃 = " + result);
1
[5] mapper 返回 null 不抛异常,元素被静默丢弃 = [a, c]

JDK 25 的 ReferencePipeline.flatMap 实现里对结果做了判空:

1
2
3
4
5
try (Stream<? extends R> result = mapper.apply(e)) {
if (result != null) {
// 把展开后的元素传给下游
}
}

返回 null 不是错误,表示「这个元素没有展开结果」。很多人以为它会 NPE,跑出来的结果是少元素——一个更隐蔽的问题,因为不会报错。

修法:Stream.ofNullable(JDK 9)

1
2
3
4
5
List<OrderLine> lines = orders.stream()
.flatMap(o -> Stream.ofNullable(o.lines()))
.flatMap(List::stream)
.toList();
System.out.println("[4] " + lines.size());
1
[4] 1

Stream.ofNullable(T) 自 JDK 9 引入(@since 9):null 返回空流,非 null 返回单元素流。等价写法是先 filter(Objects::nonNull)(同样 @since 9),两者都能用:

1
2
3
4
5
6
List<String> source = new ArrayList<>();
source.add("a");
source.add(null);
source.add("b");
System.out.println("[6] " + source.stream().flatMap(Stream::ofNullable).toList());
System.out.println("[7] " + source.stream().filter(Objects::nonNull).toList());
1
2
[6] [a, b]
[7] [a, b]

误用五:短路语义与无限流

findFirst 与 findAny 的差别不是性能,是稳定性

1
2
3
4
5
6
7
OptionalInt any = IntStream.range(0, 1_000_000).parallel()
.filter(n -> n % 100_000 == 7)
.findAny();
OptionalInt first = IntStream.range(0, 1_000_000).parallel()
.filter(n -> n % 100_000 == 7)
.findFirst();
System.out.println("[4] findAny = " + any.getAsInt() + ", findFirst = " + first.getAsInt());
1
[4] findAny = 300007, findFirst = 7

findAny 的 javadoc 写得很直接:并行下多次调用「may not return the same result」,要稳定结果就用 findFirst。反过来说,只要「随便给我一个满足条件的元素」,findAny 允许实现拿到就返回,不必等前面分片。

匹配操作都是短路的,所以无限流也能用

1
2
3
4
boolean anyEven = IntStream.range(1, Integer.MAX_VALUE).anyMatch(n -> n % 2 == 0);
boolean allPositive = IntStream.range(1, Integer.MAX_VALUE).allMatch(n -> n > 0);
boolean noneNegative = IntStream.range(1, Integer.MAX_VALUE).noneMatch(n -> n < 0);
System.out.println("[3] anyMatch=" + anyEven + ", allMatch=" + allPositive + ", noneMatch=" + noneNegative);
1
[3] anyMatch=true, allMatch=true, noneMatch=true

二十多亿个元素理论上跑不完,但三个方法都在前几个元素上得出结论并停止上游。filter + findFirst 同理:

1
2
3
4
5
6
7
static boolean isPrime(long n) {
if (n < 2) return false;
for (long i = 2; i * i <= n; i++) {
if (n % i == 0) return false;
}
return true;
}
1
2
3
4
5
int firstPrimeAbove1000 = IntStream.iterate(1001, n -> n + 1)
.filter(n -> isPrime(n))
.findFirst()
.orElseThrow();
System.out.println("[1] 1000 之后的第一个素数 = " + firstPrimeAbove1000);
1
[1] 1000 之后的第一个素数 = 1009

Stream.iterate(seed, f) 生成的是无限流,但 findFirst 是短路终结操作,管道拿到结果就不再向上游要元素。

takeWhile / dropWhile(JDK 9)

1
2
3
List<Integer> prefix = Stream.of(2, 4, 6, 7, 8, 10).takeWhile(n -> n % 2 == 0).toList();
List<Integer> dropped = Stream.of(2, 4, 6, 7, 8, 10).dropWhile(n -> n % 2 == 0).toList();
System.out.println("[2] takeWhile = " + prefix + ", dropWhile = " + dropped);
1
[2] takeWhile = [2, 4, 6], dropWhile = [7, 8, 10]

takeWhile 遇到第一个不满足条件的元素就整体停止;dropWhile 丢掉开头连续满足条件的部分,剩下的全部保留——这也意味着 dropWhile 之后不再可能短路。

无限流 + 非短路中间操作 = 永不返回

1
2
3
4
5
6
// 反例,不要运行:sorted() 是 stateful 操作,必须先把上游全部元素读进内存
long first = Stream.iterate(2L, n -> n + 1)
.filter(n -> isPrime(n))
.sorted()
.findFirst()
.orElseThrow();

filterfindFirst 都允许短路,但中间夹了 sorted():排序必须看到全部元素才能产出第一个,而全量元素永远读不完,所以它会一直转下去。同类操作还有 distinct()、以及任何需要在终结前缓冲全部元素的操作。

正确做法是在被消费前把「无限」变成「有限」,或者去掉不必要的 sorted

1
2
3
4
long first = Stream.iterate(2L, n -> n + 1)
.filter(n -> isPrime(n))
.findFirst()
.orElseThrow();

误用六:parallelStream 的四个坑

坑一:小数据量下并行只增加开销

1
2
3
4
5
6
7
8
List<Integer> numbers = IntStream.range(0, 1_000).boxed().toList();
long t0 = System.nanoTime();
long seq = numbers.stream().mapToInt(n -> n * n).sum();
long t1 = System.nanoTime();
long par = numbers.parallelStream().mapToInt(n -> n * n).sum();
long t2 = System.nanoTime();
System.out.printf("[1] 顺序 %d (%d µs) / 并行 %d (%d µs)%n",
seq, (t1 - t0) / 1_000, par, (t2 - t1) / 1_000);
1
[1] 顺序 332833500 (480 µs) / 并行 332833500 (515 µs)

多次运行的顺序/并行耗时分别是 431/524480/515993/1668 微秒,并行一次都没赢。并行流要付的固定成本包括:切分数据源、提交到线程池、每个分片各建一套 Sink、最后合并结果。元素本身只做一次乘法时,这些成本远大于计算本身。

坑二:有状态 lambda

1
2
3
List<Integer> acc = new ArrayList<>();
IntStream.range(0, 100_000).parallel().forEach(acc::add);
System.out.println("[2] 期望 100000,实际 " + acc.size());
1
[2] 期望 100000,实际 13475

ArrayList 不是线程安全的,addsize 递增与写入数组之间存在竞态,最终只留下一部分元素,而且不报错。同样的代码在顺序流里每次都是 100000。另一种表现是元素重复:两个线程读到相同的 size,写到同一个槽位。

修法:能用 collect 就别用 forEach(框架保证每个分片有独立容器);确实要并发累加,用 LongAdderConcurrentHashMap 这类并发容器。

坑三:forEach 没有顺序,累加容器还可能被写坏

1
2
3
4
5
6
7
StringBuilder sb = new StringBuilder();
IntStream.range(0, 8).parallel().forEach(i -> sb.append(i).append(','));
System.out.println("[3] forEach 副作用 = " + sb);
String ordered = IntStream.range(0, 8).parallel()
.mapToObj(String::valueOf)
.collect(Collectors.joining(",", "[", "]"));
System.out.println("[4] collect(joining) = " + ordered);
1
2
[3] forEach 副作用 = 5,2,3,6,0,7,1,4,
[4] collect(joining) = [0,1,2,3,4,5,6,7]

[3] 有两个问题叠在一起:并行 forEach 不保证顺序(输出是 5,2,3,...),StringBuilder 也不是线程安全的。collect 会按顺序合并各分片的结果,所以 [4] 永远有序。

确实需要按遇到顺序执行副作用就用 forEachOrdered;并行、需要顺序、有副作用三者同时出现,基本说明这里不该并行。

坑四:commonPool 是全局共享的,阻塞任务会拖垮别人

1
2
3
4
5
6
7
8
static String blockingCall(String id) {
try {
Thread.sleep(200);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return id;
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
System.out.println("[0] commonPool 并行度 = " + ForkJoinPool.getCommonPoolParallelism());

long t0 = System.nanoTime();
List<String> viaCommonPool = IntStream.range(0, 20)
.parallel()
.mapToObj(i -> blockingCall("task-" + i))
.toList();
long t1 = System.nanoTime();
try (ForkJoinPool pool = new ForkJoinPool(20)) {
List<String> viaOwnPool = pool.submit(() -> IntStream.range(0, 20)
.parallel()
.mapToObj(i -> blockingCall("task-" + i))
.toList()).join();
long t2 = System.nanoTime();
System.out.printf("[5] commonPool 里 20 次 200ms 阻塞 = %d ms(并行度 %d),专用池 = %d ms,结果条数 %d/%d%n",
(t1 - t0) / 1_000_000, ForkJoinPool.getCommonPoolParallelism(),
(t2 - t1) / 1_000_000, viaCommonPool.size(), viaOwnPool.size());
}
1
2
[0] commonPool 并行度 = 9
[5] commonPool 里 20 次 200ms 阻塞 = 411 ms(并行度 9),专用池 = 209 ms,结果条数 20/20

parallelStream() 默认使用 ForkJoinPool.commonPool(),并行度是 Runtime.getRuntime().availableProcessors() - 1(JDK 25 源码里就是 Math.max(1, Runtime.getRuntime().availableProcessors() - 1))。可以通过系统属性 java.util.concurrent.ForkJoinPool.common.parallelism 调整,但它的 javadoc 明确写了「Usage is discouraged」;ForkJoinPool.setParallelism(int) 自 JDK 19 起提供(@since 19)。

这个池是整个 JVM 共用的,不为单个流服务。因此:

  • 阻塞任务(IO、Thread.sleep、等锁)放进 commonPool 会占住工作线程,把同进程其它并行流一起拖慢;
  • 上例中 20 个 200ms 的任务在 commonPool 上花了 411ms,换专用池 207ms;
  • 更糟的是常驻阻塞任务把 commonPool 占满后,其它并行流会退化到接近串行,且很难定位。

修法是把阻塞任务放进专用池,或用 CompletableFuture 配合显式 Executor——后者的线程归属更明确:

1
2
3
4
5
6
try (ForkJoinPool pool = new ForkJoinPool(20)) {
List<String> ids = pool.submit(() -> IntStream.range(0, 20)
.parallel()
.mapToObj(i -> blockingCall("task-" + i))
.toList()).join();
}

什么时候用 groupingByConcurrent

groupingByConcurrent 返回 ConcurrentMap,是无序收集器,并行时不必为合并每个 key 的中间 Map 付出代价,代价是结果不保序:

1
2
3
4
5
6
Map<String, List<String>> g = ORDERS.parallelStream()
.collect(Collectors.groupingBy(Order::city, Collectors.mapping(Order::id, Collectors.toList())));
Map<String, List<String>> c = ORDERS.parallelStream()
.collect(Collectors.groupingByConcurrent(Order::city, Collectors.mapping(Order::id, Collectors.toList())));
System.out.println("[6] groupingBy 保序 = " + g.get("bj"));
System.out.println("[7] groupingByConcurrent 不保序 = " + c.get("bj"));
1
2
[6] groupingBy 保序 = [A-0, A-3, A-6, A-9, A-12, A-15, A-18]
[7] groupingByConcurrent 不保序 = [A-0, A-9, A-3, A-18, A-15, A-6, A-12]

另外,groupingByConcurrent 的源码对「下游收集器不是并发安全」的情况加了 synchronized (resultContainer),此时并发收益会明显缩水。而且并行收集本身未必更快——十万个元素按 8 个 key 分组,顺序 groupingBy 18ms,并行 22ms,合并 Map 的成本盖过了并行遍历的收益。

判断标准

该用并行需要四条同时成立:数据量足够大、单个元素的计算足够重、lambda 无共享可变状态、结果不依赖顺序。任何一条不成立,顺序流或普通循环都更划算。至于并行度,别去调 commonPool 的属性,把阻塞任务放进专用池才有用。

误用七:装箱与原始类型流

开销在「装箱」这个动作本身

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
int n = 10_000_000;
long raw = 0;
long boxed = 0;
long preBoxed = 0;
for (int round = 0; round < 4; round++) {
long t0 = System.nanoTime();
int a = IntStream.range(0, n).sum();
long t1 = System.nanoTime();
int b = IntStream.range(0, n).boxed().mapToInt(Integer::intValue).sum();
long t2 = System.nanoTime();
List<Integer> list = IntStream.range(0, n).boxed().toList();
long t3 = System.nanoTime();
int c = list.stream().mapToInt(Integer::intValue).sum();
long t4 = System.nanoTime();
if (a != b || b != c) throw new AssertionError("三条路径结果应当一致");
if (round > 0) {
raw += t1 - t0;
boxed += t2 - t1;
preBoxed += t4 - t3;
}
}
System.out.printf("[1] 1000 万次求和:IntStream.sum() %.1f ms / boxed 再拆箱 %.1f ms / 已是 List<Integer> %.1f ms%n",
raw / 3_000_000.0, boxed / 3_000_000.0, preBoxed / 3_000_000.0);

预热后取三次平均:

1
[1] 1000 万次求和:IntStream.sum() 3.5 ms / boxed 再拆箱 53.4 ms / 已是 List<Integer> 9.1 ms

三个数字要分开看:

  • 全程原始类型,不分配对象,3.5ms;
  • 中间插一次 boxed(),等于额外创建一千万个 Integer 再拆回来,53.4ms;
  • 已经是 List<Integer> 时不再重复装箱,只多一次拆箱和对象访问,9.1ms。

所以「用原始类型流」的收益主要来自避免分配,而不是省掉一次虚方法调用。数值聚合(sumaveragemax)、大数组处理、range 生成序列都适合 IntStream;一旦需要集合语义(distinct、带对象比较器的 sortedList<Integer> 结果)就绕不开装箱,此时用 IntStream 只是把装箱点推迟。

sum()average() 的返回值语义不同

1
2
3
4
System.out.println("[2] 空流 sum() = " + IntStream.empty().sum());
OptionalDouble empty = IntStream.empty().average();
System.out.println("[3] 空流 average() = " + empty + ", isPresent = " + empty.isPresent());
System.out.println("[4] 非空 average() = " + IntStream.of(1, 2, 4).average().orElseThrow());
1
2
3
[2] 空流 sum() = 0
[3] 空流 average() = OptionalDouble.empty, isPresent = false
[4] 非空 average() = 2.3333333333333335

sum() 在空流上返回单位元 0,这是合法结果,所以没有 Optional 包装。average() 在空流上没有定义,返回 OptionalDouble.empty()——不是 0,也不是 NaN。把它当 0 处理会静默产出错误的业务数据(「平均耗时 0 毫秒」通常意味着数据丢了)。

sum() 会静默溢出

1
2
3
int overflowed = IntStream.of(Integer.MAX_VALUE, 1).sum();
long exact = IntStream.of(Integer.MAX_VALUE, 1).asLongStream().sum();
System.out.println("[5] int sum = " + overflowed + ", long sum = " + exact);
1
[5] int sum = -2147483648, long sum = 2147483648

IntStream.sum() 返回 int,溢出不抛异常。金额、计数、字节数这类会变大的量,用 LongStreammapToLongasLongStream)或 BigDecimal

还有一个类型陷阱:Collectors.joining 的元素类型是 CharSequenceIntStream.rangeClosed(1, 10).boxed() 得到 Stream<Integer>,直接收集编译不过,必须先 mapToObj(String::valueOf)

1
2
3
System.out.println("[6] " + IntStream.rangeClosed(1, 10)
.mapToObj(String::valueOf)
.collect(Collectors.joining(",")));
1
[6] 1,2,3,4,5,6,7,8,9,10

误用八:拿 reduce 做字符串拼接与集合累加

拼接字符串:reduce + String::concat 是 O(n²)

1
2
3
4
5
6
7
8
9
10
List<String> words = IntStream.range(0, 20_000).mapToObj(i -> "w" + i + ";").toList();
long t0 = System.nanoTime();
String viaReduce = words.stream().reduce("", String::concat);
long t1 = System.nanoTime();
String viaJoining = words.stream().collect(Collectors.joining());
long t2 = System.nanoTime();
System.out.printf("[1] reduce+concat %d 字符 %d ms / joining %d 字符 %d ms / 结果一致 %b%n",
viaReduce.length(), (t1 - t0) / 1_000_000,
viaJoining.length(), (t2 - t1) / 1_000_000,
viaReduce.equals(viaJoining));
1
[1] reduce+concat 128890 字符 88 ms / joining 128890 字符 3 ms / 结果一致 true

String::concat 每次都新建字符串并把左边整体复制过去,2 万次拼接就是 2 万次全量拷贝。Collectors.joining() 内部用 StringBuilder,一次分配搞定。需要分隔符和前后缀时用 joining(delimiter, prefix, suffix)

1
System.out.println("[4] " + words.stream().collect(Collectors.joining(", ", "<", ">")));
1
[4] <alpha, beta, gamma>

可变 identity:并行流下结果直接翻车

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
List<Integer> seq = IntStream.range(0, 10).boxed()
.reduce(new ArrayList<>(), (list, i) -> {
list.add(i);
return list;
}, (a, b) -> {
a.addAll(b);
return a;
});
List<Integer> par = IntStream.range(0, 10).boxed()
.parallel()
.reduce(new ArrayList<>(), (list, i) -> {
list.add(i);
return list;
}, (a, b) -> {
a.addAll(b);
return a;
});
System.out.println("[2] 顺序流 size = " + seq.size() + ", 并行流 size = " + par.size());
1
[2] 顺序流 size = 10, 并行流 size = 1808

顺序流看起来完全正确,所以这段代码很容易被复制到并行场景。但 reduce 要求 identity 是单位元combiner(identity, u) 必须等于 u。上面这个 new ArrayList<>() 不是单位元——各分片以它为起点,合并时把内容追加进去,元素被重复计入,10 个变成了 1808 个。同一个可变 identity 还会导致多个分片共享同一个 ArrayList 实例(identity 只创建一次),并发写入时行为未定义。

非结合操作:并行结果由分片方式决定

1
2
3
4
int sequential = IntStream.rangeClosed(1, 100).reduce(0, (a, b) -> a - b);
int parallel = IntStream.rangeClosed(1, 100).parallel().reduce(0, (a, b) -> a - b);
System.out.println("[3] 减法不是结合律:顺序流 = " + sequential + ",并行流 = " + parallel
+ "(并行结果由分片方式决定,契约上未定义)");
1
[3] 减法不是结合律:顺序流 = -5050,并行流 = 0(并行结果由分片方式决定,契约上未定义)

reduce 还要求累加函数满足结合律(a op (b op c) == (a op b) op c)。减法、除法、字符串「前插」都不满足:顺序流给出 -5050,并行流按分片归并给出 0——契约上这是未定义行为,给不出正确结果,但不会报错

正确写法

需求 该用的收集器/方法
拼接字符串 Collectors.joining() / String.join / StringBuilder
收集到 List Stream.toList()(JDK 16,不可变)或 Collectors.toCollection(ArrayList::new)
收集到 Set/Map Collectors.toSet() / Collectors.toMap(...)
数值聚合 IntStream.sum()/average()/summaryStatistics()
需要自定义容器 collect(Supplier, BiConsumer, BiConsumer)
1
2
3
4
5
System.out.println("[5] " + IntStream.range(0, 5).boxed()
.collect(Collectors.toCollection(ArrayList::new)).getClass().getSimpleName());
String built = IntStream.range(0, 5).collect(StringBuilder::new,
(sb, i) -> sb.append(i), StringBuilder::append).toString();
System.out.println("[6] " + built);
1
2
[5] ArrayList
[6] 01234

最后这个 collect(Supplier, accumulator, combiner)可变归约:契约是「把元素并入已有容器」,每个分片各自调用一次 supplier 拿到独立容器,所以并行安全。这和 reduce 的「不可变累加」是两套 API,混用就是上面那些事故的来源。

误用九:sorted 与 Comparator

1
2
3
4
5
6
static final List<Order> ORDERS = List.of(
new Order("A-1", "alice", "Beijing", 120, List.of()),
new Order("A-2", "bob", "Beijing", 80, List.of()),
new Order("A-3", "carol", "Shanghai", 80, List.of()),
new Order("A-4", "dave", "Beijing", 120, List.of()),
new Order("A-5", null, "Beijing", 60, List.of()));

稳定性

1
2
3
4
5
List<String> byCity = ORDERS.stream()
.sorted(Comparator.comparing(Order::city))
.map(Order::id)
.toList();
System.out.println("[1] 按城市排序(同城市保持原有顺序)= " + byCity);
1
[1] 按城市排序(同城市保持原有顺序)= [A-1, A-2, A-4, A-5, A-3]

Stream#sorted 的 javadoc:For ordered streams, the sort is stable. 四个 Beijing 保持了它们在源里的相对顺序。多级排序里「主键相等时保持原有次序」是可以依赖的语义(前提是源有序,比如从 List 开的流)。

链式比较器

1
2
3
4
Comparator<Order> byCityThenAmountThenId = Comparator.comparing(Order::city)
.thenComparingInt(Order::amount)
.thenComparing(Order::id);
System.out.println("[2] " + ORDERS.stream().sorted(byCityThenAmountThenId).map(Order::id).toList());
1
[2] [A-5, A-2, A-1, A-4, A-3]

comparing 的第一个键相等时才走 thenComparing,依次往下。三个细节:

  • 数值键用 thenComparingInt/Long/Double,避免装箱和二级函数调用;
  • 主键比较器如果自己写 compare,用 Integer.compare(a, b) 而不是 a - b(后者会溢出);
  • thenComparing(Order::id) 这类「最后的兜底键」能让结果完全确定,调试时省事。

null 键

1
System.out.println("[3] " + ORDERS.stream().sorted(Comparator.comparing(Order::customer)).map(Order::id).toList());
1
[3] java.lang.NullPointerException(comparing 对 null 键调用 compareTo)

Comparator.comparing(keyExtractor) 生成的比较器直接调 key.compareTo(otherKey),遇到 null 键就 NPE。用 nullsFirst / nullsLast 把 null 显式安排到一端:

1
2
3
4
5
List<String> nullSafe = ORDERS.stream()
.sorted(Comparator.comparing(Order::customer, Comparator.nullsLast(Comparator.naturalOrder())))
.map(Order::id)
.toList();
System.out.println("[4] nullsLast = " + nullSafe);
1
[4] nullsLast = [A-1, A-2, A-3, A-4, A-5]

注意两种写法的层次差别:comparing(keyExtractor, keyComparator) 是把 null 检查加在上;nullsLast(comparator) 是把 null 检查加在整个比较器外层。要按多个字段排且其中多个字段可能为 null,就得给每个键分别套一层。

反序与稳定性的组合

1
2
3
4
System.out.println("[5] 金额降序(相同金额仍是稳定次序)= " + ORDERS.stream()
.sorted(Comparator.comparingInt(Order::amount).reversed())
.map(Order::id)
.toList());
1
[5] 金额降序(相同金额仍是稳定次序)= [A-1, A-4, A-2, A-3, A-5]

reversed() 反转的是比较结果,元素「相等」的事实不变,所以两个 120(A-1、A-4)和两个 80(A-2、A-3)都保持了源里的先后。

最后提醒性能:sorted() 是 stateful 中间操作,必须把上游全部元素缓冲进内存才能产出第一个元素。sorted().findFirst() 不比 sorted().toList().get(0) 省事,sorted() 之后接 limit(1) 也短路不了。

误用十:partitioningBy 与 groupingBy 的差别

键集合的语义不同

1
2
3
4
5
6
Map<Boolean, List<Order>> partition = ORDERS.stream()
.collect(Collectors.partitioningBy(o -> o.amount() >= 100));
System.out.println("[1] partitioningBy 键集合 = " + partition.keySet());
System.out.println("[2] 两个分区都存在," + partition.get(false).size() + " / " + partition.get(true).size());
Map<String, List<Order>> grouped = ORDERS.stream().collect(Collectors.groupingBy(Order::city));
System.out.println("[3] groupingBy 键集合 = " + grouped.keySet());
1
2
3
[1] partitioningBy 键集合 = [false, true]
[2] 两个分区都存在,2 / 1
[3] groupingBy 键集合 = [Beijing, Shanghai]

partitioningBy 的 javadoc 明确写了:返回的 Map always contains mappings for both false and true keys,某个分区没有元素时它的值是空 List。所以 map.get(true).size() 不需要判空。groupingBy 则只为出现过的 key 建条目——两者的 get() 返回值是否为 null,正好相反。

key 为 null

1
2
Stream.of(ORDERS.get(0), new Order("A-9", "dave", null, 10, List.of()))
.collect(Collectors.groupingBy(Order::city));
1
java.lang.NullPointerException: element cannot be mapped to a null key

分组之前过滤掉,或者把 null 映射成哨兵值。partitioningBy 的谓词返回 boolean,天然没有这个问题。

下游收集器

1
2
3
4
5
6
7
8
9
10
11
Map<String, Integer> amountByCity = ORDERS.stream()
.collect(Collectors.groupingBy(Order::city, Collectors.summingInt(Order::amount)));
Map<String, Long> countByCity = ORDERS.stream()
.collect(Collectors.groupingBy(Order::city, Collectors.counting()));
Map<String, String> idsByCity = ORDERS.stream()
.collect(Collectors.groupingBy(Order::city,
Collectors.mapping(Order::id, Collectors.joining(", "))));
Map<String, Set<String>> skusByCity = ORDERS.stream()
.collect(Collectors.groupingBy(Order::city,
Collectors.flatMapping(o -> o.lines().stream().map(OrderLine::sku), Collectors.toSet())));
System.out.println("[5] " + amountByCity + " / " + countByCity + " / " + idsByCity + " / " + skusByCity);
1
[5] {Beijing=200, Shanghai=80} / {Beijing=2, Shanghai=1} / {Beijing=A-1, A-2, Shanghai=A-3} / {Beijing=[apple, pear], Shanghai=[]}

下游收集器的类型要和需求对齐:counting() 返回 Long 而不是 IntegersummingInt 返回 Integer(可能溢出),averagingInt 返回 Doublemapping 用来在分组之后变换元素形态。JDK 9 起多了 filteringflatMapping 两个下游收集器(都是 @since 9),解决「分组之后再过滤/展平」的场景。

注意这里必须用 flatMapping 而不是先 filtergroupingBy:过滤放在分组之前会改变分组结果(有些分组会整个消失),而 filtering 作为下游收集器时分组仍然存在,只是内容为空。

指定 Map 实现

1
2
3
Map<String, Long> linked = ORDERS.stream()
.collect(Collectors.groupingBy(Order::city, LinkedHashMap::new, Collectors.counting()));
System.out.println("[6] " + linked + " -> " + linked.getClass().getSimpleName());
1
[6] {Beijing=2, Shanghai=1} -> LinkedHashMap

三参数的 groupingBy(classifier, mapFactory, downstream)mapFactory 决定外层 Map 的实现:默认 HashMap,保序用 LinkedHashMap::new,按 key 排序用 TreeMap::new,并发场景考虑 groupingByConcurrent

teeing:一次遍历拿两个结果(JDK 12)

1
2
3
4
5
6
7
8
9
record Stats(long count, double average) {}

Stats stats = ORDERS.stream().collect(Collectors.teeing(
Collectors.counting(),
Collectors.averagingInt(Order::amount),
Stats::new));
System.out.println("[7] teeing = " + stats);
System.out.println("[8] summarizingInt = " + ORDERS.stream()
.collect(Collectors.summarizingInt(Order::amount)));
1
2
[7] teeing = Stats[count=3, average=93.33333333333333]
[8] summarizingInt = IntSummaryStatistics{count=3, sum=280, min=80, average=93.333333, max=120}

Collectors.teeing 自 JDK 12 引入(@since 12,没有独立 JEP),把同一条流交给两个下游收集器再用 merger 合并,只遍历一次。纯统计需求用现成的 summarizingInt/Long/Double 更省事,teeing 的用武之地是「两个语义上无关的聚合」:

1
2
3
4
5
6
7
8
record Range(Optional<Order> cheapest, Optional<Order> priciest) {}

Range range = ORDERS.stream().collect(Collectors.teeing(
Collectors.minBy(Comparator.comparingInt(Order::amount)),
Collectors.maxBy(Comparator.comparingInt(Order::amount)),
Range::new));
System.out.println("[9] " + range.cheapest().orElseThrow().id()
+ " / " + range.priciest().orElseThrow().id());
1
[9] A-2 / A-1

误用十一:这些场景用 for 循环更好

需要「上一个元素」或滑动窗口

1
2
3
4
5
6
7
8
9
10
11
List<Reading> readings = List.of(
new Reading(0, 310), new Reading(1, 312), new Reading(2, 350), new Reading(3, 310));
List<String> suspicious = new ArrayList<>();
for (int i = 1; i < readings.size(); i++) {
Reading prev = readings.get(i - 1);
Reading next = readings.get(i);
if (prev.jumpOver30(next)) {
suspicious.add(prev.millis() + "->" + next.millis());
}
}
System.out.println("[1] 相邻突变 = " + suspicious);
1
[1] 相邻突变 = [1->2, 2->3]

用 Stream 表达「相邻两个元素的比较」需要外部可变状态,或者 IntStream.range 配合 get(i-1),都不比循环清楚。这条限制在 JDK 24 起有了原生解法(见下一节的 Gatherers.windowSliding),更早的版本老老实实写循环。

需要下标

1
2
3
4
5
6
List<Integer> values = List.of(3, -1, 4, -1, 5, -9, 2);
int index = IntStream.range(0, values.size())
.filter(i -> values.get(i) < -5)
.findFirst()
.orElse(-1);
System.out.println("[2] 第一个小于 -5 的下标 = " + index);
1
[2] 第一个小于 -5 的下标 = 5

这段能跑,也能短路,但 IntStream.range(0, size) + get(i) 是把下标循环翻译成流的样子,并没有变得更好读。循环体同时依赖下标和元素时,for 更直接。

受检异常

1
2
3
static String readFirstLine(Path path) throws IOException {
return Files.readAllLines(path).get(0);
}
1
2
3
4
Path file = Path.of("/tmp/example.txt");
List<String> content = List.of(file).stream()
.map(path -> readFirstLine(path)) // 编译不通过
.toList();

javac 的实际报错:

1
错误: 未报告的异常错误IOException; 必须对其进行捕获或声明以便抛出

Functionapply 不允许抛受检异常,这是类型系统层面的限制,不是风格问题。只有两条路:在 lambda 里就地 try/catch 转成非受检异常,或者在循环里直接抛给调用者:

1
2
3
4
5
Path file = Path.of("/tmp/example.txt");
List<String> byLoop = new ArrayList<>();
for (Path path : List.of(file)) {
byLoop.add(readFirstLine(path)); // 直接 throws IOException 给调用者
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
Path file = Path.of("/tmp/example.txt");
try {
List<String> byStream = List.of(file).stream()
.map(path -> {
try {
return readFirstLine(path);
} catch (IOException e) {
throw new UncheckedIOException(e); // 保留原始异常作为 cause
}
})
.toList();
} catch (UncheckedIOException e) {
System.out.println("包装后的异常:" + e.getCause());
}

如果选择包装,至少保留原始异常(UncheckedIOException 或自定义异常带上 cause),不要写 catch (IOException e) { throw new RuntimeException(); } 把栈丢掉。

边界

以下情况直接写循环:需要 break/continue 控制多层流程;需要同时写多个累加器(两个 List、一个计数、一个标志位);循环体里有受检异常且不想包装;需要访问下标或上一个元素且版本低于 JDK 24;循环体超过十来行。

以下情况用 Stream:数据转换与过滤的流水线、分组聚合、集合间映射、以及满足前面四个条件的并行场景。

误用十二:新 API 的正确理解

Stream.toList()(JDK 16)

1
2
3
4
5
6
7
8
9
List<String> immutable = Stream.of("a", "b").toList();
try {
immutable.add("c");
} catch (UnsupportedOperationException e) {
System.out.println("[1] " + e.getClass().getName() + "(Stream.toList 返回不可变 List)");
}
List<String> mutable = Stream.of("a", "b").collect(Collectors.toList());
mutable.add("c");
System.out.println("[2] collect(toList) 可以 add:" + mutable);
1
2
[1] java.lang.UnsupportedOperationException(Stream.toList 返回不可变 List)
[2] collect(toList) 可以 add:[a, b, c]

Stream.toList() 自 JDK 16 引入(@since 16),javadoc 写明「The returned List is unmodifiable」,默认实现是 Collections.unmodifiableList(new ArrayList<>(...)),所以在它上面 add/set/remove 都会抛 UnsupportedOperationException

Collectors.toList() 的 javadoc 恰恰相反:「There are no guarantees on the type, mutability, serializability, or thread-safety of the List returned」——当前实现返回 ArrayList(可以 add),但这是实现细节。只读结果用 toList(),要可变结果用 Collectors.toCollection(ArrayList::new) 显式写出来。同理,JDK 10 起有 Collectors.toUnmodifiableList/Set/Map(都是 @since 10),比 Collections.unmodifiableList(collect(...)) 少一层包装。

mapMulti(JDK 16)

1
2
3
4
5
6
7
List<Integer> expanded = Stream.of(1, 2, 3)
.<Integer>mapMulti((n, downstream) -> {
downstream.accept(n);
downstream.accept(n * 10);
})
.toList();
System.out.println("[3] mapMulti = " + expanded);
1
[3] mapMulti = [1, 10, 2, 20, 3, 30]

mapMulti 也自 JDK 16 引入(@since 16,同批还有 mapMultiToInt/Long/Double)。它和 flatMap 的区别:flatMap 要求 mapper 返回一个 Stream 对象,一进多出时每个元素都要新建流;mapMulti 直接给一个 Consumer,把要输出的元素逐个 accept 进去,没有中间流对象。

1
2
3
4
5
6
7
List<Number> numbers = List.<Number>of(1, 2.5, 3, 4.5);
List<Integer> integers = numbers.stream().<Integer>mapMulti((number, consumer) -> {
if (number instanceof Integer i) {
consumer.accept(i);
}
}).toList();
System.out.println("[9] " + integers);
1
[9] [1, 3]

元素展开很少(0 个或 1 个为主)时 mapMulti 更省;已经在用 flatMap 拼多个现成的流、或者需要惰性处理巨大集合时才用 flatMap

Stream Gatherers(JDK 24 转正)

自定义中间操作。版本线索要写清楚:

  • JEP 461 把它作为预览特性引入 JDK 22;
  • JEP 473 在 JDK 23 二次预览;预览期间需要 --enable-preview
  • JEP 485 在 JDK 24 正式转正,无需任何编译参数。

截至 JDK 25 的现状:java.util.stream.Gatherers 已是可直接在生产使用的正式 API。内置五个:foldscanwindowFixedwindowSlidingmapConcurrent

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
List<List<Integer>> fixed = Stream.of(1, 2, 3, 4, 5)
.gather(Gatherers.windowFixed(2))
.toList();
System.out.println("[4] windowFixed(2) = " + fixed);

List<Integer> running = Stream.of(1, 2, 3, 4)
.gather(Gatherers.scan(() -> 0, Integer::sum))
.toList();
System.out.println("[5] scan 前缀和 = " + running);

List<Integer> sliding = Stream.of(1, 2, 3, 4)
.gather(Gatherers.windowSliding(2))
.map(w -> w.get(0) + w.get(1))
.toList();
System.out.println("[6] windowSliding(2) 相邻和 = " + sliding);

List<String> concurrent = Stream.of("a", "bb", "ccc")
.gather(Gatherers.mapConcurrent(3, String::toUpperCase))
.toList();
System.out.println("[7] mapConcurrent = " + concurrent);

List<Integer> fold = Stream.of(1, 2, 3, 4)
.gather(Gatherers.fold(() -> 0, Integer::sum))
.toList();
System.out.println("[8] fold = " + fold);
1
2
3
4
5
[4] windowFixed(2) = [[1, 2], [3, 4], [5]]
[5] scan 前缀和 = [1, 3, 6, 10]
[6] windowSliding(2) 相邻和 = [3, 5, 7]
[7] mapConcurrent = [A, BB, CCC]
[8] fold = [10]

几个容易记混的点:

  • windowFixed(n) 最后一个窗口可能不满([5]),不会丢弃也不会补 null;
  • windowSliding(n) 输出 List,元素不足一个窗口时不生成;
  • scan 把每一步的中间结果都发下去(前缀和语义),只要最终结果用 fold
  • mapConcurrent(limit, fn) 并发调用 fn 并保持输出顺序,适合每个元素要做一次阻塞 IO 的场景——它是中间操作,比手工把 parallelStream 塞进专用池干净。

Gatherers 与 Collectors 的关系可以这样记:collect 是终结操作的扩展点,gather 是中间操作的扩展点。两者结构相似(initializer/integrator/combiner/finisher 对应 supplier/accumulator/combiner/finisher),区别在于 gatherer 的 integrator 用返回值表示「是否继续处理」,因此支持短路,也能处理无限流。

版本速查

API 版本 依据
Stream.iterate(seed, hasNext, next)takeWhiledropWhileofNullable JDK 9 @since 9
Collectors.filteringflatMapping JDK 9 @since 9
Collectors.toUnmodifiableList/Set/Map JDK 10 @since 10
Collectors.teeing JDK 12 @since 12
Stream.toList()mapMulti / mapMultiToInt JDK 16 @since 16
ForkJoinPool.setParallelism JDK 19 @since 19
Stream Gatherers JDK 22 预览(JEP 461)→ JDK 23 二次预览(JEP 473)→ JDK 24 转正(JEP 485) JEP 文档

总结

这十二类误用背后是四个根因。

把 Stream 当集合用:它是一次性管道,不能存、不能反复遍历,Supplier 只是让「重建管道」变便宜。

忽略惰性与短路规则:中间操作不执行直到遇见终结操作,而 count() 这类终结操作有权跳过整条管道;反过来,sorted() 这类 stateful 操作会让短路失效,遇到无限流就是死循环。

忽略收集器的契约:identity 必须是单位元,累加函数必须满足结合律,toMap 不允许重复 key、不允许 null value,groupingBy 不允许 null key——违反任一条,得到的要么是异常,要么是看起来正常但错误的数字。

把并行当成免费的加速:并行只在数据量大、计算重、无共享状态、不需要顺序时才有收益;阻塞任务放进 commonPool 会波及整个 JVM 的并行流。

至于什么时候不该用 Stream:需要下标、需要上一个元素、需要 break、需要写多个累加器、循环体里有受检异常的时候,直接写循环更清楚。JDK 24 转正的 Gatherers 补上了其中一部分(滑窗、前缀扫描、并发映射),不需要为此改写其余循环。

参考资料

系列索引:Java 系列,语言特性与运行时的长文集