【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文基于 Apache Beam 官方 kata(练习任务)Core Transforms/Combine/BinaryCombineFn,讲解如何用Combine.BinaryCombineFn实现一个可分布、可分治的 BigInteger 求和 combiner。读完你可以掌握:Combine 转换对合并函数的可交换性/可结合性约束及其背后的分布式原理、BinaryCombineFn各方法的职责与源码实现,以及如何用TestPipeline+PAssert验证管道输出。
一、Kata 任务背景:Combine 转换的核心约束
Combine是 Beam 中用于将一组元素(或同一 key 下的一组值)合并为单一结果的核心转换。任务说明文档 task.md 给出了两点关键约束:
- 合并函数必须是可交换且可结合的(commutative and associative)。因为输入数据(包括 value 集合)可能被分布到多个 worker 上,合并函数不保证对所有值恰好调用一次,而是会对值集合的多个子集做局部合并(partial combining),之后再合并局部结果;
BinaryCombineFn专门用于实现"天然可以表达为二元运算"的 combiner,即结果可以直接由两个同类型值运算得到(如求和、求最大/最小值)。
本 kata 的具体要求是:用Combine.BinaryCombineFn实现BigInteger的求和。配套代码文件为 Task.java,元信息见 task-info.yaml(type: edu,在applyTransform方法中留有TODO()占位符待学习者填写)。
二、任务代码:从输入到输出的完整管道
kata 的Task类结构非常简单,分三部分:构造输入、应用 Combine 转换、输出日志。仓库中 Task.java 已给出完整实现(对应 kata 中填写完 TODO 的最终形态):
public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline = Pipeline.create(options); PCollection<BigInteger> numbers = pipeline.apply( Create.of( BigInteger.valueOf(10), BigInteger.valueOf(20), BigInteger.valueOf(30), BigInteger.valueOf(40), BigInteger.valueOf(50) )); PCollection<BigInteger> output = applyTransform(numbers); output.apply(Log.ofElements()); // 打印结果,Log 见 kata util 包 pipeline.run(); } static PCollection<BigInteger> applyTransform(PCollection<BigInteger> input) { return input.apply(Combine.globally(new SumBigIntegerFn())); } static class SumBigIntegerFn extends BinaryCombineFn<BigInteger> { @Override public BigInteger apply(BigInteger left, BigInteger right) { return left.add(right); } }几个值得注意的要点:
- 输入:
Create.of(...)产生 5 个BigInteger(10、20、30、40、50),期望全局求和结果为150; Combine.globally(fn):对整个PCollection做全局归并,产出单元素PCollection。与它对应的是Combine.perKey(fn),用于KV集合上按 key 归并;- combiner 实现:
SumBigIntegerFn继承BinaryCombineFn<BigInteger>,只需覆写apply(left, right)一个方法。加法天然满足可交换性与可结合性,因此分布式局部合并是安全的; - 输出日志:
Log.ofElements()来自 kata 的工具包 Log.java,仅用于在直接运行时打印结果。
三、源码剖析:BinaryCombineFn 如何把二元运算映射为四步合并协议
BinaryCombineFn定义在 Combine.java 中,其类声明揭示了它与通用CombineFn的映射关系:
public abstract static class BinaryCombineFn<V> extends CombineFn<V, Holder<V>, V> {- 输入类型与输出类型都是
V(归并不改变元素类型); - 累加器类型是
Holder<V>——即"单值 + 是否存在标志(present)"的轻量容器。这正是BinaryCombineFn的精髓:它避免了IterableCombineFn那样把全部值缓冲进列表,内存占用只与单个累加值相关。
从 Combine.java 源码 看,BinaryCombineFn替用户自动实现了CombineFn协议的全部方法:
| 方法 | 行为(源码语义) |
|---|---|
abstract V apply(V left, V right) | 唯一必须覆写的方法:对两个操作数执行二元运算 |
V identity()(可覆写,默认null) | 空输入集合时的返回值。求和场景下覆写为BigInteger.ZERO更严谨 |
createAccumulator() | 返回新的Holder<V> |
addInput(Holder<V>, V) | 若累加器已有值则apply(accumulator.value, input),否则直接把input设为累加值 |
mergeAccumulators(Iterable<Holder<V>>) | 逐个遍历局部累加器,跳过不存在(!present)者,用apply串联归并;全部不存在时返回空累加器 |
extractOutput(Holder<V>) | 累加器有值则返回该值,否则返回identity() |
getAccumulatorCoder(...) | 使用HolderCoder<V>对累加器编解码 |
getDefaultOutputCoder(...) | 输出 coder 直接复用输入 coder |
这套"create → addInput → mergeAccumulators → extractOutput"的协议,正是第一节提到的"对子集做局部合并、再树状合并局部结果"在代码上的落地:每个 worker 先用addInput折叠出自己负责的数据子集,Runner 再在需要的层调用mergeAccumulators折叠多个局部结果。apply方法本身因此可以被调用多次、且操作数顺序不固定——这就是为什么它必须是可交换、可结合的。
此外,BinaryCombineFn还提供了一个静态工厂方法(Combine.java 第 507-514 行),可以在不写子类的情况下用 lambda 直接构造 combiner:
// 等价的函数式写法 Combine.globally(BinaryCombineFn.of((BigInteger a, BigInteger b) -> a.add(b)));同文件中还有面向基本类型优化空间的特化版本,如BinaryCombineIntegerFn、BinaryCombineLongFn、BinaryCombineDoubleFn(累加器为原生int[]/long[]/double[]),以及更通用的IterableCombineFn与AccumulatingCombineFn。对于像本 kata 这样"结果即输入类型、运算为二元运算"的场景,BinaryCombineFn是最省代码、也最省内存的选择。
四、内置转换印证:Min / Max 都是 BinaryCombineFn 的产物
BinaryCombineFn并非仅为 kata 服务——SDK 自带的常用转换就是它的直接消费者。例如 Min.java 中:
public static Combine.Globally<Integer, Integer> integersGlobally() { return Combine.globally(new MinIntegerFn()); }其中MinIntegerFn就是一个继承BinaryCombineFn的类。类似的,Max、Distinct、View等核心转换也基于BinaryCombineFn构建(可对照 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/ 目录下的Max.java、Distinct.java、View.java)。这说明"求和/求极值类聚合 = 二元运算 combiner"是 Beam SDK 内部的通用实现模式,kata 让你手写一次,本质上就是在复刻 SDK 的这套做法。
五、验证方式:TestPipeline 与 PAssert 断言
kata 的验收测试为 TaskTest.java,它不依赖main方法,而是通过@Rule注入TestPipeline,只调用静态方法Task.applyTransform(numbers),最后断言全局输出恰好是 150:
@Rule public final transient TestPipeline testPipeline = TestPipeline.create(); @Test public void combine_binaryCombineFn() { Create.Values<BigInteger> values = Create.of( BigInteger.valueOf(10), BigInteger.valueOf(20), BigInteger.valueOf(30), BigInteger.valueOf(40), BigInteger.valueOf(50)); PCollection<BigInteger> numbers = testPipeline.apply(values); PCollection<BigInteger> results = Task.applyTransform(numbers); PAssert.that(results) .containsInAnyOrder(BigInteger.valueOf(150)); testPipeline.run().waitUntilFinish(); }这种"把转换逻辑抽成静态方法applyTransform、测试只针对该方法"的组织方式是 kata 的通用模式:既保证测试可复用任意 Runner(TestPipeline默认使用 DirectRunner 语义的验证),也让main入口可以独立运行观察日志输出。
六、如何运行这个 Kata
kata 工程自带独立的 Gradle 封装,位于 learning/katas/java/ 目录下(含gradlew、settings.gradle、course-info.yaml),与仓库主构建体系相互隔离。在该目录下用自带 wrapper 执行构建与测试即可(例如./gradlew test一类命令),无需改动仓库内任何文件。文件布局如下:
learning/katas/java/Core Transforms/Combine/BinaryCombineFn/ ├── task.md # 任务说明 ├── task-info.yaml # kata 元信息(占位符位置等) ├── src/org/apache/beam/learning/katas/coretransforms/combine/binarycombinefn/Task.java └── test/org/apache/beam/learning/katas/coretransforms/combine/binarycombinefn/TaskTest.java七、小结:什么时候该选 BinaryCombineFn
结合本 kata 与 SDK 源码,可以归纳出选型准则:
- 选
BinaryCombineFn<V>:结果类型与输入相同,且聚合可表达为二元运算(和、积、最大值、集合并等),且运算可交换、可结合。累加器只占一个Holder<V>,无缓冲开销; - 需要中间状态(如计数 + 总和求平均)时:改用
AccumulatingCombineFn,显式定义累加器类型; - 运算必须看到全部输入集合时:只能用
IterableCombineFn,但要注意它会缓冲值,空间开销大; - 空输入行为:如需明确语义,覆写
identity()返回该二元运算的单位元(本 kata 中即BigInteger.ZERO),否则空集合时输出为默认null。
Combine家族的设计哲学——把"分布式局部合并 + 树状归并"的复杂性封装进CombineFn协议,让业务代码只声明一个纯二元运算——在 Combine.java 与 Min.java 等实现中都有清晰体现,也是这个 kata 值得反复体会的原因。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Java Kata 实战:使用 Combine.BinaryCombineFn 实现 BigInteger 求和
Apache Beam Java Kata 实战:使用 Combine.BinaryCombineFn 实现 BigInteger 求和 本篇技术指南以 Apa
大数据批处理流处理数据工程whisper.cpp Vulkan 后端:三步部署 GPU 加速语音识别
whisper.cpp Vulkan 后端:三步部署 GPU 加速语音识别 whisper.cpp 的 CPU 推理在长音频上偏慢,多设备部署还需逐家适配私有接
Apache Beam Java Katas 实战:用 Lambda 与 BinaryCombineFn 实现 BigInteger 全局求和
Apache Beam Java Katas 实战:用 Lambda 与 BinaryCombineFn 实现 BigInteger 全局求和 本文围绕 Apa
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考