资讯详情

资讯详情

Apache Beam Java Top 变换详解:全局与按 Key 求 Top N 聚合

【免费下载链接】beamApache 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 个元素也可以从KVK, 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)PCollectionTPCollectionListT单元素自定义ComparatorT递减序Top.largest(count)PCollectionTT 需ComparablePCollectionListT单元素自然序递减序Top.smallest(count)PCollectionTT 需ComparablePCollectionListT单元素自然序递增序Top.perKey(count, compareFn)PCollectionKVK, VPCollectionKVK, ListV自定义ComparatorV递减序Top.largestPerKey(count)PCollectionKVK, VV 需ComparablePCollectionKVK, ListV自然序递减序Top.smallestPerKey(count)PCollectionKVK, VV 需ComparablePCollectionKVK, ListV自然序递增序关键细节如下全局求 Top的结果是一个只包含一个元素的PCollectionListT即所有极值打包进一个List返回。源码注释明确说明结果List中全部元素必须能放进单台机器的内存见 Top.java。按 Key 求 Top的结果是PCollectionKVK, ListV每个 Key 对应一个排序后的ListV。此时约束放宽为单个 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 PCollectionInteger input pipeline.apply(Create.of(1, 2, 3, 4, 5, 6)); PCollectionListInteger result input.apply(Top.largest(3));这段代码创建1, 2, 3, 4, 5, 6六个整数Top.largest(3)取出其中最大的 3 个按递减序输出即[6, 5, 4]。示例通过ParDoLogOutput把结果打印到日志见 TopExample.java。完整可编译的 Pipeline 如下可直接替换示例中的 main 方法结构运行PipelineOptions options PipelineOptionsFactory.create(); Pipeline pipeline Pipeline.create(options); // 输入一个整数集合 PCollectionInteger input pipeline.apply(Create.of(1, 2, 3, 4, 5, 6)); // 取出最大的 3 个元素输出 [6, 5, 4] PCollectionListInteger largest input.apply(Top.largest(3)); // 取出最小的 3 个元素输出 [1, 2, 3] PCollectionListInteger smallest input.apply(Top.smallest(3)); pipeline.run();注意Top.largest(3)返回的是PCollectionListInteger即整个结果是一个元素一个 List这与GroupByKey、ParDo的逐元素处理语义有明显区别。按 Key 求每个 Key 的 Top N当输入是PCollectionKVK, 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)后测试断言的结果为Keya→[3, 2]Keyb→[100, 10]应用Top.smallestPerKey(2)后结果为Keya→[1, 2]Keyb→[1, 10]断言见 TopTest.java。代码写法如下PCollectionKVString, Integer keyedValues ...; // 例如来自 ParDo KV.of(...) PCollectionKVString, ListInteger largestPerKey keyedValues.apply(Top.largestPerKey(2)); PCollectionKVString, ListInteger smallestPerKey keyedValues.apply(Top.smallestPerKey(2));值得注意largestPerKey允许重复值进入结果——测试中 Keyb存在两个值为10的元素结果[100, 10]只保留一个10说明 Top 的语义是取 N 个极值元素重复元素各自独立参与排序与计数。自定义 ComparatorTop.of 与 Top.perKey当元素类型本身不实现Comparable或者需要按业务规则排序例如按字符串长度、按对象的某个字段时使用Top.of(count, compareFn)和Top.perKey(count, compareFn)传入自定义比较器。一个直接来自测试的示例是按字符串长度排序见 TopTest.javaprivate static class OrderByLength implements ComparatorString, Serializable { Override public int compare(String a, String b) { if (a.length() ! b.length()) { return a.length() - b.length(); } else { return a.compareTo(b); } } } PCollectionListString longest input.apply(Top.of(1, new OrderByLength()));对测试输入[a, bb, c, c, z]Top.of(1, new OrderByLength())取长度为 1 时最长的元素即bb断言见 TopTest.java。使用自定义 Comparator 必须同时满足两个硬性要求实现java.util.ComparatorT决定谁更大的语义实现java.io.Serializable因为 Beam 的变换与 CombineFn 需要被序列化后分发到远端执行。源码中of与perKey的类型签名ComparatorT extends ComparatorT 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 ComparableT Combine.GloballyT, ListT largest(int count) { return Combine.globally(largestFn(count)); } public static T extends ComparableT TopCombineFnT, NaturalT largestFn(int count) { return new TopCombineFnT, NaturalT(count, new NaturalT()) {}; } public static T extends ComparableT TopCombineFnT, ReversedT smallestFn(int count) { return new TopCombineFnT, ReversedT(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并提供累加器 CoderBoundedHeapCoder基于输入元素的ListCoder对堆内容编码见 Top.java默认输出PCollection的 Coder 也是元素 Coder 的ListCoderDisplayData向监控面板暴露count标签 Top Count与comparer记录比较器类名两个指标见 Top.java测试testDisplayData验证了这一点见 TopTest.java名称覆盖变换名形如Combine.globally(Top(Natural))、Combine.perKey(Top(IntegerComparator))测试testTopGetNames对其逐一断言见 TopTest.java。边界行为空集合、count0 与窗口限制Top 的边界行为已有明确的源码约定与测试覆盖实际使用前务必掌握1. 空输入与全局窗口GlobalWindows如果输入 PCollection 使用全局窗口且为空Top 会输出一个GlobalWindow中的空ListT而不是不输出见 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有界堆的底层原理以及空输入、count0、非全局窗口等边界行为。理解这些细节后读者可以放心在排行榜、Top-N 统计、异常值筛选等场景中使用Top.of/Top.largest/Top.perKey系列并在需要时通过Combine的局部聚合能力处理超大规模数据。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Java Top 变换详解从全局 Top-N 到按 Key 聚合的完整实战指南Apache Beam Java Top 变换详解从全局 Top N 到按 Key 聚合的完整实战指南 导读 本文深入讲解 Apache Beam Java批处理流处理大数据Apache Beam Java SDK Top Transform 完全指南全局与按 Key 的最大/最小 Top-N 聚合Apache Beam Java SDK Top Transform 完全指南全局与按 Key 的最大/最小 Top N 聚合 Apache Beam 的 T大数据批处理流处理数据工程Apache Beam Java SDK 的 Sum 聚合变换全局求和与按 Key 求和实战指南Apache Beam Java SDK 的 Sum 聚合变换全局求和与按 Key 求和实战指南 本篇技术指南聚焦 Apache Beam Java SDK大数据批处理流处理数据工程上一篇QQ空间导出助手三步永久备份你的青春记忆告别数据丢失焦虑下一篇TaskbarX让Windows任务栏图标自动居中的优雅解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →