Apache Beam Kotlin 实战:使用 Sum 聚合变换计算 PCollection 元素总和
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载Apache Beam 的Sum变换用于计算PCollection中全部元素的总和全局聚合或按 Key 分组后计算每组值的总和按 Key 聚合是 Beam 官方学习路径 Katas 中通用变换Common Transforms→ 聚合Aggregation系列的核心一课。本文以仓库中learning/katas/kotlin/Common Transforms/Aggregation/Sum/task.md这一实战练习为主线结合Task.kt完整实现、TaskTest.kt单元测试以及 SDK 源码Sum.java带你彻底掌握 Kotlin 语言下 Sum 变换的用法、底层 Combine 机制与测试验证方式。一、任务背景Katas 中的 Sum 练习Apache Beam 在仓库中内置了一套面向初学者的编程练习Katas涵盖 Java、Kotlin、Go、Python 四种语言。其中 Kotlin 版本的通用变换 - 聚合课程见 lesson-info.yaml依次包含Count、Sum、Mean、Min、Max五个小节Sum 排在第二位。本小节 task.md 给出的练习描述只有一句话Kata:Compute the sum of all elements from an input.计算输入中所有元素的总和练习提示明确指出应使用Sum变换。这是一个典型的填空式练习task-info.yaml中配置了练习文件与占位符placeholder_text: TODO()位于 task-info.yaml学习者需要把占位符替换为正确的Sum调用代码才能通过随附的单元测试。二、完整实现一行代码完成全局求和本练习的完整参考实现位于 Task.kt核心代码只有一行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.PipelineOptionsFactory import org.apache.beam.sdk.transforms.Create import org.apache.beam.sdk.transforms.Sum 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() } fun applyTransform(input: PCollectionInt): PCollectionInt { return input.apply(Sum.integersGlobally()) } }逐段拆解这个程序创建 PipelinePipelineOptionsFactory.fromArgs(*args).create()从命令行参数解析运行选项Pipeline.create(options)构建流水线后续所有变换都通过pipeline.apply(...)挂载。构造输入Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)创建一个包含 1 到 10 共 10 个整数的PCollectionInt。应用核心变换applyTransform内部执行input.apply(Sum.integersGlobally())即对整条PCollection做全局求和结果是一个只包含单个元素55的PCollectionInt。输出与运行Log.ofElements()把结果元素打印到日志pipeline.run()触发执行。需要注意Sum在 Kotlin/Java SDK 中返回的是Combine.GloballyInteger, Integer类型即它本质上是Combine.globally(Sum.ofIntegers())的便捷封装下文详述。由于 Sum 变换是全局聚合它要求输入 PCollection 处于全局窗口Global Window若输入元素带有时间戳、分布在多个窗口或触发器中聚合语义会发生变化这一点在流式场景下要特别留意。三、测试验证PAssert 断言总和为 55练习的验收标准由 TaskTest.kt 定义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.junit.Rule import org.junit.Test class TaskTest { get:Rule Transient val testPipeline: TestPipeline TestPipeline.create() Test fun common_transforms_aggregation_sum() { val values Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10) val numbers testPipeline.apply(values) val results Task.applyTransform(numbers) PAssert.that(results).containsInAnyOrder(55) testPipeline.run().waitUntilFinish() } }该测试揭示了几个重要的工程实践TestPipeline 替代手动创建用TestPipeline.create()代替Pipeline.create(...)由测试框架自动管理流水线生命周期run()后立即waitUntilFinish()阻塞等待执行结束。PAssert 断言聚合结果PAssert.that(results).containsInAnyOrder(55)验证输出 PCollection 中恰好包含元素5512…1055且不关心元素顺序——这是 Beam 官方测试工具库 sdks/java/testing 中org.apache.beam.sdk.testing.PAssert的典型用法。练习的正确性闭环只有当学习者把task-info.yaml中标记的TODO()占位符替换为Sum.integersGlobally()或等价实现时该测试才会通过从而构成练习—验证的自动判题机制。四、源码级剖析Sum 变换的完整 API 家族练习中的Sum.integersGlobally()只是Sum工具类的一个静态方法。查看 SDK 核心源码 Sum.java 可知该类针对三种数值类型提供了六种聚合入口覆盖全局求和与按 Key 求和两类场景方法输入类型输出类型适用场景Sum.integersGlobally()PCollectionIntegerPCollectionInteger全部整数元素求和本练习Sum.integersPerKey()PCollectionKVK, IntegerPCollectionKVK, Integer每个 Key 对应整数值求和Sum.longsGlobally()PCollectionLongPCollectionLong全部 Long 元素求和Sum.longsPerKey()PCollectionKVK, LongPCollectionKVK, Long每个 Key 对应 Long 值求和Sum.doublesGlobally()PCollectionDoublePCollectionDouble全部 Double 元素求和Sum.doublesPerKey()PCollectionKVK, DoublePCollectionKVK, Double每个 Key 对应 Double 值求和从源码可以看出这些方法内部全部委托给Combine变换public static Combine.GloballyInteger, Integer integersGlobally() { return Combine.globally(Sum.ofIntegers()); } public static K Combine.PerKeyK, Integer, Integer integersPerKey() { return Combine.perKey(Sum.ofIntegers()); }也就是说Sum 并不是一个独立的 DoFn而是 Beam Combine 框架的一个特化实例。它通过Sum.ofIntegers()/ofLongs()/ofDoubles()返回对应的二元合并函数Combine.BinaryCombineIntegerFn/BinaryCombineLongFn/BinaryCombineDoubleFn例如private static class SumIntegerFn extends Combine.BinaryCombineIntegerFn { Override public int apply(int a, int b) { return a b; } Override public int identity() { return 0; } // equals / hashCode 基于类型实现便于分布式环境中序列化与合并 }这里有两个值得注意的设计细节identity()返回 0这是加法的单位元。Combine 框架在合并部分结果时会用到该值同时它也是空输入没有任何元素时全局聚合的兜底结果——即Sum.integersGlobally()作用在空 PCollection 上会输出0而不是报错。apply(a, b) a b二元合并函数满足结合律因此 Beam 可以把聚合分解到集群的多个节点上并行计算局部和再逐层归并出最终结果这正是Sum能高效处理大规模数据集的原因。五、从练习到实战Sum 的典型扩展用法掌握 Kata 之后可以按需把 Sum 泛化到真实业务1. 按 Key 分组求和——统计每个用户的总消费金额、每个商品的累计销量等。示例源自 Sum.java 的文档注释// 输入 PCollectionKVString, IntegerKey 为商品 IDValue 为单次销量 val input: PCollectionKVString, Integer ... val sumPerKey: PCollectionKVString, Integer input.apply(Sum.integersPerKey())2. Long / Double 精度场景涉及超大整数或浮点金额统计时改用Sum.longsGlobally()或Sum.doublesGlobally()按 Key 版本同理。3. 自定义 Combine 扩展如果需求是求积求最大差值等 Sum 未提供的聚合可以直接继承Combine.BinaryCombineIntegerFn等抽象类实现apply与identity后交给Combine.globally(...)/Combine.perKey(...)复用同一套分布式合并机制。六、如何运行与继续学习本练习位于 Kotlin Katas 工程内目录结构遵循统一的任务—测试约定每个小节都有独立的src学习者待补全的实现、test自动判题测试与task.md题目描述。Sum 练习的完整文件如下题目描述task.md待补全实现Task.kt自动判题测试TaskTest.kt练习元信息含 TODO 占位符配置task-info.yamlKotlin Katas 工程自带 Gradle 构建脚本gradlew在 learning/katas/kotlin 目录下执行./gradlew test即可运行所有小节的测试通过即代表练习完成。若想查看 Kotlin Katas 的整体结构可阅读其 README.md 与 course-info.yaml。完成 Sum 之后建议按同一路径继续学习聚合课程中的Mean、Min、Max练习见 lesson-info.yaml它们在代码结构上与 Sum 完全同构可以巩固对 Combine 聚合体系的理解。对于想深挖底层实现的读者可以从 Sum.java 出发进一步阅读Combine与CombineFn的源码理解结合律聚合 分布式归并这一 Beam 高性能聚合的核心原理。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam 实战 Kata使用 Sum 聚合变换计算 PCollection 元素总和Apache Beam 实战 Kata使用 Sum 聚合变换计算 PCollection 元素总和 导读 本文以 Apache Beam 仓库中 learni大数据批处理流处理数据工程Apache Beam Java Katas使用 Sum 聚合变换计算 PCollection 元素总和Apache Beam Java Katas使用 Sum 聚合变换计算 PCollection 元素总和 本指南以 learning/katas/java/C批处理流处理大数据Apache Beam Kotlin Kata 实战使用 Sum 变换计算元素总和Common Transforms 之 Aggregation/SumApache Beam Kotlin Kata 实战使用 Sum 变换计算元素总和Common Transforms 之 Aggregation/Sum大数据批处理流处理数据工程上一篇终极选择Obsidian-Git分支策略深度对比与实践指南下一篇WPF中的多窗口通信共享数据上下文创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

pstack-claude 实战:用 pstack 快速定位进程卡死与线程阻塞

pstack-claude 实战:用 pstack 快速定位进程卡死与线程阻塞

1. 从 pstack-claude 说起:一个被低估的进程栈排查利器第一次看到pstack-claude这个名字,很多人会以为它是某个新出的 AI 工具链,或者跟 Claude 模型有什么绑定关系。实际上,pstack本身是一个存在了二十多年的经典命令行工具&…

2026/10/9 9:39:55 阅读更多 →
t3code 代码片段索引方案:从 grep 到高效检索的工程实践

t3code 代码片段索引方案:从 grep 到高效检索的工程实践

1. 从“t3code”这个名字说起:它到底指什么第一次看到“t3code”这个词,很多人会一头雾水。它不像“React”“Vue”那样有明确的官方文档,也不像“Python”那样有庞大的社区。我在几个技术群里问了一圈,发现大家对它的理解分成好几…

2026/10/9 9:39:55 阅读更多 →
Excel打底SQL提效BI收口:数据分析完整链路实战指南

Excel打底SQL提效BI收口:数据分析完整链路实战指南

简介:一份面向数据分析初学者与业务人员的实战型课件,共94页,系统讲解如何用Excel与SQL完成数据采集、处理、分析与可视化。内容涵盖数据分析概念与前景、Excel数据导入与常用函数、数据透视表与图表、SQL数据库基础及CRUD操作,并…

2026/10/9 9:39:55 阅读更多 →

最新新闻

Altium Designer 17.0.6安装避坑指南:从环境检查到静默部署的完整方案

Altium Designer 17.0.6安装避坑指南:从环境检查到静默部署的完整方案

简介:Altium Designer 17.0.6安装教程PDF,面向电子设计工程师及PCB初学者,解决Altium Designer软件安装、破解与汉化流程不熟悉的问题。资源包内共1个pdf文件,整体大小3.03MB,内容紧凑,以图文步骤方式组织&…

2026/10/9 14:55:27 阅读更多 →
5G网络切片隔离性验证:从测试设计到pytest自动化落地

5G网络切片隔离性验证:从测试设计到pytest自动化落地

去年做5G行业专网交付的时候,客户在验收会上问了我一个很要命的问题:"你说切片隔离,那我车间里的视频监控流量和AGV控制流量在同一个基站下跑,监控业务能不能把控制业务挤垮?你拿什么证明它不会?"…

2026/10/9 14:55:27 阅读更多 →
Markdown编辑器从入门到进阶:选型、语法与工作流全攻略

Markdown编辑器从入门到进阶:选型、语法与工作流全攻略

刚拿到一个新的 .md 文件时,很多人第一反应是双击打开,然后看到满屏的 # 和 *,第一反应是:这文件是不是坏了?我当年也一样,项目文档发过来,我以为是文本乱码,差点把文件删了。后来才…

2026/10/9 14:55:27 阅读更多 →
GitHub热榜解码:技术趋势识别与工程化落地指南

GitHub热榜解码:技术趋势识别与工程化落地指南

1. 项目概述:这不是一份榜单,而是一份开源世界的实时脉搏图“GitHub 热榜项目:周榜(2026-10-04)”——看到这个标题,很多人第一反应是点开链接、扫一眼排名、记下几个耳熟的仓库名,然后关掉页面…

2026/10/9 14:55:27 阅读更多 →
GitHub日榜数据采集与验证:构建可复现的热榜观测体系

GitHub日榜数据采集与验证:构建可复现的热榜观测体系

1. 热榜不是排行榜,而是开发者的行为镜像“GitHub 日榜(2026-10-04)”这个标题乍看像一份静态榜单,但实际它是一扇实时窗口——透过它,你能看到全球开发者在这一天集体关注什么、正在解决什么真实问题、又在用什么新方…

2026/10/9 14:55:27 阅读更多 →
PHP 接口压力测试实战:用 wrk 从零搭建性能基线

PHP 接口压力测试实战:用 wrk 从零搭建性能基线

很多同学第一次给 PHP 项目做压力测试,第一反应都是:找个工具把接口打爆,看看会不会挂。我以前也这么想,直到有次新功能上线前我用 wrk 随手压了一套接口,发现 QPS 连 200 都没到,RT 的 P99 却已经飙到 1.5…

2026/10/9 14:54:26 阅读更多 →

日新闻

Java时间API实战:LocalDate、Date与ZonedDateTime的转换与避坑指南

Java时间API实战:LocalDate、Date与ZonedDateTime的转换与避坑指南

Java时间API这个话题,隔三差五就会在群里被翻出来讨论一次。上周还有个同事线上处理一个订单超时问题,排查到最后发现是ZonedDateTime序列化后时区丢了,用户在下单当天晚上看到的时间整整差了8个小时。这类问题几乎每个做Java开发的人都遇到过…

2026/10/9 0:00:49 阅读更多 →
EasyTier实践:从NAT穿透到子网代理的异地组网部署与排错

EasyTier实践:从NAT穿透到子网代理的异地组网部署与排错

前几个月我手头有好几台机器需要互相访问:办公室台式机、家里 NAS、还有一台云主机。如果只是偶尔传个文件倒还好,问题是工作场景经常要在几处环境之间来回切换,每次都先登录跳板机再层层代理,实在折腾。我先后试过端口映射、自建…

2026/10/9 0:00:49 阅读更多 →
AI Agent工程实战:从七要素到七个决策点的系统设计指南

AI Agent工程实战:从七要素到七个决策点的系统设计指南

AI Agent 这个词在过去一年里被反复提及,但真正动手搭过一套能跑起来的 Agent 系统的人都知道,从"知道它是什么"到"让它稳定干活"之间隔着一整套工程决策。我前后参与过几个 Agent 项目的落地,从最初用现成框架拼装&…

2026/10/9 0:01:50 阅读更多 →

周新闻

KT148A语音芯片外挂8002D功放的工程实践指南

KT148A语音芯片外挂8002D功放的工程实践指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/8 15:26:32 阅读更多 →
LLC谐振变换器增益公式推导:从FHA等效到完整归一化表达式

LLC谐振变换器增益公式推导:从FHA等效到完整归一化表达式

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/8 15:26:40 阅读更多 →
ARM架构深度解析:从RISC设计理念到交叉编译实战

ARM架构深度解析:从RISC设计理念到交叉编译实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/9 10:11:06 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/8 21:13:17 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/8 15:26:17 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/9 6:17:20 阅读更多 →