中间操作只注册不执行, 终端操作一次点火 — 元素逐个穿全链; parallelStream 的 ForkJoinPool.commonPool 是全 JVM 共享的
Stream 的本质是"声明式管道 + 惰性求值": filter/map/sorted 这些中间操作只是往管道上挂节点, 一行都不执行; 直到 collect/findFirst 这样的终端操作点火, 数据才开始流动。而流动方式也不是"第一阶段跑完全部元素再进第二阶段", 而是一个元素穿完整条链, 才轮到下一个元素(sorted 这类有状态操作除外)。这一个模型同时解释了三件事: 为什么漏写终端操作整段代码"白写"; 为什么 findFirst 命中第一个就不再多算一个元素; 为什么 peek 里的日志"有时不执行"。
parallelStream 只是把"切分 + 调度"交给 ForkJoinPool.commonPool — 它是全 JVM 共享的公共资源(默认核数-1 个线程, 工作窃取均衡负载)。任何一个并行流里做慢 IO, 全进程的并行流与 CompletableFuture 默认执行器一起陪葬; 所以业界铁律是 parallel 只留给 CPU 密集、无状态、无 IO 的纯计算管道。
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] // 关键: 没人要的元素永远不会被造出来 — 这就是无限流能存在的原因
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 同为短路List.of(1, 2).stream() .peek(x -> log("F{}", x)) // 打印顺序: F1 M1 F2 M2 .map(x -> x * 10) .toList(); // → [10, 20] // 关键: 逐元素垂直穿链, 不是逐阶段水平搬批
// 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[]::newgroupingBy(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, 还能再嵌 groupingByStream<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<Order> max = orders.stream().max(cmp); max.map(Order::id).orElse(-1L); // 链式兜底 // 关键: 禁止裸 get(); Optional 只做返回值, 不做字段/参数
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); // 装箱 // 关键: 百万级数值计算两者差一个数量级
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 的纯计算管道
IllegalStateException。要反复用: 先收集成集合, 或 Supplier<Stream<T>> 每次造新流。Stream<Integer> s = Stream.of(1, 2, 3); s.count(); s.count(); // → IllegalStateException: already operated upon // 关键: 要复用先 collect 成集合, 或 Supplier<Stream<T>> 每次新开
运营日报要"各城市平均客单价", 手写循环是两层 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), 分组逻辑零改动。
一单 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: 两层循环 + 外部累加变量 — 想并行时它就是最大障碍
用户中心查城市, 任何一环都可能为 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()
风控哈希爆破扫 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; 纯逻辑过滤开并行反而更慢
规则按优先级排好, 第一个命中即返回 — 短路让 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
全量用户拉回内存建 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 ...)
排查"元素到底走到哪一步掉了", 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
大表导出不能把千万行一次性 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
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 分组取首个
团队从 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(); } // 保留: 热路径
parallel 只留 CPU 密集, IO 逻辑拆独立线程池或 MQ 异步。// 错: list.parallelStream().forEach(o -> rpc.call(o)); // 慢 IO 霸占 commonPool → 全 JVM 并行任务集体排队 // 对: parallel 只留 CPU 密集; IO 拆独立线程池或 MQ 异步
collect(...) 归约产生, 不改外部状态。// 错: List<R> out = new ArrayList<>(); s.parallel().forEach(out::add); // → 非线程安全, 偶发丢元素甚至 ArrayIndexOutOfBoundsException // 对: List<R> out = s.collect(Collectors.toList()); 归约产生结果
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);
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>> 每次新开
forEach, peek 只做临时调试。// 错: list.stream().peek(x -> log("{}", x)); 链尾无终端 // → 整条管道惰性空转, 一条日志都没有 // 对: 补终端 .toList(); 副作用逻辑正式代码用 forEach
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 收口
groupingBy(f, LinkedHashMap::new, downstream), 要排序用 TreeMap 版。// 错: groupingBy(Order::city, counting()) 默认 HashMap, key 顺序每次变 // 对: groupingBy(Order::city, LinkedHashMap::new, counting()); // 要排序用 TreeMap 版 mapFactory
Stream<Integer>.reduce(0, Integer::sum) 比 for 循环慢数倍且 GC 压力大. 原因: 每个元素都装箱成堆对象. 正解: IntStream.range().sum() 原生管道; 热路径保留循环。// 错: Stream<Integer>.reduce(0, Integer::sum); 百万级装箱慢数倍 // 对: IntStream.range(0, n).sum(); 原生零装箱 // 热路径 int 求和直接保留 for 循环
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);
@SneakyThrows, 或泛型走私技巧 (见泛型页)。// 错: files.forEach(f -> Files.readAllBytes(f)); // → Unhandled exception: java.io.IOException // 对: 包成 RuntimeException 保留原始栈; 或 lombok @SneakyThrows
// 错: orders.stream().sorted(byAmount).filter(o -> o.paid()) // → 先排完 100 万再丢弃大半, 2.1s // 对: orders.stream().filter(o -> o.paid()).sorted(byAmount); 80ms
Collectors.toList() 的结果 add 抛 UnsupportedOperationException, 或以为不可变却被改了. 原因: 该方法对可变性无任何契约, 实现随版本变. 正解: 要不可变用 Stream.toList() (16+), 要可变用 toCollection(ArrayList::new)。// 错: Collectors.toList() 结果直接 add — 可变性无契约, 版本一变就抛 // 对: 不可变用 stream.toList() (JDK 16+); // 要可变用 Collectors.toCollection(ArrayList::new)
// 错: Stream<Order> s = list.stream(); s.forEach(...); s.map(...); // 第二遍 → IllegalStateException // 对: 先 collect 成 List 再多次遍历; 或每遍新开一个流
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)) { ... }
// 错: parallelStream 里读 MDC.get("traceId") // → 任务在 commonPool 线程, 不继承调用线程的 ThreadLocal // 对: lambda 内手动 set/清理 MDC; 或上下文传递框架包装提交
// 错: 并行下用 findAny() 却当"第一个"用 — 返回最先完成者, 不稳定 // 对: 语义要"第一个"用 findFirst(); 只要"存在一个"才 findAny()
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)
// 错: orders.flatMap(o -> o.items().stream()) // → 每个明细闭包都拽着整个 Order 大对象, 峰值翻倍 // 对: flatMap 前先 map 成小 DTO; 或分页分批处理
at App.lambda$main$0(App.java:12), 看不出哪一环空了. 原因: 方法引用被编译成合成 lambda, 源位置信息丢失. 正解: 关键链路用 Optional.map 逐环兜底, 编译带 -g 保留参数名辅助定位。// 错: NPE 栈只有 at App.lambda$main$0(App.java:12) — 看不出哪一环空 // 对: 关键链路 Optional.map 逐环兜底; 编译带 -g 保留参数名
forEachOrdered (牺牲并行度), 或先 collect 再统一顺序输出。// 错: list.parallelStream().forEach(l -> writer.write(l)); // → 并行 forEach 对处理顺序零保证, 行序随机错乱 // 对: 顺序敏感用 forEachOrdered; 或先 collect 再统一输出