ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam Sum 聚合详解:全局求和与按 Key 求和的 Java / Python / Go 三语言实战

Apache Beam Sum 聚合详解:全局求和与按 Key 求和的 Java / Python / Go 三语言实战 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Sum求和是 Apache Beam 最常用的聚合变换之一它既能计算整个PCollection中所有元素的全局总和也能按KV中的 Key 分别累加各自关联的值。本文以 Tour of Beam 学习路径中 Sum 单元说明 为主线结合仓库内 Java、Python、Go 三个 SDK 的源码与可运行示例带你彻底掌握Sum的 API 用法、底层实现原理与 Playground 练习方法读完即可在真实流水线中完成各类数值累加需求。一、Sum 能解决什么问题Sum transform 面向两类典型的聚合诉求全局求和把整个集合中的所有元素相加输出一个单一数值单元素PCollection例如统计销售总金额、日志条数对应的总流量按 Key 求和在KVKey, Value集合中把每个相同 Key 关联的所有 Value 相加输出KVKey, 累加和例如按商品统计总销量、按用户统计总消费。在 Tour of Beam 的 common-transforms 模块 中Sum 属于 Aggregations聚合单元与count、mean、min、max并列见 aggregation/group-info.yaml是入门聚合思想的第一站。不同 SDK 对该变换的命名略有差异但语义完全一致Java 使用Sum类静态方法Python 复用CombineGlobally/CombinePerKeyGo 提供stats.Sum/stats.SumPerKey。二、全局求和三种 SDK 的写法与输出2.1 JavaSum.doublesGlobally()与整数版本Java SDK 中全局求和通过Sum类上的xxxGlobally()静态方法完成。原文档示例PCollectionInteger input pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollectionDouble sum input.apply(Sum.doublesGlobally());输出55对应本仓库中可直接运行的单文件示例见 sum/java-example/Task.java它用Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)构造输入调用Sum.integersGlobally()后通过ParDo打印结果。从源码看这些方法并非独立实现而是Combine系列变换的类型化封装。在 Sum.java 中public static Combine.GloballyInteger, Integer integersGlobally() { return Combine.globally(Sum.ofIntegers()); } public static Combine.GloballyDouble, Double doublesGlobally() { return Combine.globally(Sum.ofDoubles()); }Sum 类实际提供三套数值类型的完整组合静态方法输入类型输出类型底层 CombineFnintegersGlobally()/integersPerKey()IntegerIntegerofIntegers()longsGlobally()/longsPerKey()LongLongofLongs()doublesGlobally()/doublesPerKey()DoubleDoubleofDoubles()每套都实现了Combine.BinaryCombineXxxFnapply(a, b)返回a bidentity()返回0见 Sum.java。identity()为 0 意味着当输入PCollection为空时全局求和的默认输出是0而不是空集合或异常——这也是将求和表达为可合并 CombineFn 的好处分布式执行时任意顺序的合并都能得到相同结果。2.2 PythonCombineGlobally(sum)Python SDK 没有单独的Sum类而是把内置函数sum直接传给CombineGlobally即可完成全局求和import apache_beam as beam with beam.Pipeline() as p: total ( p | Create numbers beam.Create([3, 4, 1, 2]) | Sum values beam.CombineGlobally(sum) | beam.Map(print))输出10为什么 Python 里传一个普通函数就行在 core.py 中CombineGlobally.__init__接受CombineFn对象或任意 callablecallable 会被CombineFn.maybe_from_callable自动包装内置sum(iterable)恰好满足 CombineFn 的累加 合并语义。CombineGlobally在expand内部会先执行补空 Key →CombinePerKey→ 去 Key三步KeyWithVoid ParDo(...) | CombinePerKey combine_per_key | UnKey Map(...)因此全局聚合在底层复用按 Key 聚合的执行路径。值得留意的是CombineGlobally的默认值语义默认has_defaults True空输入会输出 CombineFn 的默认结果对sum即0如果输入是使用非默认窗口非GlobalWindows的无界集合源码会在运行时提示改用.without_defaults()输出空集合或.as_singleton_view()作为单例侧输入使用见 core.py。此外还提供with_fanout(n)为热 Key 场景做扇出优化。这些方法在流式、空窗口等边界场景中非常实用。2.3 Gostats.SumGo SDK 将求和放在beam/transforms/stats子包中import ( github.com/apache/beam/sdks/go/pkg/beam github.com/apache/beam/sdks/go/pkg/beam/transforms/stats ) func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return stats.Sum(s, input) }可运行的完整示例见 sum/go-example/main.gobeam.Create(s, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)构造输入stats.Sum得到单元素集合再用debug.Printf输出。stats.Sum的实现sum.go最终落到combine(s, findSumFn, col)而combine在流水线构建期会根据元素的反射类型做类型分派并调用validateNonComplexNumber做校验——只允许 int、uint16、float32 等非复数数值类型否则直接 panic见 util.go。也就是说Go 版 Sum 的类型安全性在管道构建阶段就被静态保证了。三、按 Key 求和分组累加实战3.1 JavaSum.integersPerKey()当集合元素是KVString, Integer时用integersPerKey()对每个 Key 分别求和PCollectionKVString, Integer input pipeline.apply( Create.of(KV.of(, 3), KV.of(, 2), KV.of(, 1), KV.of(, 4), KV.of(, 5), KV.of(, 3))); PCollectionKVString, Integer sumPerKey input.apply(Sum.integersPerKey());输出KV{, 1} KV{, 12} KV{, 5}从 Sum.java 可见integersPerKey()返回Combine.PerKeyK, Integer, Integer即Combine.perKey(Sum.ofIntegers())只是比全局版多了一组 Key 维度longsPerKey、doublesPerKey同理。底层实际上是在每个 Key 内部执行与全局版相同的二元加法合并。3.2 PythonCombinePerKey(sum)import apache_beam as beam with beam.Pipeline() as p: totals_per_key ( p | Create produce beam.Create([ (, 3), (, 2), (, 1), (, 4), (, 5), (, 3),]) | Sum values per key beam.CombinePerKey(sum) | beam.Map(print))输出(, 5) (, 1) (, 12)注意 Python 中CombineGlobally(sum)的expand正是内部借用CombinePerKey实现先对全部元素补上同一个NoneKey所以两者的累加逻辑完全一致区别只在是否保留 Key 维度。3.3 Gostats.SumPerKey()func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return stats.SumPerKey(s, input) }SumPerKey要求输入是KVA, B且 B 为非复数数值类型combinePerKey先通过beam.ValidateKVType拆出值类型再做类型校验见 util.go随后交给beam.CombinePerKey执行分组累加。四、Playground 动手练习把全局求和改成按 Key 求和原文档在Playground exercise中给出了一个很有教学意义的变形练习保持整体结构不变仅更换输入和变换把全局求和升级为按 Key 求和。三语言对照如下Go将整数输入替换为 map 风格的 KV 输入并把stats.Sum换成stats.SumPerKeyinput : beam.ParDo(s, func(_ []byte, emit func(int, int)) { emit(1, 1) emit(1, 4) emit(2, 6) emit(2, 3) emit(2, -4) emit(3, 23) }, beam.Impulse(s))Java将PCollectionInteger替换为PCollectionKVInteger, Integer同时把Sum.integersGlobally换成Sum.integersPerKey并同步修改泛型与applyTransform签名PCollectionKVInteger, Integer input pipeline.apply( Create.of(KV.of(1, 11), KV.of(1, 36), KV.of(2, 91), KV.of(3, 33), KV.of(3, 11), KV.of(4, 33))); static PCollectionKVInteger, Integer applyTransform(PCollectionKVInteger, Integer input) { return input.apply(Sum.integersPerKey()); }Python把beam.CombineGlobally(sum)换成beam.CombinePerKey(sum)输入改为(key, value)元组列表p | beam.Create([(1, 36),(2, 91),(3, 33),(3, 11),(4, 67),]) | beam.CombinePerKey(sum)练习完成后可以看到每个 Key 一行输出同一 Key 的多个值被合并为一个累加和。这一步直观演示了同样的 Combine 函数 不同的聚合形态这一 Beam 核心思想——全局聚合与分组聚合共享同一套累加逻辑。五、思考题为什么输出的顺序看起来不稳定原文档末尾抛出一个值得深入思考的问题控制台打印出的集合元素顺序并不固定多次运行可能不同为什么这与 Beam 的执行模型直接相关。Sum/CombinePerKey在分布式运行器上会经历shuffle / 分组grouping阶段元素按 Key 被打散到不同工作节点再在每台机器上合并最终结果的汇聚顺序取决于调度、并行度与网络时序而不是输入源中的原始顺序。因此全局求和输出只有一个元素无所谓顺序按 Key 求和时不同 Key 之间的输出顺序是不保证的在流式场景下输出还会受窗口触发trigger策略影响例如 学习路径中的 windowing 与 triggers 章节 专门讨论这类时序问题。如果业务上对输出顺序有要求例如按 Key 排序后输出通常的做法是在 Sum 之后再叠加一次排序变换而不是依赖运行器的天然顺序。六、延伸学习从 Sum 到完整聚合家族Sum 只是 Tour of Beam 聚合单元的第一个成员。在 aggregation 目录 下与 Sum 并列的还有count/description.md统计元素个数mean/description.md计算均值min/description.md 与 max/description.md求最值。它们全部基于 Combine 语义学透 Sum 的全局 vs 按 Key两种形态即可举一反三。完成本单元后还可以回到 common-transforms 模块首页 继续学习 Filter、WithKeys 等常用变换或进入 tour-of-beam 学习内容总览 探索更多主题全部章节都配有可运行的 Playground 示例动手运行一遍是最好的巩固方式。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战使用 Combine.perKey 实现按 Key 聚合求和Apache Beam Java Kata 实战使用 Combine.perKey 实现按 Key 聚合求和 在 Apache Beam 中对按 Key 分大数据批处理流处理数据工程Apache Beam Go SDK 实战用 CombinePerKey 实现按 Key 分组聚合求和Apache Beam Go SDK 实战用 CombinePerKey 实现按 Key 分组聚合求和 导读 本文围绕 Apache Beam Go SDK大数据批处理流处理数据工程Apache Beam Java Katas 实战用 Lambda 与 BinaryCombineFn 实现 BigInteger 全局求和Apache Beam Java Katas 实战用 Lambda 与 BinaryCombineFn 实现 BigInteger 全局求和 本文围绕 Apa大数据批处理流处理数据工程上一篇SiYuan 加密笔记本全解析从 KEK/DEK 密钥架构到孤岛式隔离的完整实现指南下一篇告别卡顿Unity ECS物理系统入门指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进