首页/新闻资讯/正文详情

Apache Beam Kotlin Katas 实战:Aggregation 之 Count 聚合变换详解

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

相关推荐

时尚网站模板代码拆解:3招搞定建站报价避坑指南
时尚网站模板代码拆解:3招搞定建站报价避坑指南

时尚网站模板代码拆解:3招搞定建站报价避坑指南 找建站公司最怕什么?不是技术牛不牛,而是报价单上那串让你心跳加速的数字。很多老板拿着“高端大气上档次”的需求,结果收到一份报价单,价格直接劝退。这时候,懂点 时尚网站模板代码… · 2026/9/27 8:34:43

双目标定与极线校正
双目标定与极线校正

双目标定与极线校正 文章目录双目标定与极线校正1. 为什么深度公式之前必须标定与校正2. 双目标定要估计的量3. 采集与求解流程3.1 采集标定板3.2 角点检测与立体标定3.3 立体校正与映射4. 验收指标:纵向残差与基线误差5. 合成棋盘格上的可复现演示6. 常见失效与处理… · 2026/9/27 8:34:43

做公众号微网站建设方案别踩坑:5个实战细节决定成败
做公众号微网站建设方案别踩坑:5个实战细节决定成败

做公众号微网站建设方案别踩坑:5个实战细节决定成败 还在用那种一眼假的模板网站?用户点进来两秒就划走,转化率为零,这不仅是面子问题,更是真金白银的流失。很多老板觉得找个模板套一下就行,结果上线后流量惨淡,最后还得返工重做,钱花了事没办成。… · 2026/9/27 8:34:30

isomorphic-git 请求头(Headers)完全指南:Authorization、User-Agent 与自定义 Headers 的浏览器与 Node 兼容性
isomorphic-git 请求头(Headers)完全指南:Authorization、User-Agent 与自定义 Headers 的浏览器与 Node 兼容性

开发工具 【免费下载链接】isomorphic-git A pure JavaScript implementation of git for node and browsers! 项目地址: https://gitcode.com/gh_mirrors/is/isomorphic-git 点击查看 免费下载 导读 isomorphic-git 是一款纯 JavaScript 实现的 Git,可… · 2026/9/27 9:05:24

TypeScript 枚举(Enum)实战全解:深入 The Concise TypeScript Book 的编译原理与类型系统
TypeScript 枚举(Enum)实战全解:深入 The Concise TypeScript Book 的编译原理与类型系统

文档教程 【免费下载链接】typescript-book The Concise TypeScript Book: A Concise Guide to Effective Development in TypeScript. Free and Open Source. 项目地址: https://gitcode.com/gh_mirrors/typ/typescript-book 点击查看 免费下载 TypeScript 的 enu… · 2026/9/27 9:05:18

Haskell 数据分析与处理工具库实战指南:解读 learnhaskell 仓库的 libraries.md
Haskell 数据分析与处理工具库实战指南:解读 learnhaskell 仓库的 libraries.md

教程 【免费下载链接】learnhaskell Learn Haskell 项目地址: https://gitcode.com/gh_mirrors/le/learnhaskell 点击查看 免费下载 导读 learnhaskell 是一个面向初学者的 Haskell 学习路径仓库,而 libraries.md 是其中一张"工具地图"——它… · 2026/9/27 9:05:18

网站建设软件有哪些内容揭秘,最佳实践避坑指南
网站建设软件有哪些内容揭秘,最佳实践避坑指南

网站建设软件有哪些内容揭秘,最佳实践避坑指南 网站做好了没人访问,这大概是老板们最头疼的事。钱花了几万,页面上线了,流量却比脸还干净。别急,这往往不是技术不行,而是你在选型和部署阶段就埋下了“死局”。… · 2026/9/27 9:05:12

排水管网检测手段怎么选?CCTV、声呐与管道机器人
排水管网检测手段怎么选?CCTV、声呐与管道机器人

一、先问一个问题:这次检测,管道里有多少水 排水管网检测手段的选择,第一个决定因素不是设备贵不贵,而是管内水位。同样一段管道,满水和半水状态适用的手段完全不同。把这一点先定下来,选型就成功了一半。… · 2026/9/27 9:05:00

服装公司网站怎么做才不亏?源码下载避坑全攻略
服装公司网站怎么做才不亏?源码下载避坑全攻略

服装公司网站怎么做才不亏?源码下载避坑全攻略 找建站公司最怕什么?怕被坑高价,更怕钱花出去了,手里连个像样的 源码下载 链接都拿不到,或者拿到一堆加密过的“黑盒”代码。干了十年建站,见过太多服装老板在装修完官网后,因为服务器到期、供应商跑路… · 2026/9/27 9:04:53

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码