)
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载导读本篇基于 Apache Beam 仓库中的 Go 练习Kata任务 learning/katas/go/common_transforms/aggregation/sum/task.md讲解如何用 Go SDK 的stats包对 PCollection 执行求和聚合。你将掌握stats.Sum的调用方式、类型约束与内部实现原理并能结合练习配套的main.go与task_test.go完整地运行和验证一个求和 Pipeline。本任务是 Beam Katas 系列中聚合Aggregation课程的延续是理解stats统计变换家族Count、Min、Max、Mean、Sum的关键一环。Kata 任务计算 PCollection 中所有元素之和原练习文档的核心要求只有一句话Kata:Compute the sum of all elements from an input.计算输入中所有元素的总和该任务建立在前一课对stats包位于sdks/go/pkg/beam/transforms/stats的探索之上——stats包的作用是简化对 PCollection 的统计处理。本课将统计能力聚焦到一个具体操作求和。练习给出的提示是使用stats.Sum来计算 PCollection 中元素的总和。仓库中的 task-info.yaml 显示这是一个type: edu的教育型任务其中pkg/task/task.go含有一段待补全的占位符placeholder_text: TODO()学习者需要自行实现求和逻辑而test/task_test.go对学习者不可见visible: false用于隐藏的自动判题。cmd/main.go与pkg/task/task.go对学习者可见是完成练习的主要入口。参考实现封装 stats.Sum 变换任务要求实现的是task包下的ApplyTransform函数。仓库中已经给出了标准解答位于 learning/katas/go/common_transforms/aggregation/sum/pkg/task/task.gopackage task import ( github.com/apache/beam/sdks/v2/go/pkg/beam github.com/apache/beam/sdks/v2/go/pkg/beam/transforms/stats ) func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return stats.Sum(s, input) }这段实现有三个值得注意的要点函数签名规约ApplyTransform接收beam.ScopePipeline 作用域与beam.PCollection输入集合返回求和后的beam.PCollection。Katas 系列统一采用这种输入 PCollection → 输出 PCollection的变换封装风格方便在测试中被直接调用。引入 stats 包github.com/apache/beam/sdks/v2/go/pkg/beam/transforms/stats是 Go SDK 内置的统计变换包无需自行实现累加逻辑。一行核心代码stats.Sum(s, input)即完成对整个 PCollection 的求和返回一个只含单个元素的 PCollectionsingleton。完整可运行的 Pipeline 入口为了让变换真正跑起来练习提供了一个完整的主程序 learning/katas/go/common_transforms/aggregation/sum/cmd/main.gopackage main import ( beam.apache.org/learning/katas/common_transforms/aggregation/sum/pkg/task context github.com/apache/beam/sdks/v2/go/pkg/beam github.com/apache/beam/sdks/v2/go/pkg/beam/log github.com/apache/beam/sdks/v2/go/pkg/beam/x/beamx github.com/apache/beam/sdks/v2/go/pkg/beam/x/debug ) func main() { ctx : context.Background() p, s : beam.NewPipelineWithRoot() input : beam.Create(s, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10) output : task.ApplyTransform(s, input) debug.Print(s, output) err : beamx.Run(ctx, p) if err ! nil { log.Exitf(context.Background(), Failed to execute job: %v, err) } }该入口程序串起了 Beam Go Pipeline 的标准生命周期可作为后续所有统计类 Kata 的通用模板步骤代码说明创建上下文ctx : context.Background()后续用于日志与运行构建 Pipelinebeam.NewPipelineWithRoot()同时拿到*beam.Pipeline与根Scope所有变换都挂在s之下构造输入beam.Create(s, 1, 2, ..., 10)将 1 到 10 共 10 个 int 元素放进一个 PCollection应用变换task.ApplyTransform(s, input)调用待实现的求和变换打印输出debug.Print(s, output)将结果 PCollection 打印到日志执行作业beamx.Run(ctx, p)在可用的 Runner默认 Direct Runner 或环境配置的 Runner上执行 Pipeline在本例中beam.Create生成的输入为1, 2, 3, ..., 10经stats.Sum求和后输出应为55即 1 到 10 的算术和。运行方式在 Go 环境就绪后于仓库learning/katas/go目录下执行go run ./common_transforms/aggregation/sum/cmd执行成功后日志中会通过debug.Print打印出唯一结果元素55。用单元测试验证求和结果练习的隐藏测试 learning/katas/go/common_transforms/aggregation/sum/test/task_test.go 使用 Beam 自带的测试框架对ApplyTransform进行校验func TestApplyTransform(t *testing.T) { p, s : beam.NewPipelineWithRoot() tests : []struct { input beam.PCollection want []interface{} }{ { input: beam.Create(s, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10), want: []interface{}{55}, }, } for _, tt : range tests { got : task.ApplyTransform(s, tt.input) passert.Equals(s, got, tt.want...) if err : ptest.Run(p); err ! nil { t.Error(err) } } }这段测试揭示了两个关键实践passert.Equals来自github.com/apache/beam/sdks/v2/go/pkg/beam/testing/passert在 Pipeline 内声明输出必须等于期望值的断言若stats.Sum的结果不等于[55]测试即失败。ptest.Run来自github.com/apache/beam/sdks/v2/go/pkg/beam/testing/ptest在测试中真正执行该 Pipeline。由于求和是全局聚合stats.Sum的结果是单元素集合因此期望值只需一个元素55。这正体现了 Beam Go 的测试范式先声明断言再运行 Pipeline全部在内存中完成无需外部依赖。在learning/katas/go目录下执行go test ./common_transforms/aggregation/sum/...即可运行该测试。源码探秘stats.Sum 的内部实现要深入理解求和变换直接阅读其实现 sdks/go/pkg/beam/transforms/stats/sum.go// Sum returns the sum of the elements in a PCollectionA as a singleton // PCollectionA. It can only be used for numbers, such as int, uint16, // float32, etc. // // For example: // // col : beam.Create(s, 1, 11, 7, 5, 10) // sum : stats.Sum(s, col) // PCollectionint with 34 as the only element. func Sum(s beam.Scope, col beam.PCollection) beam.PCollection { s s.Scope(stats.Sum) return combine(s, findSumFn, col) } // SumPerKey returns the sum of the values per key in a PCollectionKVA,B as // a PCollectionKVA,B. It can only be used for value numbers, such as int, // uint16, float32, etc. func SumPerKey(s beam.Scope, col beam.PCollection) beam.PCollection { s s.Scope(stats.SumPerKey) return combinePerKey(s, findSumFn, col) }从源码可以确认以下几点返回单元素集合Sum将 PCollection 归约为一个 singleton PCollection元素类型与输入一致。仅支持数值类型如注释所述只能用于int、uint16、float32等数值类型字符串等非数值类型会失败。命名作用域s.Scope(stats.Sum)会在 Pipeline 图中为变换打上可读的名称便于调试与可视化。委托给通用合并原语combine(s, findSumFn, col)底层基于 Beam 的 Combine 机制分布式累加器模型实现findSumFn负责根据元素类型选择对应的累加函数。findSumFn的实现位于生成文件 sdks/go/pkg/beam/transforms/stats/sum_switch.go由sum_switch.tmpl模板配合//go:generate specialize指令生成见 sum.go 顶部的生成注解。它基于反射类型进行分派func findSumFn(t reflect.Type) any { switch t.String() { case int: return sumIntFn case int8: return sumInt8Fn case int16: return sumInt16Fn case int32: return sumInt32Fn case int64: return sumInt64Fn case uint: return sumUintFn case uint8: return sumUint8Fn // ... 以及 uint16/uint32/uint64/float32/float64 等 } }这套模板生成 类型分派的设计意味着Go SDK 为每一种内建数值类型都生成了对应的累加函数如sumIntFn、sumFloat64Fn在运行时根据 PCollection 的实际元素类型自动挑选从而在保持类型安全的同时免除用户手写类型断言。更多证据stats 包自身的单元测试求和行为在 stats 包内部也有充分的测试覆盖见 sdks/go/pkg/beam/transforms/stats/sum_test.goTestSumInt验证[]int{1, -2, 3}求和得2[]int{1, 11, 7, 5, 10}求和得34[]int{0}求和得0——覆盖负数、多元素与单元素场景。TestSumFloat验证[]float32{1, -2, 3.5}求和得2.5[]float32{0, -99.99, 1, 1}求和得-97.99——确认浮点类型同样受支持。TestSumKeyed验证SumPerKey对KV结构按 key 分组求和如alpha的1与-4合并为-3确认Sum家族对键值对数据的扩展能力。这些测试可以作为我们练习任务预期行为的权威参照stats.Sum支持整数与浮点数支持负值与零且对空/单元素输入均有明确定义。扩展按 Key 求和SumPerKey当输入不是普通数值集合而是键值对PCollectionKVA,B时stats包提供了SumPerKey变换它会按 key 分组并对每个 key 的 value 分别求和输出仍是PCollectionKVA,B。这在诸如按部门统计销售额按用户统计点击量等分组聚合场景中非常实用。其使用方式与Sum一致output : stats.SumPerKey(s, kvCollection)SumPerKey与Sum共用同一个findSumFn类型分派逻辑因此同样只支持数值类型的 value并同样基于 Beam Combine 的分布式合并语义适合大规模分布式数据集的归约计算。小结Kata 目标用stats.Sum将输入 PCollection 归约为单个求和结果元素。标准答案ApplyTransform中一行stats.Sum(s, input)参考 task.go。运行入口cmd/main.go构造 110 的输入并打印结果55参考 main.go。验证方式task_test.go通过passert.Equals断言输出为[55]stats 包自身的 sum_test.go 覆盖 int、float32 与 KV 三种场景。实现原理Sum基于 Combine 原语通过findSumFn按元素类型分派到对应累加函数仅支持数值类型返回 singleton PCollection参考 sum.go 与 sum_switch.go。完成本课后你可以继续探索同一stats包下的Min、Max、Mean、Count等聚合变换它们共享相同的 Combine 底层机制与测试范式一通百通。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Go SDK 实战使用 stats.Sum 计算 PCollection 元素总和Kata 详解Apache Beam Go SDK 实战使用 stats.Sum 计算 PCollection 元素总和Kata 详解 本文围绕 Beam Katas大数据批处理流处理数据工程Apache Beam 实战 Kata使用 Sum 聚合变换计算 PCollection 元素总和Apache Beam 实战 Kata使用 Sum 聚合变换计算 PCollection 元素总和 导读 本文以 Apache Beam 仓库中 learni大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战使用 Sum 变换计算元素总和Common Transforms 之 Aggregation/SumApache Beam Kotlin Kata 实战使用 Sum 变换计算元素总和Common Transforms 之 Aggregation/Sum大数据批处理流处理数据工程上一篇NodeGui QPushButton 完全指南在 Node.js 中创建与操控原生按钮控件下一篇Apache Pulsar 租户Tenant管理完全指南pulsar-admin / REST API / Java Admin 三端实战与源码解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考