Collector / Collectors 收集器源码精读
概述
Stream.collect(Collector) 是流的终点:把流中的元素累积成集合、字符串、统计值等。Collector 由四个函数(supplier/accumulator/combiner/finisher)+ 一组特性构成,Collectors 提供常用实现(toList/groupingBy/joining 等)。理解其四函数模型,既能读懂内置收集器,也能自定义高性能收集器。本文基于 OpenJDK 21 源码拆解。
一、Collector 的五个核心方法
java
// java.util.stream.Collector
public interface Collector<T, A, R> {
Supplier<A> supplier(); // ① 创建可变容器
BiConsumer<A, T> accumulator(); // ② 元素累积进容器
BinaryOperator<A> combiner(); // ③ 合并两个容器(并行)
Function<A, R> finisher(); // ④ 容器 → 最终结果
Set<Characteristics> characteristics(); // ⑤ 特性声明
}四函数模型:
T:流中元素类型
A:中间可变容器类型(如 ArrayList、StringBuilder)
R:最终结果类型(可与 A 相同)
执行顺序(串行):
A a = supplier.get(); // 建容器
for (T t : stream) {
accumulator.accept(a, t); // 逐个累积
}
R r = finisher.apply(a); // 转换结果
执行顺序(并行):
每段流各自建容器累积 → combiner 两两合并 → finisher二、Characteristics 的 3 个枚举
java
enum Characteristics {
CONCURRENT, // 容器可被多线程并发累积(无锁/线程安全容器)
UNORDERED, // 收集结果与元素顺序无关
IDENTITY_FINISH // finisher 是恒等函数(容器即结果,跳过转换)
}| 特性 | 语义 | 优化 |
|---|---|---|
CONCURRENT | 容器线程安全(如 ConcurrentMap) | 并行时共享同一容器,跳过 combiner |
UNORDERED | 结果无序 | 并行时无需保序合并 |
IDENTITY_FINISH | finisher 恒等(Function.identity()) | 直接返回容器,省一次转换 |
java
// 内置判断逻辑(Collectors 内部)
private static Set<Collector.Characteristics> CH_ID
= Collections.unmodifiableSet(EnumSet.of(Collector.Characteristics.IDENTITY_FINISH));三、Collectors.toList() 的实现
java
public static <T> Collector<T, ?, List<T>> toList() {
return new CollectorImpl<>(
(Supplier<List<T>>) ArrayList::new, // supplier:新建 ArrayList
List::add, // accumulator:add 元素
(left, right) -> { left.addAll(right); return left; }, // combiner:合并
CH_ID // 特性:恒等 finish
);
}toList 关键点:
A = ArrayList(可变容器)
R = List(finisher 恒等 → 直接返回容器)
combiner:并行时 left.addAll(right)
→ 底层就是"ArrayList + add + addAll"java
// CollectorImpl 统一包装(Collectors 内部)
static class CollectorImpl<T, A, R> implements Collector<T, A, R> {
private final Supplier<A> supplier;
private final BiConsumer<A, T> accumulator;
private final BinaryOperator<A> combiner;
private final Function<A, R> finisher;
private final Set<Characteristics> characteristics;
}四、groupingBy 分组收集
java
public static <T, K> Collector<T, ?, Map<K, List<T>>> groupingBy(Function<? super T, ? extends K> classifier) {
return groupingBy(classifier, toList()); // 下游收集器默认 toList
}
public static <T, K, A, D> Collector<T, ?, Map<K, D>> groupingBy(
Function<? super T, ? extends K> classifier,
Collector<? super T, A, D> downstream) {
return groupingBy(classifier, HashMap::new, downstream); // 默认 HashMap
}groupingBy 三步流程:
① classifier.apply(t):计算每个元素的分组键 key
② 容器 Map 中取该 key 的"下游容器"
不存在 → downstream.supplier().get() 新建并放入
存在 → 直接累积
accumulator.accept(subContainer, t) // 元素进下游容器
③ finisher:Map 的值逐个转最终结果(downstream.finisher)
→ 转不可变 Map(Collections.unmodifiableMap)java
// 内部核心逻辑(GroupingMapSink)
public void accept(T t) {
K key = Objects.requireNonNull(classifier.apply(t), "element cannot be mapped to a null key");
// 取/建下游容器 → 累积
((A) map.computeIfAbsent(key, k -> downstream.supplier().get()))
.accept(t); // 注意:这里实际是 downstream 的 accumulator
}java
// 用法示例
Map<String, List<User>> byCity =
users.stream().collect(Collectors.groupingBy(User::getCity));
// 分组 + 计数 / 求和(下游收集器)
Map<String, Long> cityCount =
users.stream().collect(Collectors.groupingBy(User::getCity, Collectors.counting()));
Map<String, Integer> cityAgeSum =
users.stream().collect(Collectors.groupingBy(User::getCity, Collectors.summingInt(User::getAge)));groupingBy 特点:
键分类 → 值聚合(下游收集器决定值的形态)
无并发需求用 HashMap;并发用 groupingByConcurrent
下游收集器可嵌套(groupingBy → mapping → toList)五、partitioningBy 分区收集
java
public static <T> Collector<T, ?, Map<Boolean, List<T>>> partitioningBy(Predicate<? super T> predicate) {
return partitioningBy(predicate, toList());
}
public static <T, D, A> Collector<T, ?, Map<Boolean, D>> partitioningBy(
Predicate<? super T> predicate, Collector<? super T, A, D> downstream) {
// 内部:Partition 类持有两个容器
return new CollectorImpl<>(...);
}
// 内部实现:双容器
static final class Partition<T> extends AbstractMap<Boolean, T> implements Map<Boolean, T> {
final T forTrue; // 满足谓词的容器
final T forFalse; // 不满足的容器
Partition(T forTrue, T forFalse) { ... }
public T get(Object key) { return (Boolean) key ? forTrue : forFalse; }
}partitioningBy 特点:
键固定只有 true/false 两个
内部 Partition 持有两个下游容器(forTrue/forFalse)
predicate.test(t) → true 进 forTrue,false 进 forFalse
结果 Map<Boolean, List<T>>(或指定下游收集器的结果)
对比 groupingBy:分区是二分的特例分组,键只有两个java
// 用法示例
Map<Boolean, List<Integer>> parts =
numbers.stream().collect(Collectors.partitioningBy(n -> n % 2 == 0));
// true → 偶数集合,false → 奇数集合六、joining 字符串拼接
java
public static Collector<CharSequence, ?, String> joining(CharSequence delimiter) {
return joining(delimiter, "", ""); // 前缀/后缀默认空
}
public static Collector<CharSequence, ?, String> joining(CharSequence delimiter,
CharSequence prefix,
CharSequence suffix) {
return new CollectorImpl<>(
() -> new StringJoiner(delimiter, prefix, suffix), // 容器:StringJoiner
StringJoiner::add, // 累积:追加元素
StringJoiner::merge, // 合并:拼接两个 joiner
StringJoiner::toString, // 结果:toString
CH_NOID // 非恒等(有转换)
);
}joining 实现:
容器 = StringJoiner(内部 StringBuilder + delimiter/prefix/suffix)
begin(size) → 预分配容量(SIZED 时避免扩容)
accept → joiner.add(value)(追加 + 分隔符)
combiner → merge(两个 StringJoiner 拼接)
finisher → toStringjava
// 用法示例
String joined = names.stream().collect(Collectors.joining(", ", "[", "]"));
// → "[Alice, Bob, Charlie]"七、toMap 与重复键处理
java
public static <T, K, U> Collector<T, ?, Map<K, U>> toMap(
Function<? super T, ? extends K> keyMapper,
Function<? super T, ? extends U> valueMapper) {
return toMap(keyMapper, valueMapper, throwMerger(), HashMap::new);
}
// 重复键默认抛异常
private static <T> BinaryOperator<T> throwMerger() {
return (u, v) -> { throw new IllegalStateException(String.format("Duplicate key %s", u)); };
}toMap 流程:
keyMapper.apply(t) → 键
valueMapper.apply(t) → 值
map.put(key, value):
键已存在 → 调 mergeFunction(默认 throwMerger 抛异常)
指定 mergeFunction → 合并(如取后值/求和/拼接)java
// 用法示例
Map<Integer, String> idToName =
users.stream().collect(Collectors.toMap(User::getId, User::getName));
// 重复键 → IllegalStateException
// 冲突合并:重复键时取后一个值
Map<String, Integer> cityTotalAge = users.stream().collect(
Collectors.toMap(User::getCity, User::getAge, (a, b) -> a + b));八、summarizingInt 统计
java
public static <T> Collector<T, ?, IntSummaryStatistics> summarizingInt(
ToIntFunction<? super T> mapper) {
return new CollectorImpl<>(
IntSummaryStatistics::new, // 容器:统计对象
(r, t) -> r.accept(mapper.applyAsInt(t)), // 累积:喂入一个 int
(l, r) -> { l.combine(r); return l; }, // 合并:合并统计
CH_ID // 恒等 finish
);
}java
// java.util.IntSummaryStatistics
public class IntSummaryStatistics implements IntConsumer {
private long count; // 计数
private long sum; // 总和
private int min = Integer.MAX_VALUE; // 最小值
private int max = Integer.MIN_VALUE; // 最大值
public void accept(int value) {
++count; sum += value; min = Math.min(min, value); max = Math.max(max, value);
}
public void combine(IntSummaryStatistics other) {
count += other.count; sum += other.sum;
min = Math.min(min, other.min); max = Math.max(max, other.max);
}
public double getAverage() { return count > 0 ? (double) sum / count : 0.0; }
}summarizingInt 一次遍历同时统计:
count/sum/min/max/average 全部一次算出(无需多次 collect)
并行时各段统计对象 combine 合并
对应 Long/Double 版本:summarizingLong/summarizingDoublejava
// 用法示例
IntSummaryStatistics stats =
ages.stream().collect(Collectors.summarizingInt(Integer::intValue));
System.out.println(stats.getCount() + " " + stats.getSum() + " "
+ stats.getMin() + " " + stats.getMax() + " " + stats.getAverage());九、实现要点
Collector 核心:
四函数:supplier 建容器 / accumulator 累积 / combiner 合并 / finisher 转换
三特性:CONCURRENT 共享容器 / UNORDERED 无序 / IDENTITY_FINISH 恒等
toList:ArrayList::new + List::add + addAll
groupingBy:classifier 分类 + 下游收集器聚合(默认 toList/HashMap)
partitioningBy:Partition 双容器(true/false 各一)
joining:StringJoiner(StringBuilder + 分隔符)
toMap:keyMapper/valueMapper + 默认重复键抛异常(throwMerger)
summarizingInt:IntSummaryStatistics 一次遍历全量统计
常见陷阱:
重复键 toMap 不指定 merge → 运行时抛异常
groupingBy 的 null 键 → NullPointerException(map 可空 key 除外)
并行 collect 无 CONCURRENT → 各自容器 combiner 合并开销
toList 返回值是 ArrayList 的实际类型(接口类型不可变)
统计值 long/int 溢出 → 大数据量注意精度