- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
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() | Integer | Integer | ofIntegers() |
longsGlobally()/longsPerKey() | Long | Long | ofLongs() |
doublesGlobally()/doublesPerKey() | Double | Double | ofDoubles() |
每套都实现了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.
相关推荐
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
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考