Java 基础体系 · 第 61/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java Collector 深入:归约、分组、下游收集器、并行和自定义

Collector 是 Java Stream 中描述“如何把一串元素汇总成结果”的抽象。它不仅用于把流转换成 ListSetMap,还用于表达求和、统计、分组、嵌套分组、连接字符串,以及自定义的可并行归约算法。

要真正理解 Collector,需要同时掌握四个层次:

  1. 归约:如何把多个元素逐步合并成一个结果。
  2. 收集器协议supplieraccumulatorcombinerfinisher 如何协作。
  3. 下游收集器:分组后如何继续对每个组进行映射、过滤、统计或再次分组。
  4. 并行语义:拆分后的局部结果如何合并,以及什么条件下自定义收集器才是正确的。

一、从归约理解 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 reducecollect 的区别

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

AR 不必相同,这正是 finisher 存在的原因。


2.1 supplier:创建空容器

supplier 负责创建新的累加容器:

Supplier<StringBuilder> supplier = StringBuilder::new;

在串行流中,通常只需要一个容器。在并行流中,Stream 实现可能创建多个容器,分别处理不同的数据分片。

因此,supplier 返回的容器必须代表“尚未累积任何元素”的合法初始状态。


2.2 accumulator:把一个元素加入容器

accumulator 接收两个参数:

  1. 当前累加容器
  2. 流中的一个元素

例如:

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 加入容器 a
  • C(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 countingsummingaveraging

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():元素数量
  • summingIntsummingLongsummingDouble:数值求和
  • averagingIntaveragingLongaveragingDouble:平均值

空流的平均值返回 0.0,因为这些收集器以计数和总和表示状态,空状态的计数为零,最终按照 API 定义返回零值。

如果需要区分“没有元素”和“平均值为零”,应使用其他状态表示或先判断流是否为空。


4.2 minBymaxByreducing

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 会:

  1. 调用 classifier 得到键;
  2. 查找该键对应的组;
  3. 如果组不存在,创建一个新的下游累加容器;
  4. 把当前元素交给该组的下游收集器;
  5. 最终把所有键和组结果放入 Map

groupingBy 的默认下游收集器是 toList(),因此:

groupingBy(classifier)

等价于“按键分组,每组收集成列表”。


5.1 分组键的限制

分类函数不应返回 nullgroupingBy 要求分类结果作为合法的映射键使用,标准实现通常会对 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]
}

IntSummaryStatisticstoString() 展示形式属于类的输出实现细节,业务代码不应依赖它的具体文本格式;实际使用时应调用 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 的语义不同:分区结果通常包含 truefalse 两个分区,即使其中一个分区为空。

多级分组也可以继续结合映射和统计:

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


十一、groupingBygroupingByConcurrent

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,并声明并发收集相关特征,适合允许无序处理的并行分组场景。

但需要注意几个边界:

  1. 返回类型是 ConcurrentMap 更准确地反映其语义;如果只声明为 Map,仍可以使用,但会隐藏并发特性。
  2. 结果的键遍历顺序不应被依赖。
  3. 并发 Map 不等于组内列表是任意时刻都适合外部并发修改。
  4. 下游收集器的状态仍必须能在该收集模型下被正确维护。

如果业务要求严格的遇到顺序,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

此外,accumulatorcombinerfinisher 应尽量避免修改外部状态。外部副作用会让结果依赖执行模式,导致串行时正常、并行时失败。


十三、自定义 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());

结果是空字符串,因为:

  1. supplier 创建空 StringBuilder
  2. 没有元素,因此不调用 accumulator
  3. 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

标准收集器已经覆盖了绝大多数常见需求:

  • toList
  • toSet
  • toMap
  • groupingBy
  • partitioningBy
  • mapping
  • filtering
  • flatMapping
  • counting
  • summing
  • summarizing
  • joining
  • teeing

自定义 Collector 适合以下情况:

  1. 累积状态需要多个字段;
  2. 中间状态和最终结果类型不同;
  3. 需要定义特殊的分片合并逻辑;
  4. 现有收集器组合无法清晰表达算法;
  5. 需要让算法天然适配串行和并行收集。

例如,一个“统计总数、非空数量和字符总长度”的收集器,可以使用一个专门的可变状态类:

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

如果串行正确、并行错误,应优先对 supplieraccumulatorcombiner 和特征声明进行最小可重复验证。


十七、诊断自定义 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 的选择首先应由结果语义决定:

  • 要一组元素,使用 toListtoSettoCollection
  • 要按键得到多值,使用 groupingBy
  • 要按键得到单值,使用 toMap 并明确重复键策略;
  • 要组内映射,使用 mapping
  • 要组内过滤,使用 filtering
  • 要展开嵌套集合,使用 flatMapping
  • 要多个统计结果,使用 summarizingteeing
  • 要并行共享累积,只有在确实满足并发协议时才考虑 CONCURRENT

其次要审查内存行为。groupingBy(..., toList()) 会保存所有元素;counting()summingInt()summarizingInt() 则只保存汇总状态。如果最终只需要总数,就不应先收集完整列表再计算数量。

最后要审查顺序语义。顺序不是性能选项,而是结果定义的一部分。字符串拼接、列表收集、按时间先后构造结果时,必须保持遇到顺序;只有业务明确不依赖顺序时,才可以采用无序收集语义。

Collector 的本质是一个可组合的归约协议:元素被累积到局部状态,局部状态可以被正确合并,最终状态再转换为结果。理解这套协议后,groupingBy、下游收集器、并行 Stream 和自定义 Collector 就不再是互相孤立的 API,而是同一个归约模型在不同结果结构上的具体表达。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。