ARTICLE DETAIL

资讯详情

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

Apache Beam Kotlin Katas 实战:Aggregation 之 Count 聚合变换详解

Apache Beam Kotlin Katas 实战:Aggregation 之 Count 聚合变换详解 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读本文以 Apache Beam 官方 Kotlin Katas 训练项目中的Aggregation - Count一课为骨架系统讲解 Beam 中最常用的聚合变换Count从全局计数Count.globally()的用法与输出类型到perElement()、perKey()两种按维度计数的变体再到其底层CombineFn的累加器实现原理与测试验证方式。读完本文你将能独立完成该 Kata并能在真实 Beam 管道中准确选择与使用 Count 系列变换完成各类计数统计任务。一、Kata 任务与课程定位1.1 本课在 Katas 课程中的位置在 Apache Beam 仓库的 learning/katas/kotlin 目录下Kotlin Katas 是一套用 IntelliJ Education / EduTools 插件交互式完成的 Beam 入门练习。其中 Common Transforms 下的 Aggregation 课程见 lesson-info.yaml依次编排了五个聚合变换练习Count、Sum、Mean、Min、MaxCount 是第一个、也是理解其余四个聚合变换的基石。1.2 本课的 Kata 目标task.md 给出了本课的核心任务Kata:Count the number of elements from an input.统计输入中元素的数量任务提示只有一个使用Count变换。也就是说本练习要求你接收一个包含 1 到 10 十个整数的PCollectionInt输出一个表示元素总数的PCollectionLong值为 10L。1.3 练习的运行方式本 Kata 采用填空式教学Task.kt中预留了TODO()占位符由练习者补全applyTransform函数体隐藏的TaskTest.kt作为评分器只有实现正确时测试才会通过。关于项目导入与运行环境IntelliJ Education 导入 Gradle 项目、配置 JDK、以 Course 视图浏览练习请参考 learning/katas/kotlin/README.md 中的 Setup 步骤。二、完整解法用 Count.globally() 实现全局计数2.1 答案源码本练习的标准实现位于 Task.kt其核心代码为package org.apache.beam.learning.katas.commontransforms.aggregation.count 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.Count import org.apache.beam.sdk.transforms.Create 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): PCollectionLong { return input.apply(Count.globally()) // ← 填空处TODO() } }2.2 逐步拆解管道创建管道PipelineOptionsFactory.fromArgs(*args).create()解析命令行参数生成PipelineOptions再Pipeline.create(options)创建管道实例构造输入Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)将内存中的十个整数包装为PCollectionInt这是无外部 I/O 的练习型数据源核心变换input.apply(Count.globally())对整条PCollection做全局聚合计数观察结果output.apply(Log.ofElements())将每个输出元素打印到日志。Log是 Katas 提供的工具类其实现见 Log.kt本质是一个ParDo包装的LoggingTransform在ProcessElement中通过 SLF4J 输出元素内容若窗口不是全局窗口还会额外附加窗口信息便于在流式/开窗场景下观察元素归属执行pipeline.run()提交管道本地运行默认使用 DirectRunner。运行后日志会打印10即输入的十个元素总数。2.3 为什么输出类型是 PCollection Count.globally()的返回类型是PCollectionLong。原因在底层实现Count.globally()本质是Combine.globally(new CountFnT())而CountFnT是CombineFnT, long[], Long——累加器用long[]承载计数最终extractOutput取出long值装箱为Long。这也解释了为何 Kotlin 侧函数签名写的是PCollectionLong而非PCollectionInt计数可能远超 Int 范围Beam 采用Long类型天然支持大数量级。三、Count 家族的三种形态与适用场景本练习只要求全局计数但Count变换实际提供三个入口方法见 Count.java 源码注释方法输入输出语义典型场景Count.globally()PCollectionTPCollectionLong整个集合的元素总数统计总记录条数、事件总量Count.perElement()PCollectionTPCollectionKVT, Long每个不同元素出现的次数词频统计、去重后各类别计数Count.perKey()PCollectionKVK, VPCollectionKVK, Long每个 Key 关联的 Value 数量按用户/地区/商品等维度计数三种形态分别对应聚合的三种粒度全局globally、按值perElement、按键perKey。在真实业务中perElement()可直接实现 WordCount 中的词频统计val wordCounts: PCollectionKVString, Long words.apply(Count.perElement())3.1 perElement() 的实现方式从源码看Count.java 中的PerElementT变换分两步完成先用MapElements把每个元素T映射为KVT, VoidKV.of(element, null)再交给Count.perKey()按元素本身作为 Key 计数。因此perElement()本质上是perKey()的特例。3.2 一个值得注意的约束perElement()判定元素相等的方式是把元素用输入PCollection的Coder编码后再比较字节Coder#verifyDeterministic()因此要求输入 Coder 必须是确定性的globally()与perKey()则无此约束。若输入使用了非全局窗口的窗口策略Count.globally()会直接抛出不兼容全局窗口的错误提示源码中的getIncompatibleGlobalWindowErrorMessage()明确建议改用Combine.globally(Count.TcombineFn()).withoutDefaults()来处理开窗场景下的计数。四、源码级原理CountFn 累加器如何工作Count是典型的Combine变换其精髓在于内部的CountFnT见 Count.java。它实现了CombineFn的四个核心方法createAccumulator()返回long[] {0}——刻意用长度为 1 的数组作为可变 long 的盒子规避 Java 中 long 不可变、无法原地累加的问题addInput(acc, input)对每个到达的元素执行accumulator[0] 1这是每个元素计 1的语义落点mergeAccumulators(accs)分布式环境下多个分区的部分计数在此合并——遍历所有累加器并累加各自的计数这正是 Beam 聚合能够水平扩展的关键extractOutput(acc)从最终累加器取出accumulator[0]并装箱为Long。此外getAccumulatorCoder用VarInt变长整数对累加器编码使中间结果在分布式传输时足够紧凑equals/hashCode基于类型实现保证相同变换可被合理合并优化。分布式含义由于Combine天然支持本地部分聚合 全局合并即使输入分布在数百台机器上Count 也只需在每个分区维护一个long计数器再逐级合并内存与网络开销都极小——这正是它被广泛用于流式事件计数等高频场景的原因。五、测试验证PAssert 断言输出隐藏的测试文件 TaskTest.kt 是练习的评分依据同时也是学习如何测试 Beam 聚合变换的范本class TaskTest { Transient get:Rule val testPipeline: TestPipeline TestPipeline.create() Test fun common_transforms_aggregation_count() { 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(10L) testPipeline.run().waitUntilFinish() } }测试的验证思路非常清晰用TestPipeline作为 JUnit Rule自动管理管道生命周期输入与Task.main完全一致1~10 十个整数保证练习场景一致关键断言PAssert.that(results).containsInAnyOrder(10L)校验输出集合中只含一个元素10L——注意10L是Long字面量与Count.globally()返回的PCollectionLong类型严格对应containsInAnyOrder不关心元素顺序适用于聚合结果这类单元素输出run().waitUntilFinish()确保管道执行完毕、断言生效。如果你的实现误用了Count.perElement()输出会是 10 个KVInt, Long或返回值类型写错PAssert 都会因输出与10L不匹配而失败——这正是填空式教学设计的精妙之处。六、举一反三向 Sum / Mean / Min / Max 迁移掌握 Count 后Aggregation 课程中其余四个变换见 lesson-info.yaml几乎可以零成本迁移SumSum.integersGlobally()等按数值类型区分的全局求和MeanMean.globally()计算全局均值Min / MaxMin.globally()/Max.globally()求全局最值。它们与Count.globally()一样都是Combine.globally(...)的便捷封装输出同样为单元素PCollection如PCollectionLong或数值类型测试断言模式也完全一致。理解了 Count 的CombineFn机制就理解了整个 Aggregation 课程背后的统一抽象。七、总结Count 是 Beam 聚合变换的入门第一课Count.globally()一行代码即可完成全局计数返回PCollectionLong按需扩展时perElement()做词频类统计、perKey()做维度分组计数其底层CountFn采用long[]可变累加器 分布式合并兼顾正确性与可扩展性相关实现可在 Count.java 中完整查阅配合 TaskTest.kt 的PAssert断言模式你可以把同样的测试方法复用到任何自定义聚合变换的验证中。完成本 Kata 后不妨直接打开 Aggregation 课程中下一个练习 Sum你会发现聚合的思想是统一的变的只是CombineFn内部的运算规则。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战使用 Count 聚合变换统计 PCollection 元素个数Apache Beam Java Kata 实战使用 Count 聚合变换统计 PCollection 元素个数 本指南围绕 Apache Beam 官方 J大数据批处理流处理数据工程Apache Beam Kotlin Katas 实战用 Partition 变换将 PCollection 按规则拆分为多个子集合Apache Beam Kotlin Katas 实战用 Partition 变换将 PCollection 按规则拆分为多个子集合 本篇技术指南围绕 Bea大数据批处理流处理数据工程Apache Beam Java 实战用 Min 聚合变换计算全局最小值Katas 入门篇Apache Beam Java 实战用 Min 聚合变换计算全局最小值Katas 入门篇 本文基于 Apache Beam 官方 Katas 课程中 大数据批处理流处理数据工程上一篇如何使用libimagequant生成高质量GIF掌握alpha通道处理技巧下一篇WebRTC-Experiment媒体流加密端到端加密保护通信隐私创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表