CodeGym /课程 /JAVA 25 SELF /自定义 Collector 和 Spliterator

自定义 Collector 和 Spliterator

JAVA 25 SELF
第 33 级 , 课程 4
可用

1. 自定义 Collector:何时以及如何编写

在 Java Stream API 中,使用接口 Collector 将流转换为集合或聚合结果。通常你会使用 Collectors 类中的现成收集器(toList()toMap()groupingBy() 等),但有时需求特殊——这时你可以编写自己的 Collector。

Collector 是一种对象,用来描述如何把流中的元素汇聚成最终结果。它定义了四个(实际上是五个)关键组件:

  • supplier —— 创建一个用于收集元素的新容器(例如新的列表或映射)。
  • accumulator —— 将下一个元素追加到容器中。
  • combiner —— 合并两个容器(对并行流尤为重要!)。
  • finisher —— 将容器转换为最终结果(例如将其设为不可变,或转换为其他类型)。
  • characteristics —— 一组标志,描述收集器的属性(例如是否支持并行、是否改变结果类型等)。

签名:

Collector<T, A, R>
  • T —— 流中元素的类型,
  • A —— 中间累加器的类型,
  • R —— 结果类型。

2. 示例:用于 MultiMap (Map<K, List<V>>) 的 Collector

假设你希望把 Pair<K, V> 的流收集到 Map<K, List<V>>(多重映射)中,使每个键对应一个值列表。

实现示例:

public static <K, V> Collector<Pair<K, V>, ?, Map<K, List<V>>> toMultiMap() {
    return Collector.of(
        HashMap::new, // supplier
        (map, pair) -> map.computeIfAbsent(pair.key(), k -> new ArrayList<>()).add(pair.value()), // accumulator
        (map1, map2) -> { // combiner
            map2.forEach((k, vList) -> map1.merge(k, vList, (l1, l2) -> { l1.addAll(l2); return l1; }));
            return map1;
        },
        Function.identity(), // finisher
        Collector.Characteristics.UNORDERED
    );
}

用法:

List<Pair<String, Integer>> pairs = List.of(
    new Pair<>("a", 1), new Pair<>("b", 2), new Pair<>("a", 3)
);

Map<String, List<Integer>> multiMap = pairs.stream().collect(toMultiMap());
// multiMap: {a=[1, 3], b=[2]}

3. 示例:用于 Top-N 元素的 Collector

假设你希望把流收集为包含 N 个最大元素的列表(例如按降序的 Top-5)。

实现:

public static <T> Collector<T, ?, List<T>> topN(int n, Comparator<? super T> comparator) {
    return Collector.of(
        () -> new PriorityQueue<>(n, comparator), // supplier
        (pq, t) -> {
            pq.offer(t);
            if (pq.size() > n) pq.poll(); // 移除最小值
        },
        (pq1, pq2) -> {
            pq2.forEach(t -> {
                pq1.offer(t);
                if (pq1.size() > n) pq1.poll();
            });
            return pq1;
        },
        pq -> {
            List<T> result = new ArrayList<>(pq);
            result.sort(comparator.reversed()); // 降序
            return result;
        },
        Collector.Characteristics.UNORDERED
    );
}

用法:

List<Integer> top3 = Stream.of(5, 1, 9, 3, 7, 2).collect(topN(3, Comparator.naturalOrder()));
// top3: [9, 7, 5]

4. 何时不该编写自己的 Collector

  • 如果可以通过组合标准收集器与下游操作(如 groupingBymappingflatMappingcollectingAndThen 等)来表达需求,最好使用它们
  • 只有在确实非常非标准的场景下才需要自定义 Collector(特殊数据结构、复杂聚合、Top‑N、多重映射等)。
  • 不要为写而写 Collector——这会增加维护与测试的复杂度。

示例:

// 不必为 Map<K, Set<V>> 编写自定义 Collector:
.collect(Collectors.groupingBy(
    Pair::key,
    Collectors.mapping(Pair::value, Collectors.toSet())
))

5. 自定义 Spliterator:为何以及如何实现

Spliterator 是一个用于高效遍历与拆分集合(或其他数据源)的特殊接口,尤其适合并行处理。与普通迭代器不同,Spliterator 可以“拆分”(split)集合为彼此独立的部分以进行并行处理。

关键方法:

  • tryAdvance(Consumer<? super T> action) —— 处理下一个元素。
  • trySplit() —— 尝试将集合划分为两部分(返回其中一部分的新 Spliterator)。
  • estimateSize() —— 估计剩余元素数量。
  • characteristics() —— 特性位掩码(ORDEREDSIZEDSUBSIZED 等)。

trySplit:划分策略

均衡划分 对并行流非常重要:trySplit 应返回规模大致相当的部分,以实现负载均衡。

如果无法再划分(例如元素过少)——返回 null

示例:分批读取文件的 Spliterator
假设你有一个大文件,希望每次处理 1000 行(分批),以避免将全部内容放入内存。

public class ChunkedLineSpliterator implements Spliterator<List<String>> {
    private final BufferedReader reader;
    private final int chunkSize;

    public ChunkedLineSpliterator(BufferedReader reader, int chunkSize) {
        this.reader = reader;
        this.chunkSize = chunkSize;
    }

    @Override
    public boolean tryAdvance(Consumer<? super List<String>> action) {
        List<String> chunk = new ArrayList<>(chunkSize);
        try {
            String line;
            for (int i = 0; i < chunkSize && (line = reader.readLine()) != null; i++) {
                chunk.add(line);
            }
            if (chunk.isEmpty()) return false;
            action.accept(chunk);
            return true;
        } catch (IOException e) {
            throw new UncheckedIOException(e);
        }
    }

    @Override
    public Spliterator<List<String>> trySplit() {
        // 对于文件的流式读取,划分没有意义 —— 返回 null
        return null;
    }

    @Override
    public long estimateSize() {
        return Long.MAX_VALUE; // 事先未知
    }

    @Override
    public int characteristics() {
        return ORDERED | NONNULL;
    }
}

用法:

try (BufferedReader reader = Files.newBufferedReader(Path.of("big.txt"))) {
    StreamSupport.stream(new ChunkedLineSpliterator(reader, 1000), false)
        .forEach(chunk -> processChunk(chunk));
}

Spliterator 的特性

  • ORDERED —— 元素具有确定的顺序(例如列表)。
  • SIZED —— 已知精确的元素数量。
  • SUBSIZED —— 所有通过 trySplit 得到的 Spliterator 也都是 SIZED 的。
  • IMMUTABLE —— 遍历期间数据源不会发生变化。
  • CONCURRENT —— 数据源支持安全的并行修改。
  • DISTINCTSORTEDNONNULL —— 其他附加属性。

重要:正确声明特性会影响流的优化。

6. 示例

  • 按块(chunk)读取文件 —— 允许分段处理大文件,而无需将其全部载入内存。
  • 无额外分配的解析 —— 如果你在解析字节/字符流并希望最小化临时对象创建,可以实现一个 Spliterator 来返回源数组的“窗口”或“切片”。

示例:按行解析 CSV 的 Spliterator

public class CsvLineSpliterator implements Spliterator<String[]> {
    private final BufferedReader reader;

    public CsvLineSpliterator(BufferedReader reader) {
        this.reader = reader;
    }

    @Override
    public boolean tryAdvance(Consumer<? super String[]> action) {
        try {
            String line = reader.readLine();
            if (line == null) return false;
            action.accept(line.split(","));
            return true;
        } catch (IOException e) {
            throw new UncheckedIOException(e);
        }
    }

    @Override
    public Spliterator<String[]> trySplit() {
        return null; // 顺序解析
    }

    @Override
    public long estimateSize() {
        return Long.MAX_VALUE;
    }

    @Override
    public int characteristics() {
        return ORDERED | NONNULL;
    }
}

7. 与 parallel() 集成:如何安全实现

  • 如果你的 Spliterator 支持并行划分(trySplit 返回的不是 null),并且特性包含 SIZED/SUBSIZED,那么 Stream API 能更高效地并行处理。
  • 对于流式数据源(文件、套接字),通常不支持划分——请使用顺序流。
  • 对于集合与数组——实现均衡划分(例如将数组对半拆分)。

示例:用于数组的 Spliterator

public class ArraySpliterator<T> implements Spliterator<T> {
    private final T[] array;
    private int start, end;

    public ArraySpliterator(T[] array, int start, int end) {
        this.array = array;
        this.start = start;
        this.end = end;
    }

    @Override
    public boolean tryAdvance(Consumer<? super T> action) {
        if (start < end) {
            action.accept(array[start++]);
            return true;
        }
        return false;
    }

    @Override
    public Spliterator<T> trySplit() {
        int mid = (start + end) >>> 1;
        if (mid == start) return null;
        ArraySpliterator<T> split = new ArraySpliterator<>(array, start, mid);
        start = mid;
        return split;
    }

    @Override
    public long estimateSize() {
        return end - start;
    }

    @Override
    public int characteristics() {
        return ORDERED | SIZED | SUBSIZED | IMMUTABLE;
    }
}

用法:

String[] arr = {"a", "b", "c", "d"};
StreamSupport.stream(new ArraySpliterator<>(arr, 0, arr.length), true)
    .forEach(System.out::println);
1
调查/小测验
优化集合操作第 33 级,课程 4
不可用
优化集合操作
优化集合操作
评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION