Java · Stream API 与函数式

中间操作只注册不执行, 终端操作一次点火 — 元素逐个穿全链; parallelStream 的 ForkJoinPool.commonPool 是全 JVM 共享的

惰性管道: 中间操作只注册 · 终端操作一次点火 · 逐元素穿全链 ① 管道搭建阶段 — 只是注册节点, 一行都不会执行 数据源 orders.stream() List<Order> 100 万元素躺在原地 filter (注册①) filter(o -> o.paid()) 无状态: 来一个判一个 不碰任何元素 map (注册②) map(Order::amount) 无状态: 逐元素变形 同样只是挂节点 sorted (注册③) sorted(cmp) 有状态: 必须攒全量 前 n-1 个元素出不去 collect · 终端点火 terminal = 拉动整条链 collect / forEach / findFirst collect 调用之前: 整条管道零执行 — 这就是"漏写终端操作, 代码白写"的机制根源 ✗ 常见误解: 逐阶段执行 (水平搬批) 以为先把 100 万元素全部 filter 完, 再整体进入 map filter map sorted collect 元素成批推进 — 不是 Stream 的工作方式 ✓ 真实模型: 逐元素穿全链 (垂直穿链) e1 走完 filter→map→…→collect, e2 才出发 e1 e2 e3 e4 filter map sorted collect findFirst 在 e1 命中 → e2/e3/e4 根本不出发 (短路) parallelStream — 全 JVM 共享一个 ForkJoinPool.commonPool 调用线程 orders.parallelStream() fork 切分 + 提交任务 ForkJoinPool.commonPool (进程级唯一) worker-1 ▸ filter 分片 [CPU 98%] worker-2 ▸ filter 分片 [CPU 97%] worker-3 ▸ ⚠ BLOCKED socketRead 2.3s worker-4 ▸ idle (等待窃取任务) 并行度 = availableProcessors() - 1 (8C → 7) -Djava.util.concurrent.ForkJoinPool.common.parallelism=N 同池竞争者 别的接口的 parallelStream CompletableFuture 默认执行器 也在等这 7 个 worker 🚫 IO 任务霸池 = 全局饿死 并行流里调下游 RPC / 慢 SQL worker 全部 BLOCKED 全 JVM 并行任务集体排队 铁律: CPU 密集才许 parallel 红线: parallel 管道内禁 IO / 锁 / 共享可变状态 — 公共池是全进程并行任务的独木桥, 淹没一个人整条桥塌 Legend 中间操作/无状态 有状态/警示 终端点火/阻塞 共享池 线程/误解

惰性是双刃剑

  • • 中间操作只是往管道挂节点, 没有终端操作 = 整段空转
  • • findFirst/anyMatch/limit 短路: 拿到答案就停
  • • peek 不执行/日志缺失, 十有八九是链尾没终端
  • • 无限流 Stream.generate 能存在, 正因为惰性

逐元素执行模型

  • • e1 穿完 filter→map→collect, 才轮到 e2 出发
  • • sorted/distinct 是异类: 有状态, 必须攒全量
  • • 无状态操作才能被并行流自由切分与合并
  • • 管道里打印元素顺序, 印证的就是这个模型

parallel 红线

  • • commonPool 并行度 = 核数 - 1, 全 JVM 共享
  • • 池内 IO / 锁 / ThreadLocal 三禁
  • • CPU 密集 + 大数据量 + 无状态才值得开
  • • CompletableFuture 默认执行器是同一个池

💡 一句话理解

Stream 的本质是"声明式管道 + 惰性求值": filter/map/sorted 这些中间操作只是往管道上挂节点, 一行都不执行; 直到 collect/findFirst 这样的终端操作点火, 数据才开始流动。而流动方式也不是"第一阶段跑完全部元素再进第二阶段", 而是一个元素穿完整条链, 才轮到下一个元素(sorted 这类有状态操作除外)。这一个模型同时解释了三件事: 为什么漏写终端操作整段代码"白写"; 为什么 findFirst 命中第一个就不再多算一个元素; 为什么 peek 里的日志"有时不执行"。

parallelStream 只是把"切分 + 调度"交给 ForkJoinPool.commonPool — 它是全 JVM 共享的公共资源(默认核数-1 个线程, 工作窃取均衡负载)。任何一个并行流里做慢 IO, 全进程的并行流与 CompletableFuture 默认执行器一起陪葬; 所以业界铁律是 parallel 只留给 CPU 密集、无状态、无 IO 的纯计算管道。

🧠 必知必会 必考 & 必会

中间 vs 终端操作
filter/map/flatMap/sorted/limit/peek 只返回新流并登记自己(惰性, 不碰数据); collect/forEach/reduce/findFirst/count/findAny 是终端: 触发整条链执行并"关闭"流。漏终端 = 什么都不发生。
Stream<Integer> s = list.stream().filter(x -> x > 0); // 只注册
s.count();  // 终端点火, filter 才真正执行
// 关键: 漏终端 = 整段代码白写, 什么都不发生
惰性求值
管道是蓝图不是指令: 终端点火后才逐元素pull。这让 findFirst 能"按需计算", 也让 Stream.generate/iterate 无限流成为可能 — 反正没人要的元素永远不会被造出来。
Stream<Integer> inf = Stream.iterate(1, i -> i + 2); // 无限流
inf.limit(3).toList();  // → [1, 3, 5]
// 关键: 没人要的元素永远不会被造出来 — 这就是无限流能存在的原因
无状态 vs 有状态
filter/map 无状态: 元素独立处理, 来一个算一个; sorted/distinct/limit/skip 有状态: 要攒住前面(甚至全部)元素才能放行 — 它们是并行切分与有序流的瓶颈所在。
orders.stream()
    .filter(o -> o.paid())      // 无状态: 来一个判一个
    .sorted(byAmount)            // 有状态: 必须攒全量才放行
    .limit(10);
// 关键: sorted/distinct 是并行切分与有序流的瓶颈
短路终端
findFirst/findAny/anyMatch/allMatch/noneMatch + limit: 拿到答案就停, 后面的元素根本不进入管道。规则引擎按优先级找首个命中, 靠的就是它。
Optional<Order> hit = orders.stream()
    .filter(o -> o.risk())
    .findFirst();   // e1 命中 → e2 起根本不进管道
// 关键: anyMatch/allMatch/noneMatch + limit 同为短路
逐元素执行模型
e1 穿完 filter→map→…→终端才轮到 e2(垂直穿链), 不是每阶段搬完一批再下一阶段(水平搬批)。理解了它, "peek 打印顺序"和"短路到底省了多少计算"都有答案。
List.of(1, 2).stream()
    .peek(x -> log("F{}", x))   // 打印顺序: F1 M1 F2 M2
    .map(x -> x * 10)
    .toList();                   // → [10, 20]
// 关键: 逐元素垂直穿链, 不是逐阶段水平搬批
无副作用纪律
lambda 只做"输入→输出": 不改外部变量/集合, 不发 IO。干净的管道才能被安全并行与优化; 要"效果"就用 collect 归约或 forEach 显式表达。
// lambda 只做 输入→输出: 不改外部集合/变量, 不发 IO
List<Integer> ok = src.stream().filter(x -> x > 0).toList();
// 关键: 结果用 collect 归约产生; "效果"用 forEach 显式表达
方法引用
Order::getCity / Order::new / System.out::println 是 lambda 的语法糖: 更短、不为实例捕获 this; IDE 提示转方法引用时, 也是"该抽方法了"的信号。
cities.stream()
    .map(Order::getCity)          // 等价 o -> o.getCity()
    .forEach(System.out::println);
// 关键: 构造器引用 Order::new; 数组引用 int[]::new
Collectors.groupingBy
分类器 + 下游收集器两段式: groupingBy(Order::city, counting()) / averagingDouble / mapping(..., toList()), 下游还能再嵌 groupingBy — 报表的多维聚合全靠它。
Map<String, Long> cnt = orders.stream()
    .collect(Collectors.groupingBy(Order::city,
        Collectors.counting()));  // → {SH=3, BJ=5}
// 关键: downstream 可换 averaging/summing/mapping, 还能再嵌 groupingBy
flatMap 摊平
元素→流的摊平: 一单多明细摊成 Stream<Item>, 一行拆多词摊成 Stream<String>。map 会增加流的层数, flatMap 保证永远只有一层。
List<List<Integer>> nested = List.of(List.of(1, 2), List.of(3));
nested.stream().flatMap(List::stream).toList(); // → [1, 2, 3]
// 关键: map 会加层 (Stream<Stream>), flatMap 保证只有一层
Optional 与流
findFirst/findAny/max 返回 Optional: "可能没有"进类型契约。正确姿势: map/orElse/orElseThrow 链; 禁止裸 get(), 也不做字段/参数 — 它是返回值专用。
Optional<Order> max = orders.stream().max(cmp);
max.map(Order::id).orElse(-1L);   // 链式兜底
// 关键: 禁止裸 get(); Optional 只做返回值, 不做字段/参数
IntStream 避免装箱
Stream<Integer> 每个数都是堆对象; IntStream/LongStream/DoubleStream 走原生类型, sum/average/max 零装箱。百万级数值计算差一个数量级。
int fast = IntStream.range(0, 1_000_000).sum(); // 原生零装箱
int slow = Stream.of(1, 2, 3).reduce(0, Integer::sum); // 装箱
// 关键: 百万级数值计算两者差一个数量级
parallelStream 原理
把源切分成子区间任务提交 ForkJoinPool.commonPool(并行度 = 核数-1, 全 JVM 共享), 工作窃取负载均衡。适合: CPU 密集 + 大数据量 + 无状态无 IO 无共享。
long n = IntStream.range(0, 100_000_000).parallel()
    .filter(i -> isPrime(i)).count();
// 关键: commonPool 并行度 = 核数-1 且全 JVM 共享,
// 只留给 CPU 密集 + 无状态 + 无 IO 的纯计算管道
流不可复用
一个 Stream 绑一次管道生命周期: 终端操作后即"已操作或已关闭", 再用抛 IllegalStateException。要反复用: 先收集成集合, 或 Supplier<Stream<T>> 每次造新流。
Stream<Integer> s = Stream.of(1, 2, 3);
s.count();
s.count();  // → IllegalStateException: already operated upon
// 关键: 要复用先 collect 成集合, 或 Supplier<Stream<T>> 每次新开

🏭 生产实战 real world

场景 1 · 报表聚合: groupingBy + averaging 干掉手写循环 Map

运营日报要"各城市平均客单价", 手写循环是两层 for + 三个临时 Map; Stream 一句表达意图:

// 手写循环版: ~30 行样板, 每加一个指标都要重抄一遍结构
Map<String, List<Order>> tmp = new HashMap<>();
for (Order o : orders)
    tmp.computeIfAbsent(o.city(), k -> new ArrayList<>()).add(o);
Map<String, Double> avg = new HashMap<>();
for (var e : tmp.entrySet()) {
    double s = 0;
    for (Order o : e.getValue()) s += o.amount();
    avg.put(e.getKey(), s / e.getValue().size());
}

// Stream 版: groupingBy + downstream 一步到位, 分组键和组内算法各自独立演进
Map<String, Double> avgByCity = orders.stream()
    .collect(Collectors.groupingBy(Order::city,
        Collectors.averagingDouble(Order::amount)));

改造后新增指标只需换 downstream (如 summingInt/counting), 分组逻辑零改动。

场景 2 · flatMap 摊平订单-明细算 GMV

一单 N 件商品, GMV 要在"件"的粒度算并剔除退款 — 嵌套 for 的外部累加变量是共享可变状态, flatMap 之后语义立刻清晰:

// 订单 → 多条明细: flatMap 把 Stream<Order> 摊平成 Stream<Item>
BigDecimal gmv = orders.stream()
    .flatMap(o -> o.items().stream())          // 一单 N 件 → 逐件流出
    .filter(i -> i.status() != ItemStatus.REFUNDED)   // 退款不计 GMV
    .map(Item::payAmount)
    .reduce(BigDecimal.ZERO, BigDecimal::add);   // 金额必须 BigDecimal

// 对比嵌套 for: 两层循环 + 外部累加变量 — 想并行时它就是最大障碍

场景 3 · Optional 链式取值替代四层 if null

用户中心查城市, 任何一环都可能为 null; 嵌套 if 是缩进地狱, Optional 链一样的短路语义一条读到底:

// 旧写法: 四层缩进, 每一环都要人工判空, 漏一个就是 NPE
City cityOld = null;
if (user != null) {
    Profile p = user.getProfile();
    if (p != null && p.getAddress() != null) {
        cityOld = p.getAddress().getCity();
    }
}

// Optional 链: 中途 null 整链变 empty, 不再往下走 — 与上面完全同语义
City city = Optional.ofNullable(user)
    .map(User::getProfile)
    .map(Profile::getAddress)
    .map(Address::getCity)
    .orElse(City.UNKNOWN);              // 兜底显式, 不许裸 get()

场景 4 · IntStream 并行爆破: 8C 容器 41s → 6.3s

风控哈希爆破扫 1 亿候选, 串行 41.2s; 数据无状态无共享, 是 parallel 的教科书场景:

// CPU 密集 + 无 IO + 无共享可变状态 → parallel 的三个前提全满足
long hits = IntStream.range(0, 100_000_000)
    .parallel()                          // 自动切分给 commonPool 所有核
    .filter(i -> slowHash(i) == TARGET)
    .count();

// 量化: 8C 容器串行 41.2s → 并行 6.3s (≈6.5x, 贴近核数收益上限)
// 判据: 单元素计算达微秒级才值得 parallel; 纯逻辑过滤开并行反而更慢

场景 5 · 规则引擎 findFirst: 200 条规则只跑前 3 条

规则按优先级排好, 第一个命中即返回 — 短路让 197 条规则根本不执行:

// 规则引擎: 命中第一个就停, 后面的规则连 lambda 都不会被调用
Optional<Violation> hit = rules.stream()
    .sorted(Comparator.comparingInt(Rule::priority))
    .map(r -> r.apply(order))           // 惰性: 没到终端, map 根本不执行
    .filter(Optional::isPresent)
    .map(Optional::get)
    .findFirst();                    // e1 命中 → e2 起全部不再计算

// 误用 collect(toList()).get(0) 的代价: 全部规则先跑一遍, P99 3ms → 480ms

场景 6 · toMap 构建索引表: merge 函数是保命参数

全量用户拉回内存建 userId 索引, 脏数据里存在重复注册 — 没有 merge 函数直接抛异常:

// 构建 userId → User 索引: 第三个参数防撞键, 否则线上一条脏数据炸全接口
Map<Long, User> idx = users.stream()
    .collect(Collectors.toMap(
        User::id,                        // key mapper
        Function.identity(),             // value mapper: 元素本身
        (a, b) -> a));                   // 撞键保留先到者: 重复注册吞掉并告警

// 不传 merge: 一旦重复 key 直接抛
// IllegalStateException: Duplicate key 10086 (attempted merging values ...)

场景 7 · peek 调试: 正确姿势与滥用边界

排查"元素到底走到哪一步掉了", peek 是唯一的观察窗 — 但它只是调试工具, 不是执行逻辑的地方:

// 正确: 临时调试 — 打印穿链元素, 看清"逐元素"执行顺序, 用完即删
orders.stream()
    .peek(o -> log.debug("before filter: {}", o.id()))
    .filter(o -> o.amount() > 100)
    .peek(o -> log.debug("after filter: {}", o.id()))
    .count();                            // 终端: 没有它, 上面全部不执行

// 滥用: 在 peek 里改状态/发请求 — 规范只保证"结果等价", peek 可被整体跳过
// (如 count() 走 size 优化时连 filter 都不跑); 副作用逻辑必须放 forEach

场景 8 · 千万行导出: MyBatis 游标流式处理, 512MB 容器跑通

大表导出不能把千万行一次性 load 进堆 — JDBC 侧开流式结果集, 应用侧用 Cursor 逐条消费:

// JDBC 需 fetchSize=Integer.MIN_VALUE (MySQL 流式), 否则驱动仍全量缓冲
try (Cursor<Order> cursor = orderMapper.scanByTime(from, to)) {
    cursor.forEach(order -> writer.write(order.csvLine()));
}   // Cursor 实现 AutoCloseable — 不 close 会占死连接与结果集

// JPA 同理: repository.streamByTime(...) 也必须包在 try-with-resources 里
// 效果: 千万行导出峰值堆 < 100MB; 换 list() 直接 OOMKilled

场景 9 · 按字段去重: TreeSet 下游收集器一把梭

distinct() 只认 equals, 业务上"按城市去重且留金额最大单"要靠收集器组合:

// 去重键 = city; 同城冲突取 amount 最大 — TreeSet 天然按比较器去重
List<Order> dedup = orders.stream()
    .collect(Collectors.collectingAndThen(
        Collectors.toCollection(() ->
            new TreeSet<>(Comparator.comparing(Order::city)
                .thenComparing(Comparator.comparing(Order::amount).reversed()))),
        ArrayList::new));               // TreeSet → List, 保留流序语义交给下游

// 复杂度 O(n log n); 想保持原始出现顺序再套 LinkedHashMap 分组取首个

场景 10 · 迁移评审清单: 什么时候不该改成 stream

团队从 for 循环迁 Stream 翻过车, 沉淀出四条 checklist — 满足全部才改, 否则保持循环:

// ✓ 纯转换/过滤/聚合, 无复杂 break/continue 控制流 (takeWhile 例外)
// ✓ 循环体无副作用: 不改外部变量/集合, 不发 IO 请求
// ✓ 元素量 > 数千, 或语义确实是"管道" — 30 个元素 for 循环更快更直白
// ✗ 热路径 int 求和 (装箱), 多层 break, 依赖执行顺序 → 一律保留循环

// 团队规约落成一句话: 报表类代码用 stream, 内核热路径用循环
int total = 0;
for (Order o : orders) { if (o.paid()) total += o.amountCents(); }  // 保留: 热路径

⚠️ 编码注意与常见坑 pitfalls

坑 1 · parallelStream 里调下游 RPC — 上线后无关接口的 CompletableFuture 集体超时, jstack 里 commonPool worker 全 BLOCKED 在 socketRead. 原因: 公共池容量仅 核数-1 且全 JVM 共享, 被慢 IO 霸占后所有并行任务排队. 正解: parallel 只留 CPU 密集, IO 逻辑拆独立线程池或 MQ 异步。
// 错: list.parallelStream().forEach(o -> rpc.call(o));
//     慢 IO 霸占 commonPool → 全 JVM 并行任务集体排队
// 对: parallel 只留 CPU 密集; IO 拆独立线程池或 MQ 异步
坑 2 · lambda 里往外部 ArrayList add — 并行流偶发丢元素甚至 ArrayIndexOutOfBoundsException. 原因: ArrayList 非线程安全, 且违反"无副作用"纪律 — 执行时机与线程都不受你控制. 正解: 结果用 collect(...) 归约产生, 不改外部状态。
// 错: List<R> out = new ArrayList<>(); s.parallel().forEach(out::add);
//     → 非线程安全, 偶发丢元素甚至 ArrayIndexOutOfBoundsException
// 对: List<R> out = s.collect(Collectors.toList()); 归约产生结果
坑 3 · toMap 撞键 — 运行时抛 IllegalStateException: Duplicate key .... 原因: 两个元素映射出同一 key 且没给 merge 函数. 正解: 传第三个参数 (a, b) -> a 明确保留策略, 或上游先去重。
// 错: users.collect(Collectors.toMap(User::id, Function.identity()));
//     重复 key → IllegalStateException: Duplicate key
// 对: toMap(User::id, Function.identity(), (a, b) -> a);
坑 4 · 流二次消费 — 第二次遍历抛 IllegalStateException: stream has already been operated upon or closed. 原因: 管道有生命周期, 终端操作即"用完". 正解: 要多次用先 collect 成集合, 或 Supplier<Stream<T>> 每次新开。
// 错: Stream<Integer> s = list.stream(); s.count(); s.count();
//     第二次 → IllegalStateException: already operated upon or closed
// 对: 先 collect 成集合复用; 或 Supplier<Stream<T>> 每次新开
坑 5 · peek 一条日志都没打 — 加了 peek 的调试代码毫无输出. 原因: 链尾没有终端操作, 整条管道惰性空转, 且规范允许实现直接跳过 peek. 正解: 补终端操作; 副作用逻辑正式代码用 forEach, peek 只做临时调试。
// 错: list.stream().peek(x -> log("{}", x));  链尾无终端
//     → 整条管道惰性空转, 一条日志都没有
// 对: 补终端 .toList(); 副作用逻辑正式代码用 forEach
坑 6 · skip/limit 顺序语义 — skip(2).limit(5) 出 5 个元素, limit(5).skip(2) 只出 3 个, 顺手换序结果就变; 无限流上只有 limit 能收口. 原因: 中间操作按声明顺序层层包裹, 两者都是有状态裁剪. 正解: 明确"先跳再取"次序, 无限流必配 limit 或短路终端。
List<Integer> l = List.of(1,2,3,4,5,6,7);
l.stream().skip(2).limit(5).toList(); // → [3, 4, 5, 6, 7]
l.stream().limit(5).skip(2).toList(); // → [3, 4, 5]
// 关键: 换序结果就变; 无限流必配 limit 收口
坑 7 · groupingBy 结果乱序 — 分组结果 key 顺序每次重启都变, 下游按序拼接的报表串行. 原因: 默认用 HashMap, 本就无序. 正解: groupingBy(f, LinkedHashMap::new, downstream), 要排序用 TreeMap 版。
// 错: groupingBy(Order::city, counting())  默认 HashMap, key 顺序每次变
// 对: groupingBy(Order::city, LinkedHashMap::new, counting());
//     要排序用 TreeMap 版 mapFactory
坑 8 · 装箱流求和反而慢 — 百万级 Stream<Integer>.reduce(0, Integer::sum) 比 for 循环慢数倍且 GC 压力大. 原因: 每个元素都装箱成堆对象. 正解: IntStream.range().sum() 原生管道; 热路径保留循环。
// 错: Stream<Integer>.reduce(0, Integer::sum);  百万级装箱慢数倍
// 对: IntStream.range(0, n).sum();  原生零装箱
//     热路径 int 求和直接保留 for 循环
坑 9 · parallel 下可变 reduce — 用 reduce(new ArrayList<>(), (l, e) -> { l.add(e); return l; }) 并行收集, 结果随机丢元素. 原因: 多线程共享同一个初始容器, combiner 形同虚设. 正解: collect(ArrayList::new, List::add, List::addAll) 三参收集器。
// 错: s.parallel().reduce(new ArrayList<>(), (l, e) -> { l.add(e); return l; });
//     多线程共享初始容器 → 随机丢元素
// 对: s.collect(ArrayList::new, List::add, List::addAll);
坑 10 · lambda 里抛受检异常编译不过 — "Unhandled exception: java.io.IOException". 原因: Function/Consumer 的方法签名不含受检异常. 正解: 包成 RuntimeException 保留原始栈, 或 lombok @SneakyThrows, 或泛型走私技巧 (见泛型页)。
// 错: files.forEach(f -> Files.readAllBytes(f));
//     → Unhandled exception: java.io.IOException
// 对: 包成 RuntimeException 保留原始栈; 或 lombok @SneakyThrows
坑 11 · 全量 sorted 排完再 filter — 100 万元素先 sorted 再 filter 耗时 2.1s, 换序 80ms. 原因: sorted 有状态 O(n log n), 必须排完才轮到 filter 丢弃大半. 正解: filter 尽量前置, 让 sorted 只处理幸存者。
// 错: orders.stream().sorted(byAmount).filter(o -> o.paid())
//     → 先排完 100 万再丢弃大半, 2.1s
// 对: orders.stream().filter(o -> o.paid()).sorted(byAmount);  80ms
坑 12 · toList 可变性误会 — 对 Collectors.toList() 的结果 add 抛 UnsupportedOperationException, 或以为不可变却被改了. 原因: 该方法对可变性无任何契约, 实现随版本变. 正解: 要不可变用 Stream.toList() (16+), 要可变用 toCollection(ArrayList::new)。
// 错: Collectors.toList() 结果直接 add — 可变性无契约, 版本一变就抛
// 对: 不可变用 stream.toList() (JDK 16+);
//     要可变用 Collectors.toCollection(ArrayList::new)
坑 13 · 想遍历两次 — 第一遍 forEach 后第二遍 map 抛 IllegalStateException. 原因: 一条管道只伴随一次遍历, 流不是集合. 正解: 先收集成 List 再多次遍历, 或每遍新开一个流。
// 错: Stream<Order> s = list.stream(); s.forEach(...); s.map(...);
//     第二遍 → IllegalStateException
// 对: 先 collect 成 List 再多次遍历; 或每遍新开一个流
坑 14 · Files.lines 没关 — 日志分析作业跑几次后 "Too many open files". 原因: 流持有文件句柄, 不 close 不释放. 正解: try (Stream<String> lines = Files.lines(path)); JPA Stream / MyBatis Cursor 同理。
// 错: Files.lines(path).filter(...).count();  句柄不释放
//     → 跑几次后 Too many open files
// 对: try (Stream<String> lines = Files.lines(path)) { ... }
坑 15 · 并行流里 ThreadLocal 丢 — 日志 traceId 断断续续缺失, MDC 上下文取不到. 原因: 任务跑在 commonPool 工作线程, 不继承调用线程的 ThreadLocal. 正解: lambda 内手动 set/清理 MDC, 或用上下文传递框架包装提交。
// 错: parallelStream 里读 MDC.get("traceId")
//     → 任务在 commonPool 线程, 不继承调用线程的 ThreadLocal
// 对: lambda 内手动 set/清理 MDC; 或上下文传递框架包装提交
坑 16 · findAny 当 findFirst 用 — 单测偶尔过偶尔不过, 返回值不稳定. 原因: 并行下 findAny 返回"最先完成者"(非确定), findFirst 才保证 encounter 顺序(代价是顺序约束). 正解: 语义要"第一个"就 findFirst, 只要"存在一个"才用 findAny。
// 错: 并行下用 findAny() 却当"第一个"用 — 返回最先完成者, 不稳定
// 对: 语义要"第一个"用 findFirst(); 只要"存在一个"才 findAny()
坑 17 · BigDecimal 求和差一分 — 金额汇总对账不平. 原因: new BigDecimal(0.1) 的 double 构造带二进制误差, equals 还连 scale 一起比 (2.0 与 2.00 不等). 正解: BigDecimal.valueOf 或字符串构造, 比较用 compareTo, 求和 reduce(BigDecimal.ZERO, BigDecimal::add)。
// 错: new BigDecimal(0.1);  double 构造带二进制误差;
//     equals 连 scale 一起比: 2.0 与 2.00 不等
// 对: BigDecimal.valueOf(0.1); 比较用 compareTo;
//     求和 reduce(BigDecimal.ZERO, BigDecimal::add)
坑 18 · flatMap 后内存翻倍 OOM — 订单流 flatMap 出明细再 collect, 千万级时 OOMKilled. 原因: 中间流的每个元素都通过闭包拽着整个外层大对象, 峰值 = 全量明细 × 外层对象. 正解: flatMap 前先把外层投影成小 DTO, 或分页分批处理。
// 错: orders.flatMap(o -> o.items().stream())
//     → 每个明细闭包都拽着整个 Order 大对象, 峰值翻倍
// 对: flatMap 前先 map 成小 DTO; 或分页分批处理
坑 19 · 方法引用 NPE 难定位 — 线上 NPE 栈只有 at App.lambda$main$0(App.java:12), 看不出哪一环空了. 原因: 方法引用被编译成合成 lambda, 源位置信息丢失. 正解: 关键链路用 Optional.map 逐环兜底, 编译带 -g 保留参数名辅助定位。
// 错: NPE 栈只有 at App.lambda$main$0(App.java:12) — 看不出哪一环空
// 对: 关键链路 Optional.map 逐环兜底; 编译带 -g 保留参数名
坑 20 · parallel forEach 写文件乱序 — 并行 forEach 逐行写文件, 行序随机错乱. 原因: 并行 forEach 对处理顺序零保证. 正解: 顺序敏感用 forEachOrdered (牺牲并行度), 或先 collect 再统一顺序输出。
// 错: list.parallelStream().forEach(l -> writer.write(l));
//     → 并行 forEach 对处理顺序零保证, 行序随机错乱
// 对: 顺序敏感用 forEachOrdered; 或先 collect 再统一输出