Java 基础体系 · 第 61/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java Collector 深入:归约、分组、下游收集器、并行和自定义
Collector 是 Java Stream 中描述“如何把一串元素汇总成结果”的抽象。它不仅用于把流转换成 List、Set 或 Map,还用于表达求和、统计、分组、嵌套分组、连接字符串,以及自定义的可并行归约算法。
要真正理解 Collector,需要同时掌握四个层次:
- 归约:如何把多个元素逐步合并成一个结果。
- 收集器协议:
supplier、accumulator、combiner、finisher如何协作。 - 下游收集器:分组后如何继续对每个组进行映射、过滤、统计或再次分组。
- 并行语义:拆分后的局部结果如何合并,以及什么条件下自定义收集器才是正确的。
一、从归约理解 Collector
1.1 归约的基本形式
归约(reduction)是把多个值按照某种规则合并成一个值。
例如,对整数求和:
int result = Stream.of(1, 2, 3, 4)
.reduce(0, Integer::sum);
计算过程可以写成:
0 + 1 = 1
1 + 2 = 3
3 + 3 = 6
6 + 4 = 10
这里有三个重要对象:
- 初始值
0 - 单元素累加规则
Integer::sum - 最终结果
10
Stream.reduce 适合“单值、直接合并”的场景,但很多结果并不是简单的一个不可变值。例如:
- 收集成
List - 按城市分组
- 每组计算平均值
- 同时得到数量、总和、最小值和最大值
- 每组收集成一个新的集合
这时就需要一个可变的中间容器,以及描述容器如何创建、更新、合并和转换的 Collector。
1.2 reduce 与 collect 的区别
reduce 通常围绕“值到值”的合并:
int sum = Stream.of(1, 2, 3, 4)
.reduce(0, Integer::sum);
collect 则围绕“元素到可变结果容器”的累积:
List<Integer> numbers = Stream.of(1, 2, 3, 4)
.collect(Collectors.toList());
两者的核心区别不是“一个能并行、一个不能并行”,而是结果状态的表达方式不同:
| 维度 | reduce |
collect |
|---|---|---|
| 典型结果 | 一个值 | 可变容器再转换为结果 |
| 状态更新 | 通过返回新值 | 修改累加容器 |
| 常见用途 | 求和、最值、字符串合并 | 列表、映射、分组、统计 |
| 并行关键 | 合并值必须满足结合律 | 容器合并必须满足 Collector 协议 |
例如,使用 reduce 拼接字符串:
String result = Stream.of("A", "B", "C")
.reduce("", String::concat);
逻辑上可行,但每次字符串拼接可能创建新对象。更适合使用收集器让每个局部结果写入 StringBuilder,最后再转换成 String。
二、Collector 的五个组成部分
Collector<T, A, R> 的三个类型参数分别表示:
T:流中元素类型A:中间累加容器类型R:最终结果类型
其核心接口可抽象为:
public interface Collector<T, A, R> {
Supplier<A> supplier();
BiConsumer<A, T> accumulator();
BinaryOperator<A> combiner();
Function<A, R> finisher();
Set<Characteristics> characteristics();
}
例如,一个把字符串收集成逗号分隔文本的收集器,可以是:
T = String
A = StringBuilder
R = String
A 和 R 不必相同,这正是 finisher 存在的原因。
2.1 supplier:创建空容器
supplier 负责创建新的累加容器:
Supplier<StringBuilder> supplier = StringBuilder::new;
在串行流中,通常只需要一个容器。在并行流中,Stream 实现可能创建多个容器,分别处理不同的数据分片。
因此,supplier 返回的容器必须代表“尚未累积任何元素”的合法初始状态。
2.2 accumulator:把一个元素加入容器
accumulator 接收两个参数:
- 当前累加容器
- 流中的一个元素
例如:
BiConsumer<StringBuilder, String> accumulator =
(builder, value) -> {
if (builder.length() > 0) {
builder.append(',');
}
builder.append(value);
};
它直接修改容器,而不是返回一个新容器。
需要注意:使用可变容器并不意味着可以在多个线程之间无保护地共享它。普通收集器的并行执行通常是“每个线程使用自己的容器,最后再合并”,而不是所有线程同时修改同一个容器。
2.3 combiner:合并局部容器
并行流会把输入拆成多个分片。假设输入是:
[A, B, C, D]
可能被拆分为:
左分片:[A, B]
右分片:[C, D]
两个分片分别累积后得到:
左容器:A,B
右容器:C,D
combiner 需要把它们合并:
BinaryOperator<StringBuilder> combiner =
(left, right) -> {
if (left.length() > 0 && right.length() > 0) {
left.append(',');
}
left.append(right);
return left;
};
结果为:
A,B,C,D
combiner 不只是并行优化点,它是 Collector 正确性的核心组成部分。即使当前代码使用串行流,错误的 combiner 也会在未来切换到并行流时暴露。
2.4 finisher:把中间容器转换成最终结果
如果中间容器就是最终结果,可以直接返回它;否则需要 finisher:
Function<StringBuilder, String> finisher = StringBuilder::toString;
例如:
A、B、C
在收集完成后转换为不可变的 String。
finisher 通常只在所有累积和合并完成后调用一次。它不应该重新读取原始流,也不应该依赖元素处理顺序之外的隐含状态。
2.5 characteristics:声明收集器能力
常见特征有三个:
IDENTITY_FINISH
表示中间类型 A 可以直接作为结果类型 R 返回,等价于:
A == R
例如收集成 ArrayList 时,通常可以直接把累加容器作为结果返回。
如果存在类似 StringBuilder::toString 的转换,则不能声明 IDENTITY_FINISH。
UNORDERED
表示结果不依赖流的遇到顺序(encounter order)。
例如,集合中的元素顺序不重要时,可以声明此特征。但“结果类型是无序集合”并不自动意味着所有算法都可以忽略顺序;必须确认业务语义确实不依赖顺序。
CONCURRENT
表示多个线程可以同时调用同一个结果容器的累加操作。
这是一项很强的保证。仅仅因为容器是并发容器,例如 ConcurrentHashMap,并不代表整个收集器就可以声明 CONCURRENT。累加逻辑、容器状态和下游操作都必须满足并发要求。
对于有序并行流,即使收集器声明了 CONCURRENT,Stream 也不一定会让多个线程并发更新同一个容器;遇到顺序和 UNORDERED 特征会共同影响执行策略。
三、Collector 必须满足的正确性条件
Collector 的正确性可以通过“串行累积”和“拆分后累积再合并”是否等价来理解。
设:
S():supplier创建的空容器A(a, t):accumulator把元素t加入容器aC(a1, a2):combiner合并两个容器F(a):finisher产生最终结果
3.1 空容器必须是单位状态
把一个空容器合并到已有容器中,结果应该等价于原容器:
C(a, S()) ≈ a
C(S(), a) ≈ a
这里的“≈”表示最终结果等价,不一定要求对象身份相同。
如果一个空容器会向结果中添加额外分隔符、虚拟元素或默认业务数据,就可能破坏这个条件。
3.2 合并必须具有结合性
对于三个局部容器:
C(C(a, b), c) ≈ C(a, C(b, c))
这就是结合律。
例如整数求和满足:
(1 + 2) + 3 = 1 + (2 + 3)
而减法不满足:
(10 - 5) - 2 = 3
10 - (5 - 2) = 7
因此,不能把普通减法直接作为并行归约的合并规则。
3.3 累积与合并必须相容
把一个元素先加入某个局部容器,再与另一个容器合并,应该等价于先合并容器,再加入这个元素:
C(A(a, t), b) ≈ A(C(a, b), t)
这条规则说明了为什么 combiner 不能随意定义。它必须理解 accumulator 所维护的状态,否则局部结果合并后可能不再代表原始数据的真实汇总。
3.4 顺序敏感的收集器必须保持左到右
对有序流而言,拆分合并通常需要保持分片的先后关系:
C(左分片结果, 右分片结果)
不能反过来调用:
C(右分片结果, 左分片结果)
例如字符串连接:
[A, B] + [C, D] = A,B,C,D
反向合并则会得到:
C,D,A,B
所以字符串连接收集器不能声明 UNORDERED。
四、标准归约收集器
Collectors 提供了多种归约工具。
4.1 counting、summing 与 averaging
long count = Stream.of(10, 20, 30)
.collect(Collectors.counting());
int sum = Stream.of(10, 20, 30)
.collect(Collectors.summingInt(Integer::intValue));
double average = Stream.of(10, 20, 30)
.collect(Collectors.averagingInt(Integer::intValue));
它们分别表示:
counting():元素数量summingInt、summingLong、summingDouble:数值求和averagingInt、averagingLong、averagingDouble:平均值
空流的平均值返回 0.0,因为这些收集器以计数和总和表示状态,空状态的计数为零,最终按照 API 定义返回零值。
如果需要区分“没有元素”和“平均值为零”,应使用其他状态表示或先判断流是否为空。
4.2 minBy、maxBy 与 reducing
Optional<Integer> maximum = Stream.of(4, 2, 9, 1)
.collect(Collectors.maxBy(Integer::compareTo));
结果是 Optional<Integer>,因为空流没有最大值。
reducing 可表达更一般的归约:
int sum = Stream.of(1, 2, 3, 4)
.collect(Collectors.reducing(0, Integer::sum));
也可以指定映射函数:
int totalLength = Stream.of("Java", "Stream", "Collector")
.collect(Collectors.reducing(
0,
String::length,
Integer::sum
));
这里每个字符串先映射为长度,再按照整数加法归约。
不过,顶层流的简单归约通常直接使用 Stream.reduce 更清晰;Collectors.reducing 的重要价值在于作为下游收集器嵌入分组操作。
4.3 summarizingInt:一次维护多个统计量
IntSummaryStatistics statistics = Stream.of(10, 20, 30)
.collect(Collectors.summarizingInt(Integer::intValue));
System.out.println(statistics.getCount()); // 3
System.out.println(statistics.getSum()); // 60
System.out.println(statistics.getMin()); // 10
System.out.println(statistics.getMax()); // 30
System.out.println(statistics.getAverage()); // 20.0
它的内部状态可以理解为:
count = 3
sum = 60
min = 10
max = 30
并行时,每个分片各自维护一组统计量,然后通过:
count 相加
sum 相加
min 取较小值
max 取较大值
进行合并。这些操作分别满足适合并行归约的结合性要求。
五、分组:groupingBy 的数据流
groupingBy 的基本形式是:
Collectors.groupingBy(classifier)
其中 classifier 是分类函数:
Function<T, K>
它把每个元素 T 映射成分组键 K。
例如:
Map<String, List<String>> namesByInitial =
Stream.of("Alice", "Bob", "Anna", "Bill")
.collect(Collectors.groupingBy(
name -> name.substring(0, 1)
));
结果的逻辑内容是:
A -> [Alice, Anna]
B -> [Bob, Bill]
数据流可以表示为:
flowchart LR
E[流元素] --> C[分类函数 classifier]
C --> K[分组键]
K --> M[Map 中对应的组]
E --> D[下游收集器]
D --> V[该组的结果]
M --> V
对每个元素,groupingBy 会:
- 调用
classifier得到键; - 查找该键对应的组;
- 如果组不存在,创建一个新的下游累加容器;
- 把当前元素交给该组的下游收集器;
- 最终把所有键和组结果放入
Map。
groupingBy 的默认下游收集器是 toList(),因此:
groupingBy(classifier)
等价于“按键分组,每组收集成列表”。
5.1 分组键的限制
分类函数不应返回 null。groupingBy 要求分类结果作为合法的映射键使用,标准实现通常会对 null 分类结果拒绝处理,而不是自动创建一个 null 分组。
如果业务允许空键,应显式转换:
Map<String, List<String>> grouped =
values.stream()
.collect(Collectors.groupingBy(
value -> value == null ? "<null>" : value
));
这比依赖具体 Map 实现对 null 的处理更明确。
5.2 指定 Map 类型
默认返回的 Map 类型不应被业务代码当作固定实现依赖。如果需要保持分组键的插入顺序,可以指定 Map 工厂:
Map<String, List<String>> grouped =
Stream.of("B1", "A1", "B2")
.collect(Collectors.groupingBy(
value -> value.substring(0, 1),
LinkedHashMap::new,
Collectors.toList()
));
此时结果的键顺序通常按首次出现顺序保存:
B -> [B1, B2]
A -> [A1]
这里必须区分两件事:
LinkedHashMap保持映射的迭代顺序;- 分组列表是否保持流的遇到顺序,还取决于流和下游收集器的语义。
六、下游收集器:分组后继续计算
下游收集器是 groupingBy 的第二个核心参数:
groupingBy(classifier, downstream)
它不再决定“如何分组”,而是决定“每个组内部如何收集”。
6.1 每组计数
Map<String, Long> countByInitial =
Stream.of("Alice", "Bob", "Anna", "Bill")
.collect(Collectors.groupingBy(
name -> name.substring(0, 1),
Collectors.counting()
));
逻辑结果:
A -> 2
B -> 2
与先分组为 List 再调用 size() 相比,下游 counting() 不需要为每个组保存所有元素,状态更小。
6.2 每组求和
record Order(String city, String id, int amount) {}
List<Order> orders = List.of(
new Order("Beijing", "A001", 120),
new Order("Beijing", "A002", 80),
new Order("Shanghai", "B001", 200),
new Order("Shanghai", "B002", 50)
);
Map<String, Integer> totalByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.summingInt(Order::amount)
));
结果:
Beijing -> 200
Shanghai -> 250
这里每个城市的组状态是一个整数求和状态,而不是订单列表。
6.3 每组统计最小值、最大值和平均值
Map<String, IntSummaryStatistics> statisticsByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.summarizingInt(Order::amount)
));
IntSummaryStatistics beijing =
statisticsByCity.get("Beijing");
System.out.println(beijing.getCount()); // 2
System.out.println(beijing.getSum()); // 200
System.out.println(beijing.getMin()); // 80
System.out.println(beijing.getMax()); // 120
System.out.println(beijing.getAverage()); // 100.0
这相当于为每个分组维护独立的统计状态:
Beijing:
count = 2
sum = 200
min = 80
max = 120
6.4 mapping:先映射,再收集
如果每组只需要订单编号,而不是完整订单,可以使用:
Map<String, List<String>> idsByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.mapping(
Order::id,
Collectors.toList()
)
));
结果:
Beijing -> [A001, A002]
Shanghai -> [B001, B002]
mapping 的处理顺序是:
Order
-> Order::id
-> 下游 toList()
因此不必先构造:
Map<String, List<Order>>
再遍历每个列表转换 ID。
6.5 filtering:在组内过滤
Map<String, List<String>> highValueIdsByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.filtering(
order -> order.amount() >= 100,
Collectors.mapping(
Order::id,
Collectors.toList()
)
)
));
结果:
Beijing -> [A001]
Shanghai -> [B001]
filtering 与在分组前调用 filter 并不总是等价。
例如:
orders.stream()
.filter(order -> order.amount() >= 100)
.collect(Collectors.groupingBy(Order::city));
如果某个城市没有高价值订单,这个城市的键会完全不存在。
而使用下游 filtering:
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.filtering(
order -> order.amount() >= 100,
Collectors.toList()
)
));
该城市仍可能存在,只是对应空列表。这种差异在报表和统计接口中非常重要。
6.6 flatMapping:每个元素展开多个值
假设订单有多个标签:
record TaggedOrder(String city, List<String> tags) {}
可以把每个城市下所有订单的标签展平成一个集合:
Map<String, Set<String>> tagsByCity =
taggedOrders.stream()
.collect(Collectors.groupingBy(
TaggedOrder::city,
Collectors.flatMapping(
order -> order.tags().stream(),
Collectors.toSet()
)
));
数据流为:
TaggedOrder
-> tags 列表
-> 多个 tag 流元素
-> 当前城市的 Set
如果直接使用 mapping,得到的会是 List<List<String>> 或类似嵌套结构,而不是扁平的标签集合。
6.7 collectingAndThen:收集完成后转换
例如每组收集后变成不可修改列表:
Map<String, List<String>> immutableIdsByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.collectingAndThen(
Collectors.mapping(Order::id, Collectors.toList()),
List::copyOf
)
));
处理顺序为:
订单
-> 按城市分组
-> 每组映射为 ID
-> 收集为 List
-> List.copyOf 转换为不可修改列表
collectingAndThen 的结果通常不能声明 IDENTITY_FINISH,因为下游还存在最后的转换函数。
6.8 teeing:同一组同时计算两个结果
teeing 可以把两个下游收集器应用到同一批元素,再合并两个结果:
record CitySummary(long count, int total) {}
Map<String, CitySummary> summaryByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.teeing(
Collectors.counting(),
Collectors.summingInt(Order::amount),
CitySummary::new
)
));
结果逻辑上是:
Beijing -> CitySummary(count=2, total=200)
Shanghai -> CitySummary(count=2, total=250)
它避免了先收集完整订单列表,再单独进行两次遍历。
teeing 的两个下游都会看到每个组中的每个元素,但它们各自维护独立状态。两个结果最终由第三个合并函数组合成最终结果。
七、一个可运行的端到端示例
下面的程序可使用 Java 25 编译运行,不依赖预览特性:
import java.util.IntSummaryStatistics;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.stream.Collectors;
public class CollectorDemo {
record Order(String city, String id, int amount) {}
record CitySummary(long count, int total) {}
public static void main(String[] args) {
List<Order> orders = List.of(
new Order("Beijing", "A001", 120),
new Order("Beijing", "A002", 80),
new Order("Shanghai", "B001", 200),
new Order("Shanghai", "B002", 50)
);
Map<String, IntSummaryStatistics> statisticsByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
LinkedHashMap::new,
Collectors.summarizingInt(Order::amount)
));
Map<String, List<String>> highValueIdsByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
LinkedHashMap::new,
Collectors.filtering(
order -> order.amount() >= 100,
Collectors.mapping(
Order::id,
Collectors.toList()
)
)
));
Map<String, CitySummary> summaryByCity =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
LinkedHashMap::new,
Collectors.teeing(
Collectors.counting(),
Collectors.summingInt(Order::amount),
CitySummary::new
)
));
System.out.println(statisticsByCity);
System.out.println(highValueIdsByCity);
System.out.println(summaryByCity);
Objects.requireNonNull(statisticsByCity.get("Beijing"));
}
}
典型输出的逻辑内容为:
{
Beijing=IntSummaryStatistics{count=2, sum=200, min=80, average=100.000000, max=120},
Shanghai=IntSummaryStatistics{count=2, sum=250, min=50, average=125.000000, max=200}
}
{
Beijing=[A001],
Shanghai=[B001]
}
{
Beijing=CitySummary[count=2, total=200],
Shanghai=CitySummary[count=2, total=250]
}
IntSummaryStatistics 的 toString() 展示形式属于类的输出实现细节,业务代码不应依赖它的具体文本格式;实际使用时应调用 getCount()、getSum() 等方法。
八、多级分组
下游收集器本身也可以是另一个 groupingBy。
例如,先按城市,再按金额是否达到 100:
Map<String, Map<Boolean, List<Order>>> result =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.partitioningBy(
order -> order.amount() >= 100
)
));
逻辑结构:
城市
-> true / false
-> 订单列表
结果类似:
Beijing:
true -> [A001]
false -> [A002]
Shanghai:
true -> [B001]
false -> [B002]
partitioningBy 专门处理布尔分类。它与普通 groupingBy 的语义不同:分区结果通常包含 true 和 false 两个分区,即使其中一个分区为空。
多级分组也可以继续结合映射和统计:
Map<String, Map<Boolean, Integer>> totals =
orders.stream()
.collect(Collectors.groupingBy(
Order::city,
Collectors.partitioningBy(
order -> order.amount() >= 100,
Collectors.summingInt(Order::amount)
)
));
此时最终结构是:
Map<String, Map<Boolean, Integer>>
不是:
Map<String, Map<Boolean, List<Order>>>
每一层下游收集器都会改变对应层级的状态类型。
九、Map 收集器与重复键
分组并不是唯一产生 Map 的方式。对于“一对一”或“显式处理重复键”的场景,应考虑 toMap。
Map<String, Integer> amountById =
orders.stream()
.collect(Collectors.toMap(
Order::id,
Order::amount
));
如果两个订单产生相同 ID,默认实现会抛出 IllegalStateException。这不是异常行为,而是 toMap 在无法确定如何处理重复键时主动拒绝丢失数据。
可以显式提供合并函数:
Map<String, Integer> totalAmountById =
orders.stream()
.collect(Collectors.toMap(
Order::id,
Order::amount,
Integer::sum
));
此时重复 ID 会执行:
旧值 + 新值
也可以指定 Map 工厂:
Map<String, Integer> ordered =
orders.stream()
.collect(Collectors.toMap(
Order::id,
Order::amount,
Integer::sum,
LinkedHashMap::new
));
应区分:
groupingBy:一个键天然对应一组值;toMap:一个键期望对应一个值,重复键必须由调用者决定如何处理。
如果把本应使用 groupingBy 的问题写成 toMap,常见失败表现就是重复键异常;如果为了避免异常随意使用覆盖策略,则可能静默丢失数据。
十、并行 Collector 的执行模型
10.1 串行执行
串行流大致可以理解为:
supplier()
-> accumulator(element1)
-> accumulator(element2)
-> accumulator(element3)
-> finisher()
例如:
String result = Stream.of("A", "B", "C")
.collect(joiningCollector());
通常只有一个累加容器。
10.2 并行执行
并行流可能执行为:
输入:[A, B, C, D]
分片一:[A, B] -> 容器一
分片二:[C, D] -> 容器二
容器一 + 容器二 -> combiner
最终容器 -> finisher
更大数据量下可能形成树状合并:
[A] + [B] -> AB
[C] + [D] -> CD
AB + CD -> ABCD
因此,不能只验证“从左到右串行执行时结果正确”,还必须验证任意合法拆分和合并都正确。
10.3 并行不等于更快
并行收集会产生额外成本:
- 拆分源数据;
- 创建多个累加容器;
- 合并局部结果;
- 线程调度;
- 可能的排序或顺序保持。
如果输入规模小、下游操作简单,额外成本可能大于收益。
此外,某些收集器的结果合并本身很昂贵。例如每次都把右侧大列表复制到左侧,可能导致大量内存复制。并行执行是否有收益,取决于数据规模、拆分能力、累积成本、合并成本和线程池资源,而不能仅凭 parallel() 判断。
10.4 有序结果与并行
对于有序流和顺序敏感的收集器,Stream 实现通常需要维护分片的相对顺序。
例如:
String result = Stream.of("A", "B", "C", "D")
.parallel()
.collect(joiningCollector());
如果收集器没有声明 UNORDERED,正确结果应保持:
A,B,C,D
但这不是因为 combiner 可以任意交换左右参数,而是因为有序流的合并必须保留分片顺序。
如果业务不关心顺序,可以使用无序语义;这可能为实现更多并发优化空间,但不能为了“可能更快”而错误声明 UNORDERED。
十一、groupingBy 与 groupingByConcurrent
11.1 groupingBy
Map<String, List<Order>> grouped =
orders.parallelStream()
.collect(Collectors.groupingBy(Order::city));
它可以用于并行流,但并不意味着所有线程都会直接并发写入同一个 Map。实现可以先构造局部结果,再进行 Map 合并。
对于分组数量较多、下游处理较重的场景,局部构造和合并可能是合理的;对于某些数据分布,合并大量 Map 也可能成为瓶颈。
11.2 groupingByConcurrent
Map<String, List<Order>> grouped =
orders.parallelStream()
.collect(Collectors.groupingByConcurrent(
Order::city
));
它返回并发 Map,并声明并发收集相关特征,适合允许无序处理的并行分组场景。
但需要注意几个边界:
- 返回类型是
ConcurrentMap更准确地反映其语义;如果只声明为Map,仍可以使用,但会隐藏并发特性。 - 结果的键遍历顺序不应被依赖。
- 并发 Map 不等于组内列表是任意时刻都适合外部并发修改。
- 下游收集器的状态仍必须能在该收集模型下被正确维护。
如果业务要求严格的遇到顺序,groupingByConcurrent 通常不是首选。
十二、线程安全与副作用
下面的写法存在风险:
List<Integer> shared = new ArrayList<>();
List<Integer> result = Stream.of(1, 2, 3, 4)
.parallel()
.peek(shared::add)
.collect(Collectors.toList());
shared 被多个线程修改,可能出现:
- 数据丢失;
- 竞态条件;
- 内部结构损坏;
- 结果偶发不一致。
正确做法是让 Collector 管理自己的局部容器:
List<Integer> result = Stream.of(1, 2, 3, 4)
.parallel()
.collect(Collectors.toList());
普通 Collector 的设计假设是:
线程一 -> 容器一
线程二 -> 容器二
线程三 -> 容器三
最后统一 combiner
而不是:
线程一、线程二、线程三 -> 共享普通 ArrayList
此外,accumulator、combiner 和 finisher 应尽量避免修改外部状态。外部副作用会让结果依赖执行模式,导致串行时正常、并行时失败。
十三、自定义 Collector:完整实现
下面实现一个把字符串连接成逗号分隔文本的收集器:
import java.util.List;
import java.util.Objects;
import java.util.stream.Collector;
import java.util.stream.Stream;
public class CustomCollectorDemo {
static Collector<String, StringBuilder, String> joiningWithComma() {
return Collector.of(
StringBuilder::new,
(builder, value) -> {
Objects.requireNonNull(value, "value");
if (builder.length() > 0) {
builder.append(',');
}
builder.append(value);
},
(left, right) -> {
if (left.length() > 0 && right.length() > 0) {
left.append(',');
}
left.append(right);
return left;
},
StringBuilder::toString
);
}
public static void main(String[] args) {
String serial = Stream.of("Java", "25", "Stream")
.collect(joiningWithComma());
String parallel = Stream.of("Java", "25", "Stream")
.parallel()
.collect(joiningWithComma());
System.out.println(serial); // Java,25,Stream
System.out.println(parallel); // Java,25,Stream
}
}
这个收集器的生命周期如下:
supplier:
创建 StringBuilder
accumulator:
依次把单个字符串追加到 StringBuilder
combiner:
把右侧 StringBuilder 追加到左侧,并补充分隔符
finisher:
StringBuilder.toString()
它没有声明任何特征:
- 不能声明
IDENTITY_FINISH,因为StringBuilder需要转换为String; - 不能声明
UNORDERED,因为字符串顺序影响结果; - 不能声明
CONCURRENT,因为单个StringBuilder不是并发安全容器,而且该算法也没有声明共享累积语义。
13.1 为什么这个自定义收集器可以并行
假设数据被拆成:
左:[Java, 25]
右:[Stream]
局部容器分别为:
左:Java,25
右:Stream
调用 combiner(left, right) 后:
Java,25,Stream
只要合并始终按照左分片、右分片的顺序进行,就保持了有序流的遇到顺序。
如果数据为空:
String empty = Stream.<String>empty()
.collect(joiningWithComma());
结果是空字符串,因为:
supplier创建空StringBuilder;- 没有元素,因此不调用
accumulator; finisher对空容器调用toString()。
13.2 一个错误的 combiner
下面的合并方式是错误的:
(left, right) -> {
left.insert(0, right);
return left;
}
假设:
左:A,B
右:C,D
结果变成:
C,D,A,B
串行执行时可能完全看不出问题,因为串行执行通常不会调用 combiner。一旦切换到并行流,顺序就会错误。
另一个错误是忘记处理空容器:
(left, right) -> left.append(',').append(right)
可能产生:
,A,B
或:
A,B,
这违反了空容器应作为单位状态的要求。
十四、IDENTITY_FINISH 的误用
假设自定义收集器使用:
StringBuilder
作为中间容器,使用:
String
作为最终结果:
A = StringBuilder
R = String
就不能写成:
Collector.Characteristics.IDENTITY_FINISH
因为 StringBuilder 不能直接当作 String 返回。
错误声明可能导致实现跳过本应执行的 finisher,最终返回错误类型,或者产生未定义的行为。特征不是注释,而是对 Stream 实现的语义承诺;声明错误比不声明更危险。
十五、何时使用自定义 Collector
标准收集器已经覆盖了绝大多数常见需求:
toListtoSettoMapgroupingBypartitioningBymappingfilteringflatMappingcountingsummingsummarizingjoiningteeing
自定义 Collector 适合以下情况:
- 累积状态需要多个字段;
- 中间状态和最终结果类型不同;
- 需要定义特殊的分片合并逻辑;
- 现有收集器组合无法清晰表达算法;
- 需要让算法天然适配串行和并行收集。
例如,一个“统计总数、非空数量和字符总长度”的收集器,可以使用一个专门的可变状态类:
static final class TextStats {
long total;
long nonEmpty;
long characters;
void add(String value) {
total++;
if (!value.isEmpty()) {
nonEmpty++;
}
characters += value.length();
}
TextStats merge(TextStats other) {
total += other.total;
nonEmpty += other.nonEmpty;
characters += other.characters;
return this;
}
}
其合并规则是字段级别的结合运算:
total = total1 + total2
nonEmpty = nonEmpty1 + nonEmpty2
characters = characters1 + characters2
这类状态比先收集全部字符串再遍历统计更节省内存,也更适合并行。
十六、常见误解与失败表现
16.1 collect(toList()) 一定返回可修改列表
不应依赖 Collectors.toList() 返回的具体列表类型或可修改性。API 只保证收集成列表,不保证具体实现和修改行为。
如果需要不可修改列表,可以使用:
List<String> result = stream.collect(Collectors.toUnmodifiableList());
或者:
List<String> result = stream.toList();
如果需要明确构造一个可修改的 ArrayList:
List<String> result = stream.collect(Collectors.toCollection(ArrayList::new));
这三者的语义不同,不能仅根据方法名互换。
16.2 toSet 会保留顺序
Collectors.toSet() 不应被用来表达顺序要求。需要保持插入顺序时:
Set<String> result = stream.collect(
Collectors.toCollection(LinkedHashSet::new)
);
需要排序时:
Set<String> result = stream.collect(
Collectors.toCollection(TreeSet::new)
);
“去重”和“保持顺序”是两个不同要求。
16.3 groupingBy 会自动处理 null 键
分类函数返回 null 可能导致收集失败。诊断时应首先检查:
Function<T, K> classifier
是否对所有输入都返回非空键,而不是先怀疑 Map 实现。
16.4 parallel() 能修复错误的 Collector
不能。并行只会更快暴露以下问题:
combiner不满足结合性;- 合并左右顺序错误;
- 累加器共享外部可变状态;
- 错误声明
CONCURRENT; - 下游容器被不安全地共享;
- 计算依赖元素处理顺序,却声明了
UNORDERED。
如果串行正确、并行错误,应优先对 supplier、accumulator、combiner 和特征声明进行最小可重复验证。
十七、诊断自定义 Collector 的方法
可以构造同一输入的串行和并行结果进行比较:
List<String> input = List.of("A", "B", "C", "D", "E");
String serial = input.stream()
.collect(joiningWithComma());
String parallel = input.parallelStream()
.collect(joiningWithComma());
if (!serial.equals(parallel)) {
throw new AssertionError(
"串行和并行结果不一致: " + serial + " / " + parallel
);
}
还应测试边界输入:
List.of()
List.of("A")
List.of("A", "B")
包含重复值
包含空字符串
包含 null(如果业务允许)
尤其要检查:
- 空容器是否产生额外字符;
- 单元素是否多出分隔符;
- 多次合并是否满足结合律;
- 并行结果是否保持所需顺序;
- 异常是否能从累加函数和合并函数正常传播。
如果 accumulator 抛出异常,collect 不会把它转换成业务成功结果;异常会向调用方传播。并行执行中,其他任务可能已经完成部分累积,但这些局部容器不会作为正常结果交付给调用方。不要把 Collector 当作事务机制,也不要在累加器中执行无法回滚的外部副作用。
十八、生产取舍
Collector 的选择首先应由结果语义决定:
- 要一组元素,使用
toList、toSet或toCollection; - 要按键得到多值,使用
groupingBy; - 要按键得到单值,使用
toMap并明确重复键策略; - 要组内映射,使用
mapping; - 要组内过滤,使用
filtering; - 要展开嵌套集合,使用
flatMapping; - 要多个统计结果,使用
summarizing或teeing; - 要并行共享累积,只有在确实满足并发协议时才考虑
CONCURRENT。
其次要审查内存行为。groupingBy(..., toList()) 会保存所有元素;counting()、summingInt() 和 summarizingInt() 则只保存汇总状态。如果最终只需要总数,就不应先收集完整列表再计算数量。
最后要审查顺序语义。顺序不是性能选项,而是结果定义的一部分。字符串拼接、列表收集、按时间先后构造结果时,必须保持遇到顺序;只有业务明确不依赖顺序时,才可以采用无序收集语义。
Collector 的本质是一个可组合的归约协议:元素被累积到局部状态,局部状态可以被正确合并,最终状态再转换为结果。理解这套协议后,groupingBy、下游收集器、并行 Stream 和自定义 Collector 就不再是互相孤立的 API,而是同一个归约模型在不同结果结构上的具体表达。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java 并发集合:ConcurrentHashMap、CopyOnWrite、BlockingQueue 和边界
- 下一篇:Java 线程生命周期:创建、状态、Interrupt、守护线程和协作退出
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论