Stream API 高级详解
大翔收到了一封邮件——Frank 总监要求他对飞翔科技第一季度所有员工的绩效数据进行多维度分析。
大翔打开数据文件,CSV 里有两千多行员工数据。按以往的写法,他要写十几层 for 循环和 if 判断。
白歌凑过来:"大翔哥,用 Stream API 啊!filter 过滤,groupingBy 分组,summarizingInt 汇总,一条 pipeline 搞定。"
大翔眼睛一亮:"JDK 8 的 Stream?我只用过简单的 forEach..."
白歌拉过椅子:"来,我教你。Stream 的精髓在于——声明式编程。你告诉它'做什么',而不是'怎么做'。"
小崔在旁边记笔记:"这跟 SQL 很像啊!SELECT WHERE GROUP BY ORDER BY..."
朱璐一拍桌子:"对!Stream 就是 Java 的 SQL!"
一、定义表
1.1 中间操作(Intermediate Operations)——惰性求值
| 操作 | 方法签名 | 说明 | 是否短路 |
|---|---|---|---|
filter | Stream<T> filter(Predicate<T>) | 过滤元素,保留满足条件的 | 否 |
map | <R> Stream<R> map(Function<T,R>) | 元素一对一映射 | 否 |
flatMap | <R> Stream<R> flatMap(Function<T,Stream<R>>) | 元素一对多扁平化映射 | 否 |
distinct | Stream<T> distinct() | 去重(依赖 equals) | 否 |
sorted | Stream<T> sorted(Comparator<T>) | 排序(有状态操作) | 否 |
peek | Stream<T> peek(Consumer<T>) | 调试/中间消费,不改变流 | 否 |
limit | Stream<T> limit(long n) | 截取前 n 个 | 是(短路) |
skip | Stream<T> skip(long n) | 跳过前 n 个 | 否 |
1.2 终止操作(Terminal Operations)——触发计算
| 操作 | 方法签名 | 说明 | 返回类型 |
|---|---|---|---|
collect | <R,A> R collect(Collector<T,A,R>) | 收集到集合/Map/字符串 | R |
reduce | T reduce(T identity, BinaryOperator<T>) | 归约/聚合 | T |
forEach | void forEach(Consumer<T>) | 遍历每个元素 | void |
count | long count() | 计数 | long |
findFirst | Optional<T> findFirst() | 找第一个元素 | Optional |
findAny | Optional<T> findAny() | 找任意元素(并行友好) | Optional |
anyMatch | boolean anyMatch(Predicate<T>) | 是否有任意匹配 | boolean |
allMatch | boolean allMatch(Predicate<T>) | 是否全部匹配 | boolean |
noneMatch | boolean noneMatch(Predicate<T>) | 是否无一匹配 | boolean |
max/min | Optional<T> max/min(Comparator<T>) | 最大/最小值 | Optional |
1.3 Collectors 常用工厂方法
| 方法 | 说明 | 示例 |
|---|---|---|
toList() | 收集到 List | .collect(Collectors.toList()) |
toSet() | 收集到 Set(去重) | .collect(Collectors.toSet()) |
toMap(keyMapper, valueMapper) | 收集到 Map | .collect(Collectors.toMap(E::getId, E::getName)) |
groupingBy(classifier) | 按条件分组 → Map<K, List<T>> | .collect(Collectors.groupingBy(E::getDept)) |
partitioningBy(predicate) | 按布尔条件分区 | .collect(Collectors.partitioningBy(e -> e.salary > 10000)) |
joining(delimiter) | 字符串拼接 | .collect(Collectors.joining(", ")) |
summarizingInt/Long/Double | 统计信息(count/sum/min/max/avg) | .collect(Collectors.summarizingInt(E::getAge)) |
averagingInt | 平均值 | .collect(Collectors.averagingInt(E::getAge)) |
mapping(func, downstream) | 映射后再收集 | .collect(Collectors.mapping(E::getName, toList())) |
reducing(identity, mapper, op) | 归约收集 | .collect(Collectors.reducing(0, E::getScore, Integer::sum)) |
二、Mermaid 流程图:Stream Pipeline 执行模型
三、Mermaid 类图:Collectors 关键方法体系
四、Stream 三大特性深入
4.1 惰性求值(Lazy Evaluation)
// 中间操作不会执行,直到终止操作触发
Stream<String> stream = names.stream()
.filter(n -> {
System.out.println("filter: " + n); // 此时不打印
return n.length() > 2;
})
.map(n -> {
System.out.println("map: " + n); // 此时不打印
return n.toUpperCase();
});
// 到此为止,没有一行打印输出!
// 终止操作触发整个 pipeline 执行
List<String> result = stream.collect(Collectors.toList());
// 现在才依次输出 filter 和 map
4.2 短路操作(Short-Circuiting)
limit(n), findFirst(), findAny(), anyMatch(), allMatch(), noneMatch() 是短路操作——它们不需要处理全部元素就能得到结果。
// limit(2) 短路:只处理前两个元素就停止
Stream.of(1, 2, 3, 4, 5)
.peek(x -> System.out.println("处理: " + x))
.limit(2)
.collect(Collectors.toList());
// 输出:
// 处理: 1
// 处理: 2
// (3, 4, 5 不会被处理!)
4.3 并行流注意事项
// parallelStream 使用 ForkJoinPool.commonPool()
// 适合 CPU 密集型任务,不适合 I/O 密集型
// 注意线程安全问题:避免在并行流中使用非线程安全集合
// 错误示例:
List<String> list = new ArrayList<>(); // 非线程安全!
stream.parallel().forEach(list::add); // 可能导致数据丢失或异常
// 正确示例:
List<String> list = stream.parallel().collect(Collectors.toList()); // 线程安全
五、完整代码示例
示例1:飞翔科技员工数据复杂聚合分析
import java.util.*;
import java.util.stream.*;
import java.util.function.*;
/**
* 飞翔科技 2024 Q1 员工绩效数据分析
* 演示 Stream API 各种中间操作和终止操作
*/
class FeixiangEmployee {
private int id;
private String name;
private String department; // 研发部、市场部、财务部、人事部
private String city; // 北京、上海、深圳、杭州
private double salary;
private int age;
private String level; // P4, P5, P6, P7, P8
private double performanceScore; // 绩效评分 0-100
public FeixiangEmployee(int id, String name, String department, String city,
double salary, int age, String level, double performanceScore) {
this.id = id; this.name = name; this.department = department;
this.city = city; this.salary = salary; this.age = age;
this.level = level; this.performanceScore = performanceScore;
}
// Getters
public int getId() { return id; }
public String getName() { return name; }
public String getDepartment() { return department; }
public String getCity() { return city; }
public double getSalary() { return salary; }
public int getAge() { return age; }
public String getLevel() { return level; }
public double getPerformanceScore() { return performanceScore; }
@Override
public String toString() {
return String.format("%s[%s|%s|¥%.2f|%d岁|%s|%.1f分]",
name, department, city, salary, age, level, performanceScore);
}
}
public class FeixiangStreamDemo {
public static void main(String[] args) {
System.out.println("========== 飞翔科技 Q1 员工数据分析 ==========\n");
// 构造测试数据
List<FeixiangEmployee> employees = Arrays.asList(
new FeixiangEmployee(1001, "大翔", "研发部", "北京", 18888.88, 30, "P7", 92.5),
new FeixiangEmployee(1002, "白歌", "研发部", "深圳", 15888.88, 28, "P6", 88.0),
new FeixiangEmployee(1003, "小崔", "研发部", "北京", 8888.88, 24, "P4", 75.5),
new FeixiangEmployee(1004, "朱璐", "研发部", "杭州", 12888.88, 26, "P5", 85.0),
new FeixiangEmployee(1005, "Frank", "市场部", "上海", 28888.88, 38, "P8", 95.0),
new FeixiangEmployee(1006, "黄俪", "人事部", "北京", 10888.88, 32, "P6", 78.5),
new FeixiangEmployee(1007, "李眉", "财务部", "上海", 13888.88, 29, "P6", 90.0),
new FeixiangEmployee(1008, "孔蓝", "市场部", "深圳", 16888.88, 31, "P6", 82.0),
new FeixiangEmployee(1009, "刘洋", "研发部", "北京", 9888.88, 25, "P5", 70.0),
new FeixiangEmployee(1010, "张伟", "财务部", "杭州", 11888.88, 27, "P5", 86.5)
);
// ===== 1. 基础过滤与映射 =====
System.out.println("【1> 研发部薪资高于10000的员工(按薪资降序)】");
employees.stream()
.filter(e -> "研发部".equals(e.getDepartment()))
.filter(e -> e.getSalary() > 10000)
.sorted(Comparator.comparingDouble(FeixiangEmployee::getSalary).reversed())
.forEach(e -> System.out.printf(" %s - ¥%.2f\n", e.getName(), e.getSalary()));
// ===== 2. map + distinct =====
System.out.println("\n【2> 所有部门列表(去重排序)】");
String departments = employees.stream()
.map(FeixiangEmployee::getDepartment)
.distinct()
.sorted()
.collect(Collectors.joining(" | "));
System.out.println(" " + departments);
// ===== 3. groupingBy 单级分组 =====
System.out.println("\n【3> 按部门分组】");
Map<String, List<FeixiangEmployee>> byDept = employees.stream()
.collect(Collectors.groupingBy(FeixiangEmployee::getDepartment));
byDept.forEach((dept, emps) -> {
String names = emps.stream().map(FeixiangEmployee::getName)
.collect(Collectors.joining(", "));
System.out.printf(" %s: [%s] (共 %d 人)\n", dept, names, emps.size());
});
// ===== 4. groupingBy 多级分组(部门 + 城市) =====
System.out.println("\n【4> 按部门再按城市二级分组】");
Map<String, Map<String, List<FeixiangEmployee>>> byDeptCity = employees.stream()
.collect(Collectors.groupingBy(
FeixiangEmployee::getDepartment,
Collectors.groupingBy(FeixiangEmployee::getCity)
));
byDeptCity.forEach((dept, cityMap) -> {
System.out.println(" " + dept + ":");
cityMap.forEach((city, emps) -> {
System.out.printf(" ├─ %s: %s\n", city,
emps.stream().map(FeixiangEmployee::getName)
.collect(Collectors.joining(", ")));
});
});
// ===== 5. 按部门统计薪资汇总 =====
System.out.println("\n【5> 各部门薪资统计汇总】");
Map<String, DoubleSummaryStatistics> salaryStats = employees.stream()
.collect(Collectors.groupingBy(
FeixiangEmployee::getDepartment,
Collectors.summarizingDouble(FeixiangEmployee::getSalary)
));
System.out.printf(" %-8s %8s %10s %10s %10s %10s\n",
"部门", "人数", "总薪资", "平均薪资", "最高", "最低");
System.out.println(" " + "-".repeat(58));
salaryStats.forEach((dept, stats) -> {
System.out.printf(" %-8s %8d %10.2f %10.2f %10.2f %10.2f\n",
dept, stats.getCount(),
stats.getSum(), stats.getAverage(),
stats.getMax(), stats.getMin());
});
// ===== 6. partitioningBy 分区 =====
System.out.println("\n【6> 按薪资是否超过15000分区】");
Map<Boolean, List<FeixiangEmployee>> partitioned = employees.stream()
.collect(Collectors.partitioningBy(e -> e.getSalary() > 15000));
System.out.printf(" 高薪资 (>¥15000): %s\n",
partitioned.get(true).stream()
.map(e -> e.getName() + " ¥" + e.getSalary())
.collect(Collectors.joining(", ")));
System.out.printf(" 普通薪资: %s\n",
partitioned.get(false).stream()
.map(e -> e.getName() + " ¥" + e.getSalary())
.collect(Collectors.joining(", ")));
// ===== 7. reduce 归约 =====
System.out.println("\n【7> 公司总薪资池与平均薪资】");
double totalSalary = employees.stream()
.map(FeixiangEmployee::getSalary)
.reduce(0.0, Double::sum);
double avgSalary = totalSalary / employees.size();
System.out.printf(" 总薪资池: ¥%.2f | 公司平均: ¥%.2f\n", totalSalary, avgSalary);
// ===== 8. flatMap 示例 =====
System.out.println("\n【8> flatMap: 每个部门的技能关键词扁平化】");
Map<String, List<String>> deptSkills = new LinkedHashMap<>();
deptSkills.put("研发部", Arrays.asList("Java", "Spring", "MySQL", "Docker"));
deptSkills.put("市场部", Arrays.asList("SEO", "SEM", "数据分析"));
deptSkills.put("财务部", Arrays.asList("Excel", "SAP", "税务"));
deptSkills.put("人事部", Arrays.asList("招聘", "培训", "绩效"));
// 获取所有部门的所有技能(去重)
List<String> allSkills = deptSkills.values().stream()
.flatMap(Collection::stream)
.distinct()
.sorted()
.collect(Collectors.toList());
System.out.println(" 全公司技能池: " + allSkills);
// ===== 9. peek 调试 =====
System.out.println("\n【9> peek 调试流处理过程】");
long highPerfCount = employees.stream()
.filter(e -> e.getPerformanceScore() > 85)
.peek(e -> System.out.printf(" [peek] 高绩效: %s (%.1f分)\n",
e.getName(), e.getPerformanceScore()))
.filter(e -> "研发部".equals(e.getDepartment()))
.peek(e -> System.out.printf(" [peek] 研发部高绩效: %s\n", e.getName()))
.count();
System.out.println(" 研发部高绩效人数: " + highPerfCount);
// ===== 10. Collectors.mapping + groupingBy 组合 =====
System.out.println("\n【10> 各部门绩效最高员工姓名】");
Map<String, String> topPerformerByDept = employees.stream()
.collect(Collectors.groupingBy(
FeixiangEmployee::getDepartment,
Collectors.collectingAndThen(
Collectors.maxBy(
Comparator.comparingDouble(FeixiangEmployee::getPerformanceScore)
),
opt -> opt.map(FeixiangEmployee::getName).orElse("无")
)
));
topPerformerByDept.forEach((dept, name) ->
System.out.printf(" %s: %s\n", dept, name));
// ===== 11. IntStream 基本类型流 =====
System.out.println("\n【11> IntStream 年龄统计分析】");
IntSummaryStatistics ageStats = employees.stream()
.mapToInt(FeixiangEmployee::getAge)
.summaryStatistics();
System.out.printf(" 年龄范围: %d ~ %d | 平均: %.1f岁 | 总和: %d\n",
ageStats.getMin(), ageStats.getMax(),
ageStats.getAverage(), ageStats.getSum());
System.out.println("\n========== End ==========");
}
}
运行输出:
========== 飞翔科技 Q1 员工数据分析 ==========
【1> 研发部薪资高于10000的员工(按薪资降序)】
大翔 - ¥18888.88
白歌 - ¥15888.88
朱璐 - ¥12888.88
【2> 所有部门列表(去重排序)】
人事部 | 市场部 | 研发部 | 财务部
【3> 按部门分组】
研发部: [大翔, 白歌, 小崔, 朱璐, 刘洋] (共 5 人)
市场部: [Frank, 孔蓝] (共 2 人)
人事部: [黄俪] (共 1 人)
财务部: [李眉, 张伟] (共 2 人)
【4> 按部门再按城市二级分组】
研发部:
├─ 北京: 大翔, 小崔, 刘洋
├─ 深圳: 白歌
├─ 杭州: 朱璐
市场部:
├─ 上海: Frank
├─ 深圳: 孔蓝
...
【5> 各部门薪资统计汇总】
部门 人数 总薪资 平均薪资 最高 最低
----------------------------------------------------------
研发部 5 66456.40 13291.28 18888.88 8888.88
市场部 2 45777.76 22888.88 28888.88 16888.88
人事部 1 10888.88 10888.88 10888.88 10888.88
财务部 2 25777.76 12888.88 13888.88 11888.88
...
【10> 各部门绩效最高员工姓名】
研发部: 大翔
市场部: Frank
人事部: 黄俪
财务部: 李眉
【11> IntStream 年龄统计分析】
年龄范围: 24 ~ 38 | 平均: 29.0岁 | 总和: 290
========== End ==========
示例2:飞翔科技商品订单多维分析
import java.util.*;
import java.util.stream.*;
/**
* 飞翔科技商城订单多维分析 —— 实战 Stream API
*/
class Order {
private String orderId;
private String customerName;
private String productCategory; // 电子产品、图书、服装、食品
private String productName;
private int quantity;
private double unitPrice;
private String month; // 2024-01, 2024-02, 2024-03
private String paymentMethod; // 微信、支付宝、银行卡
public Order(String orderId, String customerName, String productCategory,
String productName, int quantity, double unitPrice,
String month, String paymentMethod) {
this.orderId = orderId; this.customerName = customerName;
this.productCategory = productCategory; this.productName = productName;
this.quantity = quantity; this.unitPrice = unitPrice;
this.month = month; this.paymentMethod = paymentMethod;
}
public double getTotalAmount() { return quantity * unitPrice; }
public String getProductCategory() { return productCategory; }
public String getProductName() { return productName; }
public int getQuantity() { return quantity; }
public double getUnitPrice() { return unitPrice; }
public String getMonth() { return month; }
public String getPaymentMethod() { return paymentMethod; }
public String getCustomerName() { return customerName; }
}
public class OrderAnalysisDemo {
public static void main(String[] args) {
System.out.println("========== 飞翔科技商城订单分析 ==========\n");
List<Order> orders = Arrays.asList(
new Order("ORD001", "大翔", "电子产品", "机械键盘", 1, 188.88, "2024-01", "微信"),
new Order("ORD002", "白歌", "电子产品", "显示器", 1, 1288.88, "2024-01", "支付宝"),
new Order("ORD003", "小崔", "图书", "Java编程思想", 2, 88.88, "2024-01", "微信"),
new Order("ORD004", "朱璐", "服装", "T恤", 3, 68.88, "2024-01", "银行卡"),
new Order("ORD005", "大翔", "食品", "坚果礼盒", 2, 128.88, "2024-02", "微信"),
new Order("ORD006", "Frank","电子产品", "平板电脑", 1, 3888.88, "2024-02", "支付宝"),
new Order("ORD007", "黄俪", "服装", "衬衫", 2, 158.88, "2024-02", "微信"),
new Order("ORD008", "李眉", "图书", "Effective Java", 1, 78.88, "2024-02", "支付宝"),
new Order("ORD009", "孔蓝", "电子产品", "蓝牙耳机", 1, 288.88, "2024-03", "银行卡"),
new Order("ORD010", "大翔", "食品", "咖啡豆", 3, 58.88, "2024-03", "微信"),
new Order("ORD011", "白歌", "图书", "Clean Code", 2, 68.88, "2024-03", "支付宝"),
new Order("ORD012", "朱璐", "服装", "运动鞋", 1, 388.88, "2024-03", "微信")
);
// 1. 各品类销售额排行
System.out.println("【1> 各品类销售额排行】");
Map<String, Double> categorySales = orders.stream()
.collect(Collectors.groupingBy(
Order::getProductCategory,
Collectors.summingDouble(Order::getTotalAmount)
));
categorySales.entrySet().stream()
.sorted(Map.Entry<String, Double>comparingByValue().reversed())
.forEach(e -> System.out.printf(" %s: ¥%.2f\n", e.getKey(), e.getValue()));
// 2. 月度销售额趋势
System.out.println("\n【2> 各月度销售额趋势】");
Map<String, Double> monthlySales = orders.stream()
.collect(Collectors.groupingBy(
Order::getMonth,
TreeMap::new, // 按月份排序
Collectors.summingDouble(Order::getTotalAmount)
));
monthlySales.forEach((month, sales) ->
System.out.printf(" %s: ¥%.2f\n", month, sales));
// 3. 支付方式占比
System.out.println("\n【3> 支付方式分布】");
double totalAmount = orders.stream()
.mapToDouble(Order::getTotalAmount).sum();
Map<String, Double> paymentDistribution = orders.stream()
.collect(Collectors.groupingBy(
Order::getPaymentMethod,
Collectors.summingDouble(Order::getTotalAmount)
));
paymentDistribution.forEach((method, amount) ->
System.out.printf(" %s: ¥%.2f (%.1f%%)\n",
method, amount, amount / totalAmount * 100));
// 4. 客户消费排行(Top 3)
System.out.println("\n【4> 客户消费排行 Top 3】");
Map<String, Double> customerSpending = orders.stream()
.collect(Collectors.groupingBy(
Order::getCustomerName,
Collectors.summingDouble(Order::getTotalAmount)
));
customerSpending.entrySet().stream()
.sorted(Map.Entry<String, Double>comparingByValue().reversed())
.limit(3)
.forEach(e -> System.out.printf(" %s: ¥%.2f\n", e.getKey(), e.getValue()));
// 5. 每个品类下价格最高的商品
System.out.println("\n【5> 每个品类下最贵的商品】");
Map<String, String> mostExpensive = orders.stream()
.collect(Collectors.groupingBy(
Order::getProductCategory,
Collectors.collectingAndThen(
Collectors.maxBy(Comparator.comparingDouble(Order::getUnitPrice)),
opt -> opt.map(o -> o.getProductName() + " ¥" + o.getUnitPrice()).orElse("无")
)
));
mostExpensive.forEach((cat, prod) ->
System.out.printf(" %s: %s\n", cat, prod));
// 6. 订单金额分段统计
System.out.println("\n【6> 订单金额分段(<200 / 200~500 / >500)】");
Map<String, Long> amountSegments = orders.stream()
.collect(Collectors.groupingBy(order -> {
double amt = order.getTotalAmount();
if (amt < 200) return "小额 (<200)";
else if (amt <= 500) return "中额 (200~500)";
else return "大额 (>500)";
}, Collectors.counting()));
amountSegments.forEach((seg, cnt) ->
System.out.printf(" %s: %d 笔\n", seg, cnt));
// 7. 是否存在单笔超过3000的订单
System.out.println("\n【7> 是否存在大额订单(>¥3000)?】");
boolean hasLarge = orders.stream()
.anyMatch(o -> o.getTotalAmount() > 3000);
Optional<Order> largeOrder = orders.stream()
.filter(o -> o.getTotalAmount() > 3000)
.findFirst();
largeOrder.ifPresent(o ->
System.out.printf(" 是!%s 购买了 %s,金额 ¥%.2f\n",
o.getCustomerName(), o.getProductName(), o.getTotalAmount()));
System.out.println("\n========== 飞翔科技 Q1 数据分析报告完毕 ==========");
}
}
运行输出:
========== 飞翔科技商城订单分析 ==========
【1> 各品类销售额排行】
电子产品: ¥5655.52
服装: ¥885.40
食品: ¥434.40
图书: ¥404.40
【2> 各月度销售额趋势】
2024-01: ¥1895.28
2024-02: ¥4305.28
2024-03: ¥1178.88
【3> 支付方式分布】
微信: ¥3767.28 (51.1%)
支付宝: ¥2837.28 (38.4%)
银行卡: ¥774.88 (10.5%)
【4> 客户消费排行 Top 3】
Frank: ¥3888.88
大翔: ¥683.52
白歌: ¥557.76
...
========== End ==========
六、易错场景
易错1:Stream 不可重用
// 错误示范
Stream<String> stream = list.stream();
stream.forEach(System.out::println); // 第一次消费 OK
stream.forEach(System.out::println); // ← IllegalStateException: stream has already been operated upon or closed
// 正确做法:重新获取 Stream
list.stream().forEach(System.out::println);
list.stream().forEach(System.out::println);
// 或者将中间操作链赋值给 Supplier
Supplier<Stream<String>> supplier = () -> list.stream().filter(s -> s.length() > 2);
supplier.get().forEach(System.out::println); // OK
supplier.get().count(); // OK
易错2:Collectors.toMap 键冲突
// 如果 key 重复,toMap 默认抛异常
List<Order> orders = Arrays.asList(
new Order("ORD001", "大翔", "电子产品", "键盘", 1, 188.88, "2024-01", "微信"),
new Order("ORD001", "大翔", "图书", "Java书", 1, 88.88, "2024-01", "微信")
);
// 错误:重复 key "大翔" → IllegalStateException
// Map<String, Order> map = orders.stream()
// .collect(Collectors.toMap(Order::getCustomerName, Function.identity()));
// 正确:提供 merge function
Map<String, Order> map = orders.stream()
.collect(Collectors.toMap(
Order::getCustomerName,
Function.identity(),
(existing, replacement) -> replacement // 保留后者
));
易错3:并行流中使用非线程安全集合
// 错误示范
List<Integer> list = new ArrayList<>(); // 非线程安全
IntStream.range(0, 10000)
.parallel()
.forEach(list::add); // ← 数据丢失/ArrayIndexOutOfBoundsException
System.out.println(list.size()); // 可能小于 10000
// 正确做法
List<Integer> list = IntStream.range(0, 10000)
.parallel()
.boxed()
.collect(Collectors.toList()); // 线程安全
易错4:sorted 之后 limit 的性能陷阱
// sorted 是全量排序,即使是 sorted(...).limit(5),也会对整个流排序
// 对于大数据集的 Top-N 需求,应该使用 PriorityQueue 而非 Stream
// 或者使用 Google Guava 的 Ordering.greatestOf()
七、面试考点
考点1:map 和 flatMap 的区别?各自什么场景使用?
答:
| map | flatMap | |
|---|---|---|
| 函数签名 | Stream<RStream<R> map(Function<T, R>gt; map(Function<T, R>) | Stream<RStream<R> flatMap(Function<T, Stream<R>>gt; flatMap(Function<T, Stream<R>>) |
| 转换关系 | 一对一 | 一对多(扁平化) |
| 输入输出 | T → R | T → Stream<R> |
| 典型场景 | 提取字段、类型转换 | 拆分字符串为单词、展开嵌套集合 |
// map:每个员工→姓名
employees.stream().map(Emp::getName).collect(toList());
// ["大翔", "白歌", ...]
// flatMap:每个员工→技能列表→合并为一个技能流
employees.stream().flatMap(e -> e.getSkills().stream()).distinct().collect(toList());
// ["Java", "Spring", "MySQL", ...]
考点2:Stream 的中间操作和终止操作有什么区别?惰性求值如何体现?
答:
- 中间操作(filter, map, sorted...)返回 Stream,不触发实际计算,只是在构建一个操作流水线。
- 终止操作(collect, reduce, forEach...)返回非 Stream 的结果,触发整个流水线的执行。
- 惰性求值意味着:在终止操作之前,不会有任何元素被处理。这允许 JVM 对操作链进行优化(如 filter-map 融合)。
- 短路操作(limit、findFirst 等)在满足条件后立即停止处理,与惰性求值配合可大幅提升效率。
考点3:并行流(parallelStream)的内部实现原理是什么?何时使用?
答:
- 并行流底层使用 ForkJoinPool(JDK 7 引入的工作窃取框架)。
- 默认线程数为
Runtime.getRuntime().availableProcessors() - 1,可通过-Djava.util.concurrent.ForkJoinPool.common.parallelism=N调整。 - 数据被 Spliterator 分割为多个子任务并行处理,最后合并结果。
- 适用场景:CPU 密集型、数据量大(>1万)、无 I/O 阻塞、操作是无状态的。
- 不适用场景:I/O 密集型、数据量小(线程切换开销大于收益)、操作依赖外部状态。
考点4:Collectors.groupingBy 的多级分组如何实现?与 SQL 的 GROUP BY 有什么异同?
答: 二级分组就是嵌套 groupingBy:
Map<String, Map<String, List<Emp>>>> = employees.stream()
.collect(Collectors.groupingBy(
Emp::getDept, // 第一级:部门
Collectors.groupingBy( // 第二级:城市
Emp::getCity
)
));
还可以配合下游收集器做聚合:
// 对应 SQL: SELECT dept, AVG(salary) FROM emp GROUP BY dept
Map<String, Double> = employees.stream()
.collect(Collectors.groupingBy(
Emp::getDept,
Collectors.averagingDouble(Emp::getSalary)
));
与 SQL GROUP BY 的差异:Stream 的 groupingBy 支持更灵活的下游收集器组合(mapping、reducing、collectingAndThen),而不仅是聚合函数。
白歌关掉 IDE:"Stream API 是 JDK 8 最革命性的变化。它把 Java 带入了函数式编程的世界。"
大翔深有感触:"以前写一堆 for 循环 + if 判断 + 临时集合,现在一条 pipeline 清晰明了。filter → map → sorted → collect,读代码就像读需求。"
小崔兴奋地跳起来:"我想到一个用法!用 groupingBy 统计每个客户的订单量,然后用 filter 筛选出 VIP 客户!"
朱璐也举手:"还有 toMap 可以用来构建缓存——把 ID 映射到对象,不用再写 for 循环!"
Frank 总监走过来:"看来大家收获不小。下个月的报表系统重构,就用 Stream API 吧。谁来做?"
全组举手:"我来!"
JDK 8 Stream 最佳实践
- 优先使用
collect(Collectors.toList())而非forEach(list::add)- 对于大结果集的 toMap,始终提供 merge function
- 使用
mapToInt/mapToDouble/mapToLong避免装箱开销peek仅用于调试,不要在生产代码中依赖 peek 的副作用- 避免在 Stream 操作中修改外部状态(函数式编程的无副作用原则)
- 并行流需要评估收益——不是所有场景都能更快
JDK 9 Stream 增强(详见第17章「Stream 增强」)
JDK 9 为 Stream 新增了四个方法,填补了 JDK 8 的边界控制空白:
takeWhile(Predicate)— 从开头取元素,直到遇到第一个不满足条件的停止(适合已排序数据截取前缀)dropWhile(Predicate)— 从开头丢弃元素,直到遇到第一个不满足条件的开始保留(适合跳过数据头部的无用行)iterate(seed, Predicate, UnaryOperator)— 带终止条件的无限流生成器,替代limit()的硬编码ofNullable(T)— 创建 0 或 1 个元素的流,null 时为空流(配合flatMap优雅过滤 null)这些方法在处理已排序数据、日志解析、配置过滤等场景中非常实用。
JDK 12/16 Stream 增强(详见第18章「Stream增强」)
JDK 12 和 JDK 16 为 Stream API 新增了三个高频工具方法:
Collectors.teeing(Collector, Collector, BiFunction)— 双下游收集器(JDK 12),一条流同时产出两项统计(如计数 + 求和)Stream.toList()— 直接返回不可变 List(JDK 16),替代.collect(Collectors.toList())Stream.mapMulti(BiConsumer)— 一对多映射替代方案(JDK 16),避免创建中间 Stream,适合每个元素映射到少量结果的场景配合 JDK 9 的
takeWhile/dropWhile和 JDK 8 的 Stream 基础,Stream API 的工具链至此已相当完整。