☰
Apache Beam 分支(Branching)实战:让同一个 PCollection 被多个 Transform 复用
2026/10/9 1:22:35 网站建设 项目流程

【免费下载链接】beam

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

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

本篇文章以 Apache Beam 官方 Katas 训练集中 "Branching" 一课为核心,讲解 Beam 中分支(Branching)的核心思想:同一个PCollection可以同时作为多个 Transform 的输入,而不会消耗或改变原始数据。你将通过一个"数字分别乘以 5 和乘以 10"的完整 Kata 练习,掌握分支式管道(Pipeline)的编写方式、参考实现、测试验证方法,并从源码层面理解分支背后的不可变数据结构与惰性执行原理。

什么是 Branching(分支)

在 Apache Beam 的编程模型中,PCollection是不可变的分布式数据集。PTransform作用于PCollection并产生新的PCollection,这个过程不会修改输入数据本身。

Branching(分支)指的就是:同一个PCollection可以作为多个 Transform 的输入,每个 Transform 各自独立地消费这份数据,互不影响,也不会消耗或改变原始数据。这是 Beam 官方设计指南 "Design Your Pipeline" 中明确提出的模式之一——"Multiple transforms process the same PCollection"。

从数据流图(DAG)的角度看,一个PCollection节点可以拥有多个出边,指向多个不同的 Transform。由于 Beam 的PCollection是不可变的,所有分支读取到的都是同一份数据内容,天然具备并行执行的可能。

本 Kata 位于 learning/katas/java/Core Transforms/Branching/Branching/task.md,隶属于 Java 版 Katas 的 "Core Transforms" 章节,与之并列的还有 Map、GroupByKey、Combine、Flatten、Partition、CoGroupByKey、Side Input、Side Output 等核心变换练习(见 learning/katas/java/Core Transforms/section-info.yaml 所在目录)。

Kata 任务解读

原始任务(task.md)的要求非常明确:

你可以将同一个PCollection用作多个 Transform 的输入,而不会消耗或改变它。

Kata:将数字分支到两个不同的 Transform:一个 Transform 将每个数字乘以 5,另一个 Transform 将每个数字乘以 10。

任务本身的语义分两层:

  1. 概念层:理解"同一个 PCollection 可被多个 Transform 复用"这一分支模式;
  2. 实践层:基于给定的numbers(数字 1~5)PCollection,写出两个独立的 Transform,分别产出乘以 5 和乘以 10 的结果 PCollection。

任务的结构信息记录在 learning/katas/java/Core Transforms/Branching/Branching/task-info.yaml 中:这是一个type: edu的教学型任务,Task.java中预留了两个TODO()占位符(placeholder),对应"乘以 5"与"乘以 10"两个需要学习者补全的 Transform;测试文件TaskTest.java对学习者不可见(visible: false),用于隐藏校验逻辑。任务还被标注为complexity: BASIC(基础难度),并带有branch、transforms、numbers等标签,便于在 Playground / Katas 平台检索定位。

参考实现:分支式管道完整代码

Katas 的参考解法位于 learning/katas/java/Core Transforms/Branching/Branching/src/org/apache/beam/learning/katas/coretransforms/branching/Task.java,完整实现如下:

import static org.apache.beam.sdk.values.TypeDescriptors.integers; import org.apache.beam.learning.katas.util.Log; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.transforms.MapElements; import org.apache.beam.sdk.values.PCollection; public class Task { public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline = Pipeline.create(options); PCollection<Integer> numbers = pipeline.apply(Create.of(1, 2, 3, 4, 5)); PCollection<Integer> mult5Results = applyMultiply5Transform(numbers); PCollection<Integer> mult10Results = applyMultiply10Transform(numbers); mult5Results.apply("Log multiply 5", Log.ofElements("Multiplied by 5: ")); mult10Results.apply("Log multiply 10", Log.ofElements("Multiplied by 10: ")); pipeline.run(); } static PCollection<Integer> applyMultiply5Transform(PCollection<Integer> input) { return input.apply("Multiply by 5", MapElements.into(integers()).via(num -> num * 5)); } static PCollection<Integer> applyMultiply10Transform(PCollection<Integer> input) { return input.apply("Multiply by 10", MapElements.into(integers()).via(num -> num * 10)); } }

逐段拆解这个参考实现的要点:

① 创建数据源

PCollection<Integer> numbers = pipeline.apply(Create.of(1, 2, 3, 4, 5));

使用Create.of(1, 2, 3, 4, 5)从内存中构造一个包含整数 1~5 的PCollection<Integer>,作为整个管道的起点。

② 关键的分支操作——同一个输入被两个 Transform 消费

PCollection<Integer> mult5Results = applyMultiply5Transform(numbers); PCollection<Integer> mult10Results = applyMultiply10Transform(numbers);

注意:numbers这个PCollection被连续传给了两个不同的 Transform。这正是分支(Branching)模式的核心体现——numbers没有被"消耗掉",也没有被任何一方修改,两个 Transform 各自拿到的是同一份原始数据的视图,分别产出属于自己的新PCollection。在分布式执行时,这两个分支可以并行计算。

③ 两个 Transform 的实现

static PCollection<Integer> applyMultiply5Transform(PCollection<Integer> input) { return input.apply("Multiply by 5", MapElements.into(integers()).via(num -> num * 5)); } static PCollection<Integer> applyMultiply10Transform(PCollection<Integer> input) { return input.apply("Multiply by 10", MapElements.into(integers()).via(num -> num * 10)); }

两个 Transform 都基于MapElements实现:MapElements.into(integers())声明输出的元素类型为Integer(借助TypeDescriptors.integers()指定类型描述符),.via(num -> num * 5)定义逐元素的映射函数。两个 Transform 之间唯一的区别是映射函数不同(×5 与 ×10),但它们都从同一个inputPCollection 读取数据。这正是任务标题Branching想要强调的结构:把同一份数据"兵分两路"。

把 Transform 封装成独立方法(applyMultiply5Transform/applyMultiply10Transform)也是值得借鉴的习惯——既便于阅读,也让测试可以针对性地复用(下文测试就是直接调用这两个方法)。

④ 输出观察(Log)

mult5Results.apply("Log multiply 5", Log.ofElements("Multiplied by 5: ")); mult10Results.apply("Log multiply 10", Log.ofElements("Multiplied by 10: "));

两个分支的结果分别被挂上带前缀的日志 Transform,用于在运行终端中观察输出。这个Log.ofElements(String prefix)是 Katas 自带的工具类,实现位于 learning/katas/java/util/src/org/apache/beam/learning/katas/util/Log.java,其内部是一个基于ParDo+DoFn的自定义PTransform:

  • 通过@ProcessElement注解的processElement方法处理每个元素;
  • 使用 SLF4J 的LOG.info(message)输出日志,message由prefix + element.toString()拼接而成;
  • 如果当前窗口不是GlobalWindow,还会在消息末尾追加Window:<窗口信息>,方便观察窗口化数据;
  • 最后通过out.output(element)将元素原样透传,因此Log是一个"只观测、不改数据"的 Transform,非常适合在练习中打印中间结果。

预期终端输出大致为:

Multiplied by 5: 5 Multiplied by 5: 10 Multiplied by 5: 15 Multiplied by 5: 20 Multiplied by 5: 25 Multiplied by 10: 10 Multiplied by 10: 20 Multiplied by 10: 30 Multiplied by 10: 40 Multiplied by 10: 50

测试验证:用 PAssert 校验每个分支的结果

Katas 为每个任务都配了隐藏的单元测试。Branching 的测试位于 learning/katas/java/Core Transforms/Branching/Branching/test/org/apache/beam/learning/katas/coretransforms/branching/TaskTest.java:

import org.apache.beam.sdk.testing.PAssert; import org.apache.beam.sdk.testing.TestPipeline; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.values.PCollection; import org.junit.Rule; import org.junit.Test; public class TaskTest { @Rule public final transient TestPipeline testPipeline = TestPipeline.create(); @Test public void branching() { Create.Values<Integer> values = Create.of(1, 2, 3, 4, 5); PCollection<Integer> numbers = testPipeline.apply(values); PCollection<Integer> mult5Results = Task.applyMultiply5Transform(numbers); PCollection<Integer> mult10Results = Task.applyMultiply10Transform(numbers); PAssert.that(mult5Results) .containsInAnyOrder(5, 10, 15, 20, 25); PAssert.that(mult10Results) .containsInAnyOrder(10, 20, 30, 40, 50); testPipeline.run().waitUntilFinish(); } }

测试的关键点:

  1. TestPipeline+@Rule:TestPipeline.create()通过 JUnit 的@Rule自动管理测试管道的生命周期,testPipeline.run()真正执行管道,waitUntilFinish()等待执行完成;
  2. 直接复用 Task 的两个 Transform 方法:测试调用Task.applyMultiply5Transform(numbers)与Task.applyMultiply10Transform(numbers),说明参考实现将 Transform 提取为静态方法是为了可测试性;
  3. PAssert.that(...).containsInAnyOrder(...):分别对两个分支的结果进行断言——mult5Results必须恰好包含{5, 10, 15, 20, 25},mult10Results必须恰好包含{10, 20, 30, 40, 50}。containsInAnyOrder不关心元素顺序,只关心集合内容,这符合分布式 PCollection 元素顺序不确定的特性;
  4. 分支互不干扰的验证:如果某个 Transform 意外地"消费"或修改了共享输入,另一个分支的断言就会失败。因此这个测试同时验证了"分支不会消耗或改变输入"这一核心语义。

底层原理:为什么"复用输入"是安全的

从源码结构看,Branching 之所以成立,依赖 Beam 的两大设计基石:

1. PCollection 不可变(immutable)

PCollection本身不存储数据实体,而是"对一批元素的逻辑视图"。任何 Transform 应用到一个 PCollection 上,都是创建出一个新的PCollection(input.apply(...)的返回值),原始输入对象不被修改。因此在 Beam 中,一个 PCollection 被多个 Transform 引用是完全安全的,不存在"数据被抢先消费"的问题。这与迭代器(Iterator)等"消耗型"数据结构有本质区别。

2. 惰性执行(Lazy Execution)与 DAG 描述

pipeline.apply(...)只是在构建一个描述计算逻辑的有向无环图(DAG):节点是 PTransform,边是 PCollection。只有调用pipeline.run()时,图才会被真正提交给 Runner 执行。因此两个分支(×5 与 ×10)的 Transform 在被写入管道图时互不感知对方,Runner 在调度阶段再根据图的拓扑关系决定如何并行执行——分支结构天然适合并行,这也是 Beam 模型支持"一个输入、多路消费"的工程基础。

对于本 Kata 使用的Create.of(1, 2, 3, 4, 5)这种内存数据源,numbers这个 PCollection 作为两条分支的共同上游,会被下游两个 Transform 分别读取,而数据本身只有一份、不会重复创建或消耗。

如何运行这个 Kata

本 Kata 是 learning/katas/java 课程体系的一部分,按官方 README 的设置方式即可运行:

  1. 使用 IntelliJ Education(或安装了 EduTools 插件的 IntelliJ)"Open" 打开learning/katas/java目录;
  2. 按提示 "Import Gradle project" 并配置 Gradle;
  3. 等待 Gradle 构建完成;
  4. 打开 "Project Structure" 配置项目 SDK(例如 JDK 8);
  5. 在 "Project" 工具窗口切换到 "Course" 视图,即可看到包含 Branching 在内的全部 Kata 任务,逐课练习。

完成练习后,可直接运行Task.java的main方法观察两个分支的日志输出,或运行TaskTest.java中的branching()测试来校验实现是否正确。仓库中的 learning/katas/java/Core Transforms/Branching/Branching/task-info.yaml 还注明了该任务在 Beam Playground 中的元数据(名称Branching、分类Core Transforms、复杂度BASIC),同一练习也可以在上述平台上以在线方式尝试。

延伸:分支与相关 Core Transforms 的配合

分支模式是构建复杂管道的底座能力,在 Katas 的 "Core Transforms" 章节中,它常与其他变换组合使用,例如:

  • Flatten:把多个分支的PCollection合并成一个(见同目录下的Flatten/);
  • Partition:按条件把单个PCollection拆分成多个输出(与 Branching 方向相反);
  • Side Input / Side Output:分支之外的数据共享与多路输出手段;
  • CoGroupByKey / Combine:分支结果再做关联或聚合。

一个典型组合是:先把一份输入"兵分两路"做不同预处理,再用Flatten汇合,或者让一个分支的结果作为另一个分支的 Side Input。理解 Branching 这一基础模式,是掌握后续组合变换的前提。

小结

通过本 Kata,你掌握了 Apache Beam 分支(Branching)的核心要点:

  • 同一个PCollection可以作为多个 Transform 的输入,既不会被消耗,也不会被改变;
  • 分支通过"共享输入、独立 apply"实现,代码上只是把同一个PCollection变量依次传给多个 Transform;
  • 每个分支产出独立的新PCollection,可以分别Log观察、分别被PAssert断言;
  • 底层依赖PCollection的不可变性与 Beam 的惰性 DAG 执行模型,分支天然适合并行调度。

参考实现(Task.java)与测试(TaskTest.java)可以在仓库中直接对照学习,作为你编写分支式 Beam 管道的起点模板。

【免费下载链接】beam

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

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

相关推荐

上一篇:如何构建可靠的电商系统:SimplCommerce测试策略与最佳实践
下一篇:Sphinx 环境依赖记录(environment record-dependencies):源码级剖析与增量构建实战

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询