- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文以 Apache Beam 仓库中learning/katas/java/Common Transforms/Aggregation/Sum这一练习任务(Kata)为核心,讲解如何使用 Beam Java SDK 内置的Sum聚合变换,对一个PCollection<Integer>中所有元素求和。你将掌握Sum.integersGlobally()的用法、Sum变换在源码层面的实现原理(基于Combine的全局聚合与按 Key 聚合)、配套测试的编写方式,以及在 IntelliJ EduTools 环境中运行本练习的方法,最终能够独立完成并验证该练习。
任务概览:本 Kata 要求什么
本练习位于 Beam 的 Java Katas 课程目录 learning/katas/java/Common Transforms/Aggregation/Sum 下,属于Common Transforms → Aggregation(聚合)单元。该单元在 lesson-info.yaml 中定义了五个聚合练习:Count、Sum、Mean、Min、Max,本任务聚焦其中的Sum(求和)。
原始任务描述非常精炼,只有一句话:
Kata:Compute the sum of all elements from an input.(计算输入中所有元素的总和。)
提示信息明确给出了解法方向:使用org.apache.beam.sdk.transforms.Sum这个变换。
从配套的 task-info.yaml 可以看到本练习在 EduTools 课程中的结构:type: edu表明这是一个带TODO()占位符的编程练习,需要你在Task.java的占位位置补全实现(placeholder 位于offset: 1965、length: 42处),而测试文件TaskTest.java对学习者不可见(visible: false),用于自动校验答案。这正是 Beam Katas 课程"先想、再写、后验证"的教学模式。
完整解题代码:从骨架到实现
任务骨架(学习者看到的代码)
练习起始代码位于 Task.java,主体结构如下:
package org.apache.beam.learning.katas.commontransforms.aggregation.sum; import org.apache.beam.learning.katas.util.Log; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.transforms.Sum; import org.apache.beam.sdk.values.PCollection; public class Task { public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline = Pipeline.create(options); PCollection<Integer> numbers = pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollection<Integer> output = applyTransform(numbers); output.apply(Log.ofElements()); pipeline.run(); } static PCollection<Integer> applyTransform(PCollection<Integer> input) { return input.apply(Sum.integersGlobally()); // TODO() 占位处需要你写出这一行 } }代码逐段拆解:
PipelineOptionsFactory.fromArgs(args).create():解析命令行参数并创建PipelineOptions,这是 Beam 管线的标准入口写法;Pipeline.create(options):基于选项创建Pipeline实例;Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10):将一个Integer列表导入为PCollection<Integer>,作为求和输入(1 + 2 + ... + 10 = 55);applyTransform(numbers):这就是需要你补全的核心方法——对输入PCollection施加求和变换;Log.ofElements():来自 util/src/org/apache/beam/learning/katas/util/Log.java 的辅助变换,内部基于ParDo+DoFn将每个元素通过 SLF4J 日志打印出来(非GlobalWindow时还会附带窗口信息),同时把元素原样透传,方便你在控制台观察运行结果;pipeline.run():触发管线执行。
答案实现
在applyTransform方法中补全一行即可:
static PCollection<Integer> applyTransform(PCollection<Integer> input) { return input.apply(Sum.integersGlobally()); }运行后控制台会依次输出 1 到 10 这 10 个元素,并最终输出聚合结果55(因为Log.ofElements()也会打印求和后唯一的输出元素)。注意:Sum.integersGlobally()是全局聚合,无论输入元素如何分布,输出PCollection中只有一个元素——即全部元素的总和。
测试校验:PAssert 断言总和为 55
本练习的自动化测试位于 TaskTest.java,它演示了 Beam 中验证聚合结果的标准测试模式:
package org.apache.beam.learning.katas.commontransforms.aggregation.sum; import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.values.PCollection; import org.junit.Rule; import org.junit.Test; public class TaskTest { @Rule public final transient TestPipeline testPipeline = TestPipeline.create(); @Test public void sum() { Create.Values<Integer> values = Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10); PCollection<Integer> numbers = testPipeline.apply(values); PCollection<Integer> results = Task.applyTransform(numbers); PAssert.that(results) .containsInAnyOrder(55); testPipeline.run().waitUntilFinish(); } }要点分析:
TestPipeline.create()是 Beam 的测试专用Pipeline,配合 JUnit 的@Rule使用,测试结束时会自动执行并校验管线;PAssert.that(results).containsInAnyOrder(55):这是断言的核心——它声明结果PCollection中恰好包含一个元素55。containsInAnyOrder不关心元素顺序,只校验集合内容;testPipeline.run().waitUntilFinish():真正触发管线执行并等待完成;若结果与断言不符,测试将失败并给出具体差异;- 该测试直接调用
Task.applyTransform(numbers),与你补全的实现逻辑完全解耦,因此无论你如何实现求和,只要语义正确(结果为 55)即可通过测试——这也符合 Kata 练习"实现与测试分离"的设计理念。
源码剖析:Sum 变换的底层实现原理
Sum是 Beam Java SDK 中一个轻量的工厂类(PTransform集合),定义于 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sum.java。它本身并不直接实现求和逻辑,而是委托给Combine变换,并为三种数值类型(Integer、Long、Double)各提供一组静态工厂方法。
方法族一览(来自 Sum.java 源码)
| 方法 | 返回类型 | 语义 |
|---|---|---|
Sum.integersGlobally() | Combine.Globally<Integer, Integer> | 对PCollection<Integer>全局求和;空输入返回0 |
Sum.integersPerKey() | Combine.PerKey<K, Integer, Integer> | 对PCollection<KV<K, Integer>>按 Key 分组求和 |
Sum.longsGlobally() | Combine.Globally<Long, Long> | 对PCollection<Long>全局求和;空输入返回0 |
Sum.longsPerKey() | Combine.PerKey<K, Long, Long> | 按 Key 对Long值求和 |
Sum.doublesGlobally() | Combine.Globally<Double, Double> | 对PCollection<Double>全局求和;空输入返回0 |
Sum.doublesPerKey() | Combine.PerKey<K, Double, Double> | 按 Key 对Double值求和 |
Sum.ofIntegers() | Combine.BinaryCombineIntegerFn | 返回可复用的求和函数(Integer) |
Sum.ofLongs() | Combine.BinaryCombineLongFn | 返回可复用的求和函数(Long) |
Sum.ofDoubles() | Combine.BinaryCombineDoubleFn | 返回可复用的求和函数(Double) |
以integersGlobally()为例,其实现为一行委托:
public static Combine.Globally<Integer, Integer> integersGlobally() { return Combine.globally(Sum.ofIntegers()); }也就是说,Sum.integersGlobally()本质上是Combine.globally(new SumIntegerFn())——把"如何合并两个元素"的二元函数交给Combine,由Combine负责分布式聚合的调度与优化。
求和函数内部:二元合并 + 恒等元
SumIntegerFn、SumLongFn、SumDoubleFn分别继承Combine.BinaryCombineIntegerFn、BinaryCombineLongFn、BinaryCombineDoubleFn,各自只实现三个核心方法(以SumIntegerFn为例):
private static class SumIntegerFn extends Combine.BinaryCombineIntegerFn { @Override public int apply(int a, int b) { return a + b; // 二元合并:两两相加 } @Override public int identity() { return 0; // 恒等元:空输入时返回 0 } ... }这里的identity()正是前面方法表中"空输入返回0"语义的源码出处——Combine在没有任何元素可聚合时,直接以恒等元0作为结果。这种"二元函数 + 恒等元"的设计使Combine能够以树形合并方式并行计算,是 Beam 聚合变换高效处理大规模数据的关键(Combine的分布式优化、增量合并等细节由org.apache.beam.sdk.transforms.Combine实现)。
全局求和 vs 按 Key 求和
Sum同时提供了Globally(全局)与PerKey(按 Key)两类入口,这是 Beam 聚合的两个基本维度:
- 全局求和:
Sum.integersGlobally()将整个PCollection归约为单个元素,对应本 Kata 的场景; - 按 Key 求和:
Sum.integersPerKey()接收PCollection<KV<K, Integer>>,对每个不同 Key 分别求和,输出PCollection<KV<K, Integer>>(每个 Key 对应一个求和结果)。例如Sum.<String>integersPerKey()可用于统计每个用户/每个类别的数值总量,与GroupByKey相比更简洁高效,因为Combine会在洗牌前做本地预聚合。
从源码结构可以推断,Sum是 Beam "标准聚合变换"(Count、Max、Mean、Min、Sum一族)中的求和实现,开发者完全可以基于它写出自己的Combine组合逻辑。
运行与验证:在 IntelliJ 中完成练习
本练习按 Beam Katas Java 课程组织,运行方式遵循课程统一的工程化流程:
- 导入课程工程:使用 IntelliJ IDEA(教育版,或安装 EduTools 插件)打开 learning/katas/java 目录,详见 learning/katas/java/README.md 的 Setup 说明;
- 导入 Gradle 项目:按提示选择 "Import Gradle project" 并等待构建完成;随后在 "Project Structure" 中配置项目 SDK(如 JDK 8,与 course-info.yaml 中
programming_language_version: 8一致); - 进入课程视图:打开 "Project" 工具窗口,切换到 "Course" 视图,即可看到
Common Transforms → Aggregation → Sum练习; - 补全代码:在
Task.java的TODO()位置写入input.apply(Sum.integersGlobally()); - 运行与验证:直接运行
Task.main可在控制台观察输出(应为 55);点击课程的测试按钮或运行TaskTest,通过PAssert自动校验答案。若补全正确,测试通过;若结果不等于 55,测试失败并提示期望值与实际值的差异。
延伸思考:从 Kata 到真实管线
完成本练习后,你已经掌握 Beam 聚合变换的通用心智模型:
- 聚合三要素:输入
PCollection、聚合函数(如Sum内部的二元合并函数)、聚合维度(Globally全局 /PerKey按 Key / 结合Windowing按窗口); - 空输入语义:
Sum系列的全局聚合在输入为空时返回恒等元0,这由identity()保证,编写生产代码时无需额外判空; - 测试范式:
TestPipeline+PAssert是 Beam 官方推荐的断言方式,containsInAnyOrder适合聚合这类"结果为一个元素"的场景;PAssert还支持containsInAnyOrder、empty()等更丰富的断言; - 迁移路径:将输入改为真实数据源(Kafka、Pub/Sub、文件等,见仓库 learning/katas/java/IO 练习),把
Sum替换为Mean、Max、Min或自定义Combine,即可把本练习的骨架直接复用到实际聚合场景中。
参考文件索引
- 练习任务描述:task.md
- 练习起始代码与答案:Task.java
- 自动化测试:TaskTest.java
- 课程结构配置:task-info.yaml 与 lesson-info.yaml
- 日志辅助变换实现:Log.java
Sum变换源码:sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sum.java- 课程总览与工程设置:learning/katas/java/README.md、course-info.yaml
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
小爱音箱接入大模型:MiGPT 部署与配置完整指南
小爱音箱接入大模型:MiGPT 部署与配置完整指南 晚上问小爱同学"为什么天空是蓝色的",它还是那句模板式的客服腔。MiGPT 是一个把小爱音箱接入 ChatG
大数据批处理流处理数据工程zerotier-cli join 后返回 access denied 时如何确认设备已被控制器授权?
zerotier cli join 后返回 access denied 时如何确认设备已被控制器授权? 在 Unix 系统(Linux/BSD/OSX)上用 z
大数据批处理流处理数据工程Apache Beam Java 示例实战:用 Beam SQL 与 Schema Transforms 计算按键聚合指标
Apache Beam Java 示例实战:用 Beam SQL 与 Schema Transforms 计算按键聚合指标 本文基于 Apache Beam 仓
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考