- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Combine是 Apache Beam 中用于将PCollection中的一组元素或数值合并为单一结果的变换(transform),它既有作用于整个PCollection的全局变体(如CombineGlobally/Combine.globally/beam.Combine),也有针对key/value键值对PCollection的按 key 变体(CombinePerKey/Combine.perKey)。本文以 simple-function/description.md 为主线,讲解在 Java、Python、Go 三种 SDK 下如何用"简单函数"完成求和等基础聚合,并结合仓库源码剖析其底层机制与可替换的内置组合函数。
Combine 是什么
Combine是 Beam 中用于"把数据里的元素或数值合并起来"的变换。它的适用场景非常广泛:
- 计算一个
PCollection中所有整数的总和、最小值、最大值; - 把一组字符串拼接成一个字符串;
- 对
PCollection中的key/value键值对按 key 分组,再把同一 key 下的所有 value 合并(此时使用CombinePerKey变体)。
当你应用一个Combine变换时,必须提供一个包含合并逻辑的函数。Beam SDK 同时提供了若干预置的合并函数(pre-built combine functions),覆盖 sum、min、max 等常见数值运算,开箱即用,无需自己手写。
合并函数必须满足的两个性质
文档中特别强调了一个关键约束:合并函数必须是可交换(commutative)且可结合(associative)的。原因在于:
- 该函数不保证对某个 key 的所有 value 恰好只调用一次;
- 输入数据(包括 value 集合)可能被分布到多个 worker 上,因此合并函数可能被多次调用,以对 value 集合的各个子集执行部分合并(partial combining),最后再把部分结果汇总。
正是这两个数学性质,保证了无论数据在分布式环境下如何切分、合并顺序如何变化,最终结果都一致。如果合并逻辑依赖元素顺序(例如"先出现的元素优先"),就无法安全地用于分布式场景。
简单函数:求和——三种 SDK 的最小实现
对于sum这类简单合并操作,通常可以用一个**简单函数(simple function)**直接实现,无需定义复杂的累加器类。下面分别给出 Java、Python、Go 三种 SDK 的完整写法。
Go:beam.Combine+ 内联函数
Go SDK 中可以直接传入一个内联的二元函数,函数签名形如func(sum, elem int) int,第一个参数是当前累计值,第二个参数是待合并的元素:
func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.Combine(s, func(sum, elem int) int { return sum + elem }, input) }在仓库对应的可运行示例 main.go 中,完整流程是:
input := beam.Create(s, 10, 30, 50, 70, 90) output := applyTransform(s, input) debug.Print(s, output) err := beamx.Run(ctx, p)输入是10, 30, 50, 70, 90,beam.Combine将其合并为总和250,并通过debug.Print输出结果。底层 API 定义在 sdks/go/pkg/beam/combine.go:beam.Combine(s, combinefn, col, opts...)对整列元素做全局合并,beam.CombinePerKey(s, combinefn, col, opts...)则按键合并,二者均以任意函数作为合并逻辑。
Java:Combine.globally+SerializableFunction
Java SDK 中,简单函数需要实现SerializableFunction<Iterable<Integer>, Integer>接口:输入是元素的迭代集合,输出是合并后的单一结果。文档中的SumInts即为标准写法:
// Sum a collection of Integer values. The function SumInts implements the interface SerializableFunction. public static class SumInts implements SerializableFunction<Iterable<Integer>, Integer> { @Override public Integer apply(Iterable<Integer> input) { int sum = 0; for (int item : input) { sum += item; } return sum; } }使用方式(取自 Task.java):
PCollection<Integer> input = pipeline.apply(Create.of(10, 30, 50, 70, 90)); PCollection<Integer> output = input.apply(Combine.globally(new SumIntegerFn()));Combine.globally是 Java 侧全局合并的入口。从 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Combine.java 可以看到它提供了多种重载:接收SerializableFunction<Iterable<V>, V>、SerializableBiFunction<V, V, V>(二元函数,即 BinaryCombineFn 风格)以及带CombineFn的通用版本;Combine.perKey则对应按 key 合并(同文件 L165-L200)。
Python:beam.CombineGlobally+ 普通函数
Python SDK 最为灵活,直接传入普通 Python 函数即可,并且支持带默认参数的函数。文档示例通过bound参数演示了这一点:
input = [1, 10, 100, 1000] def bounded_sum(values, bound=500): return min(sum(values), bound) small_sum = input | beam.CombineGlobally(bounded_sum) # [500] large_sum = input | beam.CombineGlobally(bounded_sum, bound=5000) # [1111]small_sum使用默认bound=500,1+10+100+1000=1111被截断为500,因此输出[500];large_sum传入bound=5000,总和1111未超限,原样输出[1111]。
beam.CombineGlobally是 Python SDK 的全局合并变换,定义于 sdks/python/apache_beam/transforms/core.py。它支持丰富的调用方式:CombineGlobally(sum)直接使用函数、CombineGlobally(sum).as_singleton_view()把结果作为单例侧输入(side input)使用、CombineGlobally(sum).without_defaults()在输入为空时不输出默认值(避免空 PCollection 被填充默认结果,参考同文件 L3020-L3033 的相关说明)。
仓库中对应的可运行示例 task.py 采用同样模式:
with beam.Pipeline() as p: (p | beam.Create([1, 2, 3, 4, 5]) | beam.CombineGlobally(sum) | Output())不止数字:Combine 同样适用于字符串等任意类型
文档强调:Combine的输入数据由整数构成,但它也可以与其他类型结合使用,例如strings及其他类型。下面的练习示例把 8 个英文单词合并成一个逗号分隔的字符串。
Go 版本
input := beam.Create(s, "quick", "brown", "fox", "jumps", "over", "the", "lazy", "dog") func applyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.Combine(s, func(sum, elem string) string { return sum+","+elem }, input) }Java 版本
同样只需把SerializableFunction的泛型换成String,在累加循环中用StringBuilder拼接:
public static class ConcatenateStrings implements SerializableFunction<Iterable<String>, String> { @Override public String apply(Iterable<String> input) { StringBuilder concatenated = new StringBuilder(); for (String item : input) { concatenated.append(",").append(item); } return concatenated.toString(); } } PCollection<String> input = pipeline.apply(Create.of("quick", "brown", "fox", "jumps", "over", "the", "lazy", "dog")); PCollection<String> concatenated = input.apply(Combine.globally(new ConcatenateStrings()));Python 版本
def concat(strings): total = "" for word in strings: total += word return total with beam.Pipeline() as p: (p | beam.Create(["quick", "brown", "fox", "jumps", "over", "the", "lazy", "dog"]) | beam.CombineGlobally(concat) | LogElements())从这些例子可以看出:简单函数模式的本质是"接收整个元素的集合,返回一个合并结果"。只要该逻辑满足可交换、可结合,Beam 就能在分布式执行中安全地分批合并,最终结果与串行执行完全一致。
从简单函数到通用 CombineFn
简单函数适合sum、字符串拼接这类"输入类型与输出类型相同"的场景。一旦合并逻辑变复杂(例如计算平均值,输出类型是浮点数、中间态是"总和 + 计数"),就需要创建CombineFn子类,通过四个方法定义完整的合并协议(accumulation type 可以不同于输入/输出类型):
- Create Accumulator:创建新的"局部"累加器。求平均值的例子中,局部累加器记录已累加值的总和(最终除法中的分子)与已累加值的个数(分母),可在分布式场景下被任意多次调用;
- Add Input:把单个输入元素加入累加器并返回新的累加器(示例中更新总和并自增计数),同样可以并行调用;
- Merge Accumulators:把多个累加器合并为一个(平均值场景中合并各部分的分子与分母),其输出还可能被再次合并任意多次;
- Extract Output:执行最终计算(平均值即总和除以个数),只在最终合并后的累加器上调用一次。
关于这一主题的完整三 SDK 实现,可继续阅读同组的 combine-fn/description.md;对"每个 key 下的 value 合并"以及基于二元函数的BinaryCombineFn写法,可分别参考 combine-per-key/description.md 与 binary-combine-fn/description.md。四个单元共同构成 Combine 完整课程,见 group-info.yaml。
预置合并函数与 Runner 优化
对于 sum、min、max 等常见数值操作,不必每次手写函数:Beam SDK 提供了一系列预置合并函数(例如Sum、Min、Max家族,Java 的Combine.globally(Sum.integers())、Python 的beam.CombineGlobally(beam.combiners.Sum.IntsFn())、Go 的combine.SumInts()等),它们天然满足可交换、可结合约束,语义清晰且类型安全。
之所以强调可交换、可结合,还因为这两条性质允许 Runner 自动应用关键优化(详见 combine-fn/description.md):
- Combiner lifting(最显著的优化):在数据被 shuffle 之前,先在每个 key、每个窗口内完成部分合并,从而把需要 shuffle 的数据量减少数个数量级,也称"mapper-side combine";
- Incremental combining(增量合并):当
CombineFn能显著缩减数据大小时,在流式 shuffle 过程中边产出边合并,把合并开销摊薄到计算空闲时段,同时降低中间累加器的存储占用。
理解这些机制有助于你判断:手写简单函数时务必遵守交换律与结合律,否则 Runner 的任何并行/部分合并优化都可能改变结果语义。
小结与延伸阅读
简单函数是入门 Combine 的最佳切入点:
- 三 SDK 写法:Go 用
beam.Combine(s, fn, input)内联函数;Java 用Combine.globally(new SerializableFunction...);Python 用input | beam.CombineGlobally(fn),还支持带默认参数的函数; - 约束:合并函数必须可交换、可结合,因为它在分布式环境下可能被多次调用做部分合并;
- 通用性:不仅支持整数求和,字符串拼接等任意类型的合并同样适用;
- 进阶:复杂合并升级为
CombineFn四方法(createAccumulator / addInput / mergeAccumulators / extractOutput),按 key 合并使用CombinePerKey,二元合并可基于BinaryCombineFn; - 内置能力:优先考虑 SDK 预置的 sum / min / max 合并函数,并在理解 combiner lifting 与增量合并等 Runner 优化后合理设计自己的合并逻辑。
本文对应的可运行示例位于 simple-function 目录(含 Go 示例、Java 示例、Python 示例),欢迎在 Beam Playground 中运行并尝试修改输入数据与合并逻辑,验证 Combine 的分布式合并语义。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 核心变换实战:用 CombinePerKey 按用户名聚合游戏得分(Tour of Beam 挑战二)
Apache Beam 核心变换实战:用 CombinePerKey 按用户名聚合游戏得分(Tour of Beam 挑战二) 本篇技术指南围绕 Apache
大数据批处理流处理数据工程Apache Beam Go SDK Kata 实战:用 Combine 简单函数实现求和
Apache Beam Go SDK Kata 实战:用 Combine 简单函数实现求和 Combine 是 Apache Beam 中用于把集合中的元素或值
Apache Beam Kotlin Katas 实战:Aggregation 之 Count 聚合变换详解
Apache Beam Kotlin Katas 实战:Aggregation 之 Count 聚合变换详解 导读 本文以 Apache Beam 官方 Kot
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考