ARTICLE DETAIL

资讯详情

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

Apache Beam 实战 Kata:使用 Sum 聚合变换计算 PCollection 元素总和

Apache Beam 实战 Kata:使用 Sum 聚合变换计算 PCollection 元素总和 大数据批处理流处理数据工程【免费下载链接】beamApache 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聚合变换对一个PCollectionInteger中所有元素求和。你将掌握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求和。原始任务描述非常精炼只有一句话KataCompute 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); PCollectionInteger numbers pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollectionInteger output applyTransform(numbers); output.apply(Log.ofElements()); pipeline.run(); } static PCollectionInteger applyTransform(PCollectionInteger 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列表导入为PCollectionInteger作为求和输入1 2 ... 10 55applyTransform(numbers)这就是需要你补全的核心方法——对输入PCollection施加求和变换Log.ofElements()来自 util/src/org/apache/beam/learning/katas/util/Log.java 的辅助变换内部基于ParDoDoFn将每个元素通过 SLF4J 日志打印出来非GlobalWindow时还会附带窗口信息同时把元素原样透传方便你在控制台观察运行结果pipeline.run()触发管线执行。答案实现在applyTransform方法中补全一行即可static PCollectionInteger applyTransform(PCollectionInteger 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.ValuesInteger values Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10); PCollectionInteger numbers testPipeline.apply(values); PCollectionInteger 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.GloballyInteger, Integer对PCollectionInteger全局求和空输入返回0Sum.integersPerKey()Combine.PerKeyK, Integer, Integer对PCollectionKVK, Integer按 Key 分组求和Sum.longsGlobally()Combine.GloballyLong, Long对PCollectionLong全局求和空输入返回0Sum.longsPerKey()Combine.PerKeyK, Long, Long按 Key 对Long值求和Sum.doublesGlobally()Combine.GloballyDouble, Double对PCollectionDouble全局求和空输入返回0Sum.doublesPerKey()Combine.PerKeyK, Double, Double按 Key 对Double值求和Sum.ofIntegers()Combine.BinaryCombineIntegerFn返回可复用的求和函数IntegerSum.ofLongs()Combine.BinaryCombineLongFn返回可复用的求和函数LongSum.ofDoubles()Combine.BinaryCombineDoubleFn返回可复用的求和函数Double以integersGlobally()为例其实现为一行委托public static Combine.GloballyInteger, 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()接收PCollectionKVK, Integer对每个不同 Key 分别求和输出PCollectionKVK, Integer每个 Key 对应一个求和结果。例如Sum.StringintegersPerKey()可用于统计每个用户/每个类别的数值总量与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()保证编写生产代码时无需额外判空测试范式TestPipelinePAssert是 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.javaSum变换源码sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Sum.java课程总览与工程设置learning/katas/java/README.md、course-info.yaml赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐小爱音箱接入大模型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 仓大数据批处理流处理数据工程上一篇KMS智能激活终极指南三步永久激活Windows和Office的完整教程下一篇Mbed TLS API 文档生成体系详解基于 Sphinx Breathe Doxygen 的 API 参考文档流水线创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表