资讯详情

资讯详情

Apache Beam Java 入门实战:用 Mean 组合器计算 PCollection 元素均值

【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Apache Beam 是一个统一的批处理与流式数据处理编程模型。在学习 katas练习的 Common Transforms / Aggregation 一节中Mean 练习要求读者掌握一个非常基础且高频的聚合能力计算一个PCollectionInteger中所有元素的算术平均值并将结果输出为PCollectionDouble。完成本练习后你将熟悉 Beam Java SDK 中Mean.globally()的用法、Combine聚合框架的基本形态以及如何用PAssert对聚合结果做单元测试验证。练习目标Kata 定义本节练习的原始任务描述位于 Mean/task.md核心要求只有一句话Kata:Compute the mean/average of all elements from an input.即对输入中的所有元素计算均值/平均值。官方提示是使用 Beam 提供的现成组合器 Mean而无需手写Combine.CombineFn或DoFn累加逻辑。该练习属于 Aggregation 一课与 Count、Sum、Min、Max 共同组成了聚合类 Common Transforms 的基础训练序列参见 lesson-info.yamlMean 位于Count、Sum之后是第三个练习。完成任务修改applyTransform方法练习采用 EduToolsIntelliJ Education格式task-info.yaml声明了需要填空的位置placeholder_text: TODO()位于src/.../mean/Task.java中。也就是说你需要在自己的工作副本中把 Task.java 里applyTransform方法内的TODO()占位替换为正确的实现。完成后的参考实现如下static PCollectionDouble applyTransform(PCollectionInteger input) { return input.apply(Mean.globally()); }applyTransform的签名揭示了两个关键信息输入类型PCollectionInteger本练习的输入固定为Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)产生的 110 共 10 个整数输出类型PCollectionDouble——Mean.globally()的结果类型是Double不是输入元素的原始类型Integer。入口main方法中output.apply(Log.ofElements())会把均值结果打印到日志。Log是 katas 项目提供的一个通用 PTransform 封装实现见 Log.java它通过ParDoDoFn将每个元素若处于非全局窗口还会附带窗口信息以 SLF4J 日志输出并原样透传元素方便在控制台观察结果。结果验证PAssert 断言 5.5练习配套的测试 TaskTest.java 验证了实现正确性Test public void mean() { Create.ValuesInteger values Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10); PCollectionInteger numbers testPipeline.apply(values); PCollectionDouble results Task.applyTransform(numbers); PAssert.that(results) .containsInAnyOrder(5.5); testPipeline.run().waitUntilFinish(); }测试用TestPipelineRule构建管道复用被测的applyTransform再用PAssert.that(...).containsInAnyOrder(5.5)断言输出集合只包含一个值5.5——即 (12…10)/10 55/10 的结果。这是 Beam 官方测试框架验证PCollection内容的惯用写法聚合类 Transform 的正确性都可以用这种方式验证。对于Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)这类输入最终日志输出应显示5.5。深入原理Mean的源码实现Mean类位于 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Mean.java它本身是一个命名空间类对外提供两个静态工厂方法Mean.NumTglobally()对整条PCollectionNumT计算全局均值返回Combine.GloballyNumT, Double源码 L67-L69Mean.K, NumTperKey()按PCollectionKVK, NumT中每个 Key 分组后分别计算均值返回Combine.PerKeyK, NumT, Double源码 L82-L84。从源码结构看两者都只是把Mean.of()即MeanFn包装进Combine.globally(...)/Combine.perKey(...)底层走的是 Beam 的Combine 聚合框架——由CombineFn定义累加、合并、输出三阶段逻辑而非简单的逐元素ParDo。MeanFn继承自Combine.AccumulatingCombineFnNumT, CountSumNumT, Double其内部累加器CountSum源码 L121-L174只维护两个字段long count 0; double sum 0.0;三个关键方法决定了均值计算的精度与语义addInput(NumT element)count并把element.doubleValue()累加进sum——这意味着输入元素会先被转换为double再参与求和mergeAccumulator(CountSum accumulator)将两个累加器的count与sum分别相加这正是 Beam 在分布式执行时把多个分区/分片的局部统计合并起来的方式也是Mean可以随Combine框架并行伸缩的根本原因extractOutput()返回count 0 ? Double.NaN : sum / count源码 L148-L150——即最终均值由sum / count计算得到。这里有两个值得注意的边界语义空输入返回Double.NaN当没有任何元素时extractOutput()返回Double.NaN。而Mean.globally()的文档源码 L60-L69说明的是返回0——这指的是globally()包装层在空集合上的行为Mean.of()作为CombineFn的语义则是零元素时输出NaN。两者结合看Mean家族对空集合求均值的处理并不是简单的 0实际使用时需要结合Combine的窗口/默认值策略来确认最终结果。Double精度与溢出累加过程把每个数值转成double求和count用long统计。对于超大整数集合double求和可能出现浮点精度损失这一点从CountSum的实现可以推断对绝大多数数据集sum / count的精度足够。此外MeanFn还实现了getAccumulatorCoder()返回CountSumCoder源码 L176-L197用BigEndianLongCoder编码count、DoubleCoder编码sum保证累加器在分布式节点间序列化传输时编码是确定性的verifyDeterministic()对两个子 coder 做了确定性校验。延伸perKey 按组求均值Mean.globally()覆盖了整个 PCollection 一个均值的场景当需要按维度分组求均值时例如统计每个用户订单的平均金额应使用Mean.perKey()其用法在Mean类 Javadoc 中有示例源码 L47-L54PCollectionKVString, Integer input ...; PCollectionKVString, Double meanPerKey input.apply(Mean.String, IntegerperKey());输出为PCollectionKVK, Double每个不同的 Key 对应其所有 Value 的均值。在无界数据流或使用窗口的场景下perKey()的行为受Combine.PerKey的时间戳与分桶策略影响需结合窗口Windowing与触发Triggers一起设计。运行方式本 katas 项目采用 IntelliJ Education / EduTools 打开详见 learning/katas/java/README.md打开learning/katas/java目录并导入 Gradle 项目、等待构建完成后设置项目 SDK即可在 Course 视图中直接运行Task.main或对应的TaskTest。如果希望脱离 IDE 快速验证也可以直接运行main方法观察日志中的5.5输出。提示本练习属于 katas 系列 Common Transforms / Aggregation 单元建议按Count → Sum → Mean → Min → Max的顺序逐个完成以系统掌握 Beam 聚合类组合器的通用模式。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Java Katas 实战用 Mean 变换计算 PCollection 平均值Apache Beam Java Katas 实战用 Mean 变换计算 PCollection 平均值 本文围绕 Apache Beam 仓库中 learn批处理流处理大数据Apache Beam Go SDK 聚合实战用 stats.Mean 计算 PCollection 均值Mean详解Apache Beam Go SDK 聚合实战用 stats.Mean 计算 PCollection 均值Mean详解 导读 本文围绕 Beam KataApache Beam Python Katas 实战使用 Mean 组合器计算全局均值Aggregation - MeanApache Beam Python Katas 实战使用 Mean 组合器计算全局均值Aggregation Mean 本指南围绕 Apache Bea批处理流处理大数据上一篇open-design 中的 Shopify 设计系统暗色影院风电商设计的 Token 体系与 Agent 落地指南下一篇Jspreadsheet 入门指南从零构建基于 JavaScript 的在线数据表格创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →