- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文以 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、Max,Count 是第一个、也是理解其余四个聚合变换的基石。
1.2 本课的 Kata 目标
task.md 给出了本课的核心任务:
Kata:Count the number of elements from an input.(统计输入中元素的数量)
任务提示只有一个:使用Count变换。也就是说,本练习要求你接收一个包含 1 到 10 十个整数的PCollection<Int>,输出一个表示元素总数的PCollection<Long>(值为 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: Array<String>) { 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: PCollection<Int>): PCollection<Long> { 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)将内存中的十个整数包装为PCollection<Int>,这是无外部 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()的返回类型是PCollection<Long>。原因在底层实现:Count.globally()本质是Combine.globally(new CountFn<T>()),而CountFn<T>是CombineFn<T, long[], Long>——累加器用long[]承载计数,最终extractOutput取出long值装箱为Long。这也解释了为何 Kotlin 侧函数签名写的是PCollection<Long>而非PCollection<Int>:计数可能远超 Int 范围,Beam 采用Long类型天然支持大数量级。
三、Count 家族的三种形态与适用场景
本练习只要求全局计数,但Count变换实际提供三个入口方法(见 Count.java 源码注释):
| 方法 | 输入 | 输出 | 语义 | 典型场景 |
|---|---|---|---|---|
Count.globally() | PCollection<T> | PCollection<Long> | 整个集合的元素总数 | 统计总记录条数、事件总量 |
Count.perElement() | PCollection<T> | PCollection<KV<T, Long>> | 每个不同元素出现的次数 | 词频统计、去重后各类别计数 |
Count.perKey() | PCollection<KV<K, V>> | PCollection<KV<K, Long>> | 每个 Key 关联的 Value 数量 | 按用户/地区/商品等维度计数 |
三种形态分别对应聚合的三种粒度:全局(globally)、按值(perElement)、按键(perKey)。在真实业务中,perElement()可直接实现 WordCount 中的词频统计:
val wordCounts: PCollection<KV<String, Long>> = words.apply(Count.perElement())3.1 perElement() 的实现方式
从源码看,Count.java 中的PerElement<T>变换分两步完成:先用MapElements把每个元素T映射为KV<T, Void>(KV.of(element, null)),再交给Count.perKey()按元素本身作为 Key 计数。因此perElement()本质上是perKey()的特例。
3.2 一个值得注意的约束
perElement()判定"元素相等"的方式,是把元素用输入PCollection的Coder编码后再比较字节(Coder#verifyDeterministic()),因此要求输入 Coder 必须是确定性的;globally()与perKey()则无此约束。若输入使用了非全局窗口的窗口策略,Count.globally()会直接抛出"不兼容全局窗口"的错误提示,源码中的getIncompatibleGlobalWindowErrorMessage()明确建议改用Combine.globally(Count.<T>combineFn()).withoutDefaults()来处理开窗场景下的计数。
四、源码级原理:CountFn 累加器如何工作
Count是典型的Combine变换,其精髓在于内部的CountFn<T>(见 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()返回的PCollection<Long>类型严格对应;containsInAnyOrder不关心元素顺序,适用于聚合结果这类单元素输出; run().waitUntilFinish()确保管道执行完毕、断言生效。
如果你的实现误用了Count.perElement()(输出会是 10 个KV<Int, Long>)或返回值类型写错,PAssert 都会因输出与10L不匹配而失败——这正是"填空式"教学设计的精妙之处。
六、举一反三:向 Sum / Mean / Min / Max 迁移
掌握 Count 后,Aggregation 课程中其余四个变换(见 lesson-info.yaml)几乎可以零成本迁移:
- Sum:
Sum.integersGlobally()等按数值类型区分的全局求和; - Mean:
Mean.globally()计算全局均值; - Min / Max:
Min.globally()/Max.globally()求全局最值。
它们与Count.globally()一样,都是Combine.globally(...)的便捷封装,输出同样为单元素PCollection(如PCollection<Long>或数值类型),测试断言模式也完全一致。理解了 Count 的CombineFn机制,就理解了整个 Aggregation 课程背后的统一抽象。
七、总结
- Count 是 Beam 聚合变换的入门第一课,
Count.globally()一行代码即可完成全局计数,返回PCollection<Long>; - 按需扩展时,
perElement()做词频类统计、perKey()做维度分组计数; - 其底层
CountFn采用long[]可变累加器 + 分布式合并,兼顾正确性与可扩展性,相关实现可在 Count.java 中完整查阅; - 配合 TaskTest.kt 的
PAssert断言模式,你可以把同样的测试方法复用到任何自定义聚合变换的验证中。
完成本 Kata 后,不妨直接打开 Aggregation 课程中下一个练习 Sum,你会发现:聚合的思想是统一的,变的只是CombineFn内部的运算规则。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
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 课程中 "
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考