☰
Apache Beam 实战 Kata:使用 Sum 聚合变换计算 PCollection 元素总和
2026/9/26 8:24:20 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

导读

本文以 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 课程组织,运行方式遵循课程统一的工程化流程:

  1. 导入课程工程:使用 IntelliJ IDEA(教育版,或安装 EduTools 插件)打开 learning/katas/java 目录,详见 learning/katas/java/README.md 的 Setup 说明;
  2. 导入 Gradle 项目:按提示选择 "Import Gradle project" 并等待构建完成;随后在 "Project Structure" 中配置项目 SDK(如 JDK 8,与 course-info.yaml 中programming_language_version: 8一致);
  3. 进入课程视图:打开 "Project" 工具窗口,切换到 "Course" 视图,即可看到Common Transforms → Aggregation → Sum练习;
  4. 补全代码:在Task.java的TODO()位置写入input.apply(Sum.integersGlobally());
  5. 运行与验证:直接运行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.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

相关推荐

上一篇:KMS智能激活终极指南:三步永久激活Windows和Office的完整教程
下一篇:Mbed TLS API 文档生成体系详解:基于 Sphinx + Breathe + Doxygen 的 API 参考文档流水线

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询