news 2026/10/6 7:28:01

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Sum 聚合详解:全局求和与按 Key 求和的 Java / Python / Go 三语言实战
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache 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 求和:在KV<Key, Value>集合中,把每个相同 Key 关联的所有 Value 相加,输出KV<Key, 累加和>,例如按商品统计总销量、按用户统计总消费。

在 Tour of Beam 的 common-transforms 模块 中,Sum 属于 Aggregations(聚合)单元,与count、mean、min、max并列(见 aggregation/group-info.yaml),是入门聚合思想的第一站。不同 SDK 对该变换的命名略有差异,但语义完全一致:Java 使用Sum类静态方法,Python 复用CombineGlobally/CombinePerKey,Go 提供stats.Sum/stats.SumPerKey。

二、全局求和:三种 SDK 的写法与输出

2.1 Java:Sum.doublesGlobally()与整数版本

Java SDK 中,全局求和通过Sum类上的xxxGlobally()静态方法完成。原文档示例:

PCollection<Integer> input = pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollection<Double> 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.Globally<Integer, Integer> integersGlobally() { return Combine.globally(Sum.ofIntegers()); } public static Combine.Globally<Double, Double> doublesGlobally() { return Combine.globally(Sum.ofDoubles()); }

Sum 类实际提供三套数值类型的完整组合:

静态方法输入类型输出类型底层 CombineFn
integersGlobally()/integersPerKey()IntegerIntegerofIntegers()
longsGlobally()/longsPerKey()LongLongofLongs()
doublesGlobally()/doublesPerKey()DoubleDoubleofDoubles()

每套都实现了Combine.BinaryCombineXxxFn:apply(a, b)返回a + b,identity()返回0(见 Sum.java)。identity()为 0 意味着:当输入PCollection为空时,全局求和的默认输出是0,而不是空集合或异常——这也是将求和表达为可合并 CombineFn 的好处:分布式执行时任意顺序的合并都能得到相同结果。

2.2 Python:CombineGlobally(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对象或任意 callable,callable 会被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 Go:stats.Sum

Go 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.go:beam.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 Java:Sum.integersPerKey()

当集合元素是KV<String, Integer>时,用integersPerKey()对每个 Key 分别求和:

PCollection<KV<String, Integer>> input = pipeline.apply( Create.of(KV.of("🥕", 3), KV.of("🥕", 2), KV.of("🍆", 1), KV.of("🍅", 4), KV.of("🍅", 5), KV.of("🍅", 3))); PCollection<KV<String, Integer>> sumPerKey = input.apply(Sum.integersPerKey());

输出:

KV{🍆, 1} KV{🍅, 12} KV{🥕, 5}

从 Sum.java 可见integersPerKey()返回Combine.PerKey<K, Integer, Integer>,即Combine.perKey(Sum.ofIntegers()),只是比全局版多了一组 Key 维度;longsPerKey、doublesPerKey同理。底层实际上是在每个 Key 内部执行与全局版相同的二元加法合并。

3.2 Python:CombinePerKey(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 Go:stats.SumPerKey()

func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return stats.SumPerKey(s, input) }

SumPerKey要求输入是KV<A, B>且 B 为非复数数值类型:combinePerKey先通过beam.ValidateKVType拆出值类型再做类型校验(见 util.go),随后交给beam.CombinePerKey执行分组累加。

四、Playground 动手练习:把全局求和改成按 Key 求和

原文档在"Playground exercise"中给出了一个很有教学意义的变形练习:保持整体结构不变,仅更换输入和变换,把"全局求和"升级为"按 Key 求和"。三语言对照如下:

Go:将整数输入替换为 map 风格的 KV 输入,并把stats.Sum换成stats.SumPerKey:

input := 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:将PCollection<Integer>替换为PCollection<KV<Integer, Integer>>,同时把Sum.integersGlobally换成Sum.integersPerKey,并同步修改泛型与applyTransform签名:

PCollection<KV<Integer, 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 PCollection<KV<Integer, Integer>> applyTransform(PCollection<KV<Integer, 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 示例,动手运行一遍是最好的巩固方式。

  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载
上一篇:SiYuan 加密笔记本全解析:从 KEK/DEK 密钥架构到孤岛式隔离的完整实现指南
下一篇:告别卡顿:Unity ECS物理系统入门指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/6 7:23:06

PCB设计规则配置指南:Altium Designer对接嘉立创工艺一次过审

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

作者头像 李华
网站建设 2026/10/6 7:21:45

FPGA直连NVMe SSD:硬件协议栈实现与性能调优实战

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

作者头像 李华
网站建设 2026/10/6 7:21:44

macOS 下用 Luatools 烧录 LuatOS 固件:从驱动配置到串口调试全攻略

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

作者头像 李华
网站建设 2026/10/6 7:21:44

远程IO选型实战指南:EtherCAT/PROFINET/EtherNet/IP深度对比

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

作者头像 李华