news 2026/10/9 1:51:33

Apache Beam Java 中用 Combine.BinaryCombineFn 实现二元归并:BigInteger 求和 Kata 深度解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Java 中用 Combine.BinaryCombineFn 实现二元归并:BigInteger 求和 Kata 深度解析

【免费下载链接】beam

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

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

本文基于 Apache Beam 官方 kata(练习任务)Core Transforms/Combine/BinaryCombineFn,讲解如何用Combine.BinaryCombineFn实现一个可分布、可分治的 BigInteger 求和 combiner。读完你可以掌握:Combine 转换对合并函数的可交换性/可结合性约束及其背后的分布式原理、BinaryCombineFn各方法的职责与源码实现,以及如何用TestPipeline+PAssert验证管道输出。

一、Kata 任务背景:Combine 转换的核心约束

Combine是 Beam 中用于将一组元素(或同一 key 下的一组值)合并为单一结果的核心转换。任务说明文档 task.md 给出了两点关键约束:

  1. 合并函数必须是可交换且可结合的(commutative and associative)。因为输入数据(包括 value 集合)可能被分布到多个 worker 上,合并函数不保证对所有值恰好调用一次,而是会对值集合的多个子集做局部合并(partial combining),之后再合并局部结果;
  2. 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 源码,可以归纳出选型准则:

  1. 选BinaryCombineFn<V>:结果类型与输入相同,且聚合可表达为二元运算(和、积、最大值、集合并等),且运算可交换、可结合。累加器只占一个Holder<V>,无缓冲开销;
  2. 需要中间状态(如计数 + 总和求平均)时:改用AccumulatingCombineFn,显式定义累加器类型;
  3. 运算必须看到全部输入集合时:只能用IterableCombineFn,但要注意它会缓冲值,空间开销大;
  4. 空输入行为:如需明确语义,覆写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.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载
上一篇:鸣潮游戏自动化效率工具:从肝帝到摸鱼党的智能辅助全攻略
下一篇:PingFangSC字体包:中文字体跨平台应用新方案

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

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

C#+Sql Server网上书店管理系统:从数据库设计到事务实现

简介&#xff1a;基于C#与SQL Server开发的网上书店管理系统&#xff0c;是一份面向ASP.NET课程设计的完整项目源码包&#xff0c;适合计算机相关专业学生作为毕业设计或课程设计参考。系统采用浏览器/服务器结构&#xff0c;将前台用户与后台员工分离为两套页面模板&#xff0…

作者头像 李华
网站建设 2026/10/9 1:47:23

Multisim秒表仿真避坑指南:时序精度与物理约束实战解析

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

作者头像 李华
网站建设 2026/10/9 1:46:43

PWM控制从原理到实战:模式1与模式2区别、电机调速及故障保护

1. PWM控制的基本原理与核心价值PWM这三个字母&#xff0c;但凡搞过单片机、玩过电机、调过灯光的&#xff0c;基本都绕不开。全称Pulse Width Modulation&#xff0c;中文叫脉冲宽度调制。说白了&#xff0c;就是用一串方波去“骗”负载&#xff0c;让它以为自己拿到的是模拟量…

作者头像 李华