☰
Apache Beam Java Top 变换详解:全局与按 Key 求 Top N 聚合
2026/10/12 1:57:29 网站建设 项目流程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

Top 是 Apache Beam Java SDK 中用于求最大值/最小值集合的核心聚合变换:它既可以从一个 PCollection 中取出最大(或最小)的 N 个元素,也可以从KV<K, V>集合中取出每个 Key 对应的最大(或最小)的 N 个 Value。本文围绕官方文档 top.md 展开,结合仓库中的源码实现、测试用例与可运行示例,完整覆盖 Top 的六大 API、自定义 Comparator、窗口与空输入边界行为以及底层有界堆的聚合原理,帮助读者直接上手并在真实 Pipeline 中正确使用。

Top 变换能做什么

根据官方文档的定位,Top 提供了一组变换,用于在集合中查找最大(或最小)的一组元素,或者在KV键值对集合中查找每个 Key 对应的最大(或最小)的一组值。它与 Sample 这类"采样"变换不同:Sample 的结果是任意抽样,而 Top 的结果是严格按序排列的极值集合,例如"成绩最高的 10 名学生""每个班级成绩最高的 3 名学生""最长的 5 个会话"。

在实现层面,Top 并非一个独立的 Runner 原语,而是建立在 Beam 的 Combine 机制之上。源码 Top.java 中,全局求 Top 的of方法内部就是:

return Combine.globally(new TopCombineFn<>(count, compareFn));

而按 Key 求 Top 的perKey方法则是:

return Combine.perKey(new TopCombineFn<>(count, compareFn));

这意味着 Top 具备 Combine 的全部特性:可以在分布式执行时由 Runner 自动做局部聚合与合并(tree-combine),从而大幅减少需要传输到单机的数据量。

六大核心 API 与类型签名

Top 是一个工具类(构造器私有,不允许实例化),对外暴露的全部是静态工厂方法。按"全局 vs 按 Key"和"自然序 vs 自定义比较器"两个维度,可以划分为 6 个主要入口:

API输入输出排序方式
Top.of(count, compareFn)PCollection<T>PCollection<List<T>>(单元素)自定义Comparator<T>,递减序
Top.largest(count)PCollection<T>(T 需Comparable)PCollection<List<T>>(单元素)自然序,递减序
Top.smallest(count)PCollection<T>(T 需Comparable)PCollection<List<T>>(单元素)自然序,递增序
Top.perKey(count, compareFn)PCollection<KV<K, V>>PCollection<KV<K, List<V>>>自定义Comparator<V>,递减序
Top.largestPerKey(count)PCollection<KV<K, V>>(V 需Comparable)PCollection<KV<K, List<V>>>自然序,递减序
Top.smallestPerKey(count)PCollection<KV<K, V>>(V 需Comparable)PCollection<KV<K, List<V>>>自然序,递增序

关键细节如下:

  • 全局求 Top的结果是一个只包含一个元素的PCollection<List<T>>,即所有极值打包进一个List返回。源码注释明确说明:结果List中全部元素必须能放进单台机器的内存(见 Top.java)。
  • 按 Key 求 Top的结果是PCollection<KV<K, List<V>>>,每个 Key 对应一个排序后的List<V>。此时约束放宽为"单个 Key 关联的全部 Value 必须能放进单台机器的内存",而 Key 的总数可以远超单机容量(见 Top.java)。
  • 排序方向:largest系列按递减序(从大到小)输出,smallest系列按递增序(从小到大)输出。自然序依赖元素类型实现Comparable。
  • count 约束:count必须大于等于 0,否则构造TopCombineFn时会抛出IllegalArgumentException("count must be >= 0 (not %s)"),该校验在 Top.java 中通过checkArgument完成,并有对应测试testCountConstraint(见 TopTest.java)。

除了上述入口,Top 还提供了面向 CombineFn 的细粒度工厂,便于在Combine.perKey、Combine.globally等场景直接复用聚合逻辑:largestFn、smallestFn,以及针对基础类型的largestLongsFn、largestIntsFn、largestDoublesFn、smallestLongsFn、smallestIntsFn、smallestDoublesFn(见 Top.java)。

第一个可运行示例:全局 Top 3

仓库中的官方示例 TopExample.java 正是文档页内嵌 Playground 代码(SDK_JAVA_Top)的来源,主逻辑非常简单:

// Create numbers PCollection<Integer> input = pipeline.apply(Create.of(1, 2, 3, 4, 5, 6)); PCollection<List<Integer>> result = input.apply(Top.largest(3));

这段代码创建1, 2, 3, 4, 5, 6六个整数,Top.largest(3)取出其中最大的 3 个,按递减序输出,即[6, 5, 4]。示例通过ParDo+LogOutput把结果打印到日志(见 TopExample.java)。

完整可编译的 Pipeline 如下(可直接替换示例中的 main 方法结构运行):

PipelineOptions options = PipelineOptionsFactory.create(); Pipeline pipeline = Pipeline.create(options); // 输入一个整数集合 PCollection<Integer> input = pipeline.apply(Create.of(1, 2, 3, 4, 5, 6)); // 取出最大的 3 个元素,输出 [6, 5, 4] PCollection<List<Integer>> largest = input.apply(Top.largest(3)); // 取出最小的 3 个元素,输出 [1, 2, 3] PCollection<List<Integer>> smallest = input.apply(Top.smallest(3)); pipeline.run();

注意Top.largest(3)返回的是PCollection<List<Integer>>,即整个结果是一个元素(一个 List),这与GroupByKey、ParDo的逐元素处理语义有明显区别。

按 Key 求每个 Key 的 Top N

当输入是PCollection<KV<K, V>>时,使用perKey系列即可对每个 Key 独立求 Top。以"每个班级成绩最高的 2 名学生"为例,仓库测试 TopTest.java 中构造了如下输入:

KV.of("a", 1), KV.of("a", 2), KV.of("a", 3), KV.of("b", 1), KV.of("b", 10), KV.of("b", 10), KV.of("b", 100)

应用Top.largestPerKey(2)后,测试断言的结果为:

  • Key"a"→[3, 2]
  • Key"b"→[100, 10]

应用Top.smallestPerKey(2)后,结果为:

  • Key"a"→[1, 2]
  • Key"b"→[1, 10]

(断言见 TopTest.java。)

代码写法如下:

PCollection<KV<String, Integer>> keyedValues = ...; // 例如来自 ParDo + KV.of(...) PCollection<KV<String, List<Integer>>> largestPerKey = keyedValues.apply(Top.largestPerKey(2)); PCollection<KV<String, List<Integer>>> smallestPerKey = keyedValues.apply(Top.smallestPerKey(2));

值得注意:largestPerKey允许重复值进入结果——测试中 Key"b"存在两个值为10的元素,结果[100, 10]只保留一个10,说明 Top 的语义是"取 N 个极值元素",重复元素各自独立参与排序与计数。

自定义 Comparator:Top.of 与 Top.perKey

当元素类型本身不实现Comparable,或者需要按业务规则排序(例如按字符串长度、按对象的某个字段)时,使用Top.of(count, compareFn)和Top.perKey(count, compareFn)传入自定义比较器。

一个直接来自测试的示例是按字符串长度排序(见 TopTest.java):

private static class OrderByLength implements Comparator<String>, Serializable { @Override public int compare(String a, String b) { if (a.length() != b.length()) { return a.length() - b.length(); } else { return a.compareTo(b); } } } PCollection<List<String>> longest = input.apply(Top.of(1, new OrderByLength()));

对测试输入["a", "bb", "c", "c", "z"],Top.of(1, new OrderByLength())取长度为 1 时最长的元素,即"bb"(断言见 TopTest.java)。

使用自定义 Comparator 必须同时满足两个硬性要求:

  1. 实现java.util.Comparator<T>,决定"谁更大"的语义;
  2. 实现java.io.Serializable,因为 Beam 的变换与 CombineFn 需要被序列化后分发到远端执行。源码中of与perKey的类型签名ComparatorT extends Comparator<T> & Serializable在编译期就强制了这一约束(见 Top.java),测试testPerKeySerializabilityRequirement也验证了带自定义比较器的perKey可以正常编译执行(见 TopTest.java)。

另外两个历史遗留的内部比较器类值得了解:Top.Largest(自然序)与Top.Smallest(逆自然序)已在源码中标注@Deprecated,官方推荐改用Top.Natural与Top.Reversed(见 Top.java)。largest/smallest系列内部正是分别使用Natural与Reversed作为默认比较器:

public static <T extends Comparable<T>> Combine.Globally<T, List<T>> largest(int count) { return Combine.globally(largestFn(count)); } public static <T extends Comparable<T>> TopCombineFn<T, Natural<T>> largestFn(int count) { return new TopCombineFn<T, Natural<T>>(count, new Natural<T>()) {}; } public static <T extends Comparable<T>> TopCombineFn<T, Reversed<T>> smallestFn(int count) { return new TopCombineFn<T, Reversed<T>>(count, new Reversed<>()) {}; }

(见 Top.java。)

底层实现原理:Combine + 有界堆

Top 的"省内存、可合并"特性来自其累加器BoundedHeap(有界堆),它是 Combine 的AccumulatingCombineFn.Accumulator实现(见 Top.java)。核心数据结构是一个由PriorityQueue支撑的小顶堆/大顶堆,容量被严格限制为count:

  • 添加元素(maybeAddInput):堆未满时直接入堆;堆已满时,只有新元素与堆顶比较"更大"才替换堆顶(poll后add),否则丢弃(见 Top.java)。这样任意时刻堆中只保留当前已见元素里的 Top N,内存占用为 O(N)。
  • 合并累加器(mergeAccumulator):分布式执行时,每个并行分片先各自维护一个局部有界堆,再把局部堆的元素逐个并入全局堆;一旦某个元素进不了 Top N,则提前终止循环——因为后续元素只会更小(见 Top.java)。这是 Top 能在海量数据上高效运行的关键。
  • 输出(extractOutput/asList):把堆中元素依次poll出来得到一个"最小在前"的列表,再反转成"最大在前"的列表,从而保证输出按要求的递减/递增序排列(见 Top.java)。

TopCombineFn继承了AccumulatingCombineFn,并提供:

  • 累加器 Coder:BoundedHeapCoder基于输入元素的ListCoder对堆内容编码(见 Top.java),默认输出PCollection的 Coder 也是元素 Coder 的ListCoder;
  • DisplayData:向监控面板暴露count(标签 "Top Count")与comparer(记录比较器类名)两个指标(见 Top.java),测试testDisplayData验证了这一点(见 TopTest.java);
  • 名称覆盖:变换名形如Combine.globally(Top(Natural))、Combine.perKey(Top(IntegerComparator)),测试testTopGetNames对其逐一断言(见 TopTest.java)。

边界行为:空集合、count=0 与窗口限制

Top 的边界行为已有明确的源码约定与测试覆盖,实际使用前务必掌握:

1. 空输入与全局窗口(GlobalWindows)如果输入 PCollection 使用全局窗口且为空,Top 会输出一个GlobalWindow中的空List<T>,而不是不输出(见 Top.java)。测试testTopEmpty验证了空集合下各 API 均产生空结果(见 TopTest.java)。

2. 非全局窗口(如固定窗口)若输入使用非GlobalWindows的窗口化策略,Top 默认不支持"空窗口输出默认值",必须在变换上追加.withoutDefaults()(输入为空时输出空 PCollection)或.asSingletonView()(输入为空时输出包含空 List 的单例视图),否则会在构建时抛出IllegalStateException(见 Top.java)。测试testTopEmptyWithIncompatibleWindows使用 10 天固定窗口的Create.empty输入验证了该异常(见 TopTest.java)。真实项目 TopWikipediaSessions.java 中对会话窗口化数据求 Top 时正是使用Top.of(1, comparator).withoutDefaults()的写法。

3. count = 0Top.of(0, ...)、Top.largest(0)、Top.largestPerKey(0)等均合法:全局求 Top 结果为空的单元素 List;按 Key 求 Top 时每个 Key 仍会输出一个空 List(测试testTopZero断言KV.of("a", Arrays.asList())存在,见 TopTest.java)。只有负的 count 会触发异常。

4. count 大于元素总数如果 count 超过输入元素个数,结果会包含输入的全部元素(仍然有序),例如对[1, 2, 3]求Top.largest(10)会得到[3, 2, 1](见 Top.java)。

与相关变换的对比

文档末节将 Sample 列为 Top 的相关变换。二者的区别在于:Sample 返回集合的任意抽样(用于数据探查、降采样),而Top 返回排序后的极值集合(用于排行榜、异常检测、Top-N 报表)。实际选型建议:

  • 需要"最大的 N 个"且结果有序 → 用 Top;
  • 只需要"任意抽 N 个样本"、不关心具体是哪些 → 用 Sample;
  • 既需要 Top N 又需要后续按窗口/会话聚合 → 注意窗口限制,配合.withoutDefaults()或.asSingletonView()使用。

小结

Top 是 Beam 聚合家族中使用频率极高的变换,本文覆盖了它的完整使用面:六大 API 的语义与签名、自定义 Serializable Comparator、按 Key 聚合、BoundedHeap有界堆的底层原理,以及空输入、count=0、非全局窗口等边界行为。理解这些细节后,读者可以放心在排行榜、Top-N 统计、异常值筛选等场景中使用Top.of/Top.largest/Top.perKey系列,并在需要时通过Combine的局部聚合能力处理超大规模数据。

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载
上一篇:QQ空间导出助手:三步永久备份你的青春记忆,告别数据丢失焦虑
下一篇:TaskbarX:让Windows任务栏图标自动居中的优雅解决方案

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询