拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Beam Kotlin Katas 实战:用 Min 聚合变换求取 PCollection 最小值

Apache Beam Kotlin Katas 实战:用 Min 聚合变换求取 PCollection 最小值 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本篇技术指南围绕 Apache Beam 官方 Kotlin Katas 学习课程中 Common Transforms - Aggregation - Min 一节展开讲解如何使用Min聚合变换从PCollection中计算全局最小值。你将看到该 Katas 练习的完整任务说明、参考实现与单元测试并深入Min变换在 Beam Java SDK 核心源码 中的底层实现原理从而掌握globally与perKey两类聚合形态及空输入时的边界行为。一、Kata 任务是什么Aggregation 系列的最小值练习在 Apache Beam 的 Kotlin Katas 课程中Common Transforms 章节按功能划分为 Filter、Aggregation、WithKeys 三个小节见 section-info.yaml其中 Aggregation 小节依次包含 Count、Sum、Mean、Min、Max 五个练习见 lesson-info.yaml。本次任务对应的说明文档是 Min/task.md其核心要求只有一句话Kata:Compute the minimum of the elements from an input.即对输入 PCollection 中的全部元素求最小值。文档给出的提示是使用Min变换。这是一个典型的给定骨架、补齐核心逻辑式练习Katas 会预先搭建好Pipeline、输入数据与日志输出框架练习者只需在applyTransform函数中补上求最小值的PTransform再通过配套测试验证结果。二、参考实现拆解Task.kt 的完整代码与执行流程该练习的参考实现位于 Task.kt完整代码如下package org.apache.beam.learning.katas.commontransforms.aggregation.min import org.apache.beam.learning.katas.util.Log import org.apache.beam.sdk.Pipeline import org.apache.beam.sdk.options.PipelineOptionsFactory import org.apache.beam.sdk.transforms.Create import org.apache.beam.sdk.transforms.Min import org.apache.beam.sdk.values.PCollection object Task { JvmStatic fun main(args: ArrayString) { val options PipelineOptionsFactory.fromArgs(*args).create() val pipeline Pipeline.create(options) val numbers pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)) val output applyTransform(numbers) output.apply(Log.ofElements()) pipeline.run() } JvmStatic fun applyTransform(input: PCollectionInt): PCollectionInt { return input.apply(Min.integersGlobally()) } }2.1 骨架代码的作用PipelineOptionsFactory.fromArgs(*args).create()从命令行参数解析并创建PipelineOptions这是 Beam 管线的标准初始化方式pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10))使用Create变换将 1 到 10 的整数序列构造成内存中的PCollectionInt作为练习的输入数据output.apply(Log.ofElements())Katas 自带的日志工具把 PCollection 中每个元素通过 SLF4J 打印出来便于在终端观察结果。其实现位于 Log.kt内部是一个ParDo.of(DoFn)会将元素字符串输出为LOG.info(message)并且如果元素所在的窗口不是GlobalWindow还会附带打印窗口信息pipeline.run()提交并执行整个管线。2.2 核心一行Min.integersGlobally()练习者需要完成的正是applyTransform方法input.apply(Min.integersGlobally())。Min是 Beam SDK 提供的聚合变换工具类完整源码见 Min.javaintegersGlobally()返回一个Combine.GloballyInteger, Integer作用于整条 PCollection输出一个只包含单个元素即全局最小值的PCollectionInteger对本例输入1..10运行后输出集合中仅有一个元素1随后由Log.ofElements()打印。三、测试验证TaskTest.kt 与 PAssert 断言每个 Kata 都配套了基于TestPipeline与PAssert的单元测试本练习的测试位于 TaskTest.ktclass TaskTest { get:Rule Transient val testPipeline: TestPipeline TestPipeline.create() Test fun common_transforms_aggregation_min() { val values Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10) val numbers testPipeline.apply(values) val results applyTransform(numbers) PAssert.that(results).containsInAnyOrder(1) testPipeline.run().waitUntilFinish() } }测试要点TestPipelineBeam 的测试专用 Pipeline支持在测试环境中直接执行PAssert.that(results).containsInAnyOrder(1)断言聚合结果集合中恰好包含1这一元素。containsInAnyOrder不关心元素顺序这正是聚合类变换输出单个元素场景下的标准写法testPipeline.run().waitUntilFinish()同步等待管线执行完毕确保断言在结果产出后执行。通过该测试即可验证Min.integersGlobally()的实现是否正确——如果练习者误用Max或其它变换断言将失败。四、深入源码Min 变换的 API 全景与底层实现4.1 全局聚合Globally与按键聚合PerKey两类形态从 Min.java 的源码结构看Min提供两类变换xxxGlobally()对整条 PCollection 求最小值输出单元素 PCollectionxxxPerKey()对PCollectionKVK, V按 key 分组后分别求每个 key 对应 value 集合的最小值输出PCollectionKVK, V即每个 key 对应一行(key, 最小值)结果。源码中integersPerKey()的注释明确说明了其语义输入PCollectionKVK, Integer输出中包含将每个不同 key 映射到该 key 下所有 value 最小值的PCollectionKVK, Integer同时指出窗口与时间戳行为遵循Combine.PerKey的语义。4.2 针对数值类型的开箱即用方法源码为三种基本数值类型提供了直接可用的工厂方法方法输入类型输出类型空输入时的 identityMin.integersGlobally()/integersPerKey()Integer/KVK, IntegerInteger/KVK, IntegerInteger.MAX_VALUEMin.longsGlobally()/longsPerKey()Long/KVK, LongLong/KVK, LongLong.MAX_VALUEMin.doublesGlobally()/doublesPerKey()Double/KVK, DoubleDouble/KVK, DoubleDouble.POSITIVE_INFINITY以integersGlobally()为例其实现为Combine.globally(new MinIntegerFn())而MinIntegerFn内部的核心逻辑只有两个方法Override public int apply(int left, int right) { return Math.min(left, right); } Override public int identity() { return Integer.MAX_VALUE; }这揭示了一个容易被忽视的边界语义空输入时全局最小值是Integer.MAX_VALUELong 为Long.MAX_VALUEDouble 为Double.POSITIVE_INFINITY。因为MinIntegerFn的 identity 被定义为该类型的最大值空集合聚合后返回 identity。同理MinLongFn与MinDoubleFn分别使用Long.MAX_VALUE与Double.POSITIVE_INFINITY作为 identity这在处理可能为空的输入数据时值得特别留意。4.3 任意类型与自定义比较器除了内置数值类型Min还支持任意Comparable类型与自定义ComparatorMin.globally()/Min.perKey()要求元素类型T实现Comparable按自然顺序求最小值空输入返回nullMin.globally(comparator)/Min.perKey(comparator)允许传入自定义Comparator由调用方定义最小的判定规则Min.of(identity, comparator)/Min.of(comparator)返回可直接嵌入Combine.globally或Combine.perKey的BinaryCombineFnMin.naturalOrder()/Min.naturalOrder(identity)基于自然顺序的BinaryCombineFn可显式指定 identity 值。这些方法最终都落到私有的MinFnT类它继承BinaryCombineFnTapply(left, right)的实现是comparator.compare(left, right) 0 ? left : right即比较两元素并保留较小者当未显式提供 identity 时identity()返回null。同时MinFn通过populateDisplayData将比较器类注册进 DisplayData便于在监控与调试界面中查看所用比较器类型。4.4 与 Max 的对照Min的源码注释中给出了与Max对称的用法同目录下的 Max.java 与 Katas 练习 Max/task.md 对应。求最大值只需把Min.integersGlobally()换成Max.integersGlobally()且Max的空输入 identity 恰好相反Integer.MIN_VALUE。二者的 API 形态、PerKey 形态与窗口语义完全一致理解了Min也就掌握了这一族聚合变换的通用模式。五、运行方式与实战扩展5.1 运行与验证在 Katas 项目中运行该练习可通过 Gradle 执行Task的main方法或直接运行配套的TaskTest单元测试。正常执行后终端会通过Log.ofElements()打印出最小值1测试模式下则由PAssert.that(results).containsInAnyOrder(1)完成自动断言。5.2 从 Globally 到 PerKey 的实战迁移把全局最小值升级为按 key 分组求最小值是数据处理的常见需求。例如输入PCollectionKVString, Int如按城市分组的温度数据只需将变换替换为input.apply(Min.integersPerKey())即可得到每个 key 对应的最小值列表与Combine.PerKey的窗口、时间戳语义保持一致。该用法同样可直接参考 Min.java 中integersPerKey()的示例。5.3 聚合练习家族完成Min之后建议继续同一小节的其余练习以形成完整认知Count计数、Sum求和、Mean求均值、Max求最大值它们共同覆盖了 Beam 中最常用的聚合原语。若使用 Java 版本学习可对照 Java 版 Task.java其核心逻辑与 Kotlin 版完全一致。六、小结通过本 Kata你完成了一次读任务 → 补实现 → 跑测试的完整练习闭环并掌握了 Apache Beam 中Min聚合变换的核心知识用法Min.integersGlobally()对PCollectionInt求全局最小值输出单元素 PCollection验证使用TestPipelinePAssert.that(...).containsInAnyOrder(...)断言聚合结果原理Min底层基于Combine.globally/Combine.perKey与BinaryCombineFn数值类型 identity 分别为Integer.MAX_VALUE、Long.MAX_VALUE、Double.POSITIVE_INFINITY空输入时会返回该 identity扩展perKey形态可按 key 分组求最小值自定义Comparator与Min.of(...)/Min.naturalOrder()可支持任意元素类型。对 Kotlin 开发者而言Katas 是理解 Beam 编程模型的极佳入门路径而Min正是掌握聚合类变换的第一块基石。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam KatasKotlin实战使用 Max 变换计算聚合最大值Apache Beam KatasKotlin实战使用 Max 变换计算聚合最大值 导读 Apache Beam 提供了一组开箱即用的聚合变换用于对 P大数据批处理流处理数据工程Apache Beam Java 实战用 Min 聚合变换计算全局最小值Katas 入门篇Apache Beam Java 实战用 Min 聚合变换计算全局最小值Katas 入门篇 本文基于 Apache Beam 官方 Katas 课程中 大数据批处理流处理数据工程Apache Beam Go SDK Katas使用 stats.Min 聚合计算 PCollection 最小值Apache Beam Go SDK Katas使用 stats.Min 聚合计算 PCollection 最小值 本指南以 Apache Beam 仓库中大数据批处理流处理数据工程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表