- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本指南基于 Apache Beam 官方学习项目 Tour of Beam 中 core-transforms/branching 单元 展开,系统讲解 Beam 中「分支 PCollection(Branching PCollections)」这一核心编程范式:多个变换同时处理同一个 PCollection 而不消耗、不改变其内容。读完本文,你将掌握分支技术的基本原理、Java / Python / Go 三种 SDK 的完整写法,理解它与Partition、多输出ParDo之间的区别,并能在自己的管道中安全地对同一份数据执行多条并行处理链路。
分支的本质:变换并不"消费"PCollection
要理解 Branching,首先要纠正一个常见直觉。在 Apache Beam 中,变换(PTransform)并不会把输入PCollection"消费"掉。恰恰相反,每一个变换都是逐元素(per-element)地考虑输入PCollection中的每个元素,然后基于这些元素创建一个全新的PCollection作为输出。原始输入PCollection在变换执行后依然存在、完好无损。
正是这一设计使得"分支"成为可能:既然变换不会破坏输入,那么同一个PCollection就可以被任意多个变换同时当作输入,每条变换分支各自产出独立的输出PCollection。这种"一进多出"的扇形扩展(fan-out)结构,是 Beam 管道中实现多路并行处理的基础手段,也是 Tour of Beam 的 Core Transforms 课程 中紧随map、additional-outputs之后的核心单元(模块定义见 module-info.yaml,该分支单元的元数据定义在 unit-info.yaml)。
从源码结构也可以印证这一点:Beam 的核心抽象是"输入PCollection+ 变换 → 输出PCollection"的函数式映射关系,变换本身不持有或销毁任何数据。因此你可以放心地对同一个PCollection挂接任意数量的下游变换。
分支模式一:多个变换共享同一个输入 PCollection
这是分支技术最直接的用法。原文给出的场景如下:管道从数据库表中读出"名字"(以字符串表示)并创建一个PCollection;随后,变换 A 提取所有以字母 'A' 开头的名字,变换 B 提取所有以字母 'B' 开头的名字。变换 A 和变换 B 的输入是同一个PCollection——这正是分支的关键。
┌──> Transform A(提取 "A" 开头)──> PCollection A 输入 PCollection ────────┤ └──> Transform B(提取 "B" 开头)──> PCollection B下面分别给出三种 SDK 的写法(均以原文档代码为骨架,Python 部分用beam.Filter实现条件过滤,beam.Filter的定义见 sdks/python/apache_beam/transforms/core.py)。
Java:两个ParDo作用于同一输入
PCollection<String> input = ...; PCollection<String> aCollection = input.apply("aTrans", ParDo.of(new DoFn<String, String>(){ @ProcessElement public void processElement(ProcessContext c) { if(c.element().startsWith("A")){ c.output(c.element()); } } })); PCollection<String> bCollection = input.apply("bTrans", ParDo.of(new DoFn<String, String>(){ @ProcessElement public void processElement(ProcessContext c) { if(c.element().startsWith("B")){ c.output(c.element()); } } }));注意两处input.apply(...)都作用于同一个input对象,只是各自携带了不同的变换名称("aTrans"、"bTrans"),用于在管道图上区分这两个分支节点。
Python:两行beam.Filter
starts_with_a = input | beam.Filter(lambda x: x.startswith('A')) starts_with_b = input | beam.Filter(lambda x: x.startswith('B'))这里同一个input被两个Filter变换分别消费,互不影响。beam.Filter会保留谓词返回True的元素、丢弃其余元素,非常适合做这种"按条件分流"的分支。
Go:两个ParDo复用输入
原文档 Go 代码块中混入了一些 Java 风格的笔误(如element.startsWith(...)与返回类型不一致),这里按仓库中可运行示例的惯用写法给出修正版本:
input = ...; outputA := applyTransformA(s, input) outputB := applyTransformB(s, input) func applyTransformA(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, startWithA, input) } func applyTransformB(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, startWithB, input) } func startWithA(element string) string { if strings.HasPrefix(element, "A") { return element } return "" } func startWithB(element string) string { if strings.HasPrefix(element, "B") { return element } return "" }Go SDK 中beam.ParDo即挂接变换的入口,同一input可以被多次传入不同的ParDo,实现分支。需要强调的是:分支的所有下游分支共享输入,但每个元素会被各分支独立处理——这既意味着处理逻辑互不干扰,也意味着如果分支过多,同一元素会被重复处理多次(这是扇出模型固有的代价,属于正常语义)。
分支模式二:独立变换分别转换同一份数据
除了"条件过滤"式分支,更常见的分支形态是:对同一个PCollection分别应用不同的转换函数,各分支并行产出不同的结果集。原文的练习示例是:同一个字符串PCollection,一个分支把每个元素反转,另一个分支把每个元素转成大写。
Java:反转与大写并行
PCollection<String> input = pipeline.apply(Create.of("Apache", "Beam", "is", "an", "open", "source", "unified", "programming", "model", "To", "define", "and", "execute", "data", "processing", "pipelines", "Go", "SDK")); PCollection<String> reverseCollection = input.apply("aTrans", ParDo.of(new DoFn<String, String>(){ @ProcessElement public void processElement(ProcessContext c) { c.output(new StringBuilder(c.element()).reverse().toString()); } })); PCollection<String> upperCollection = input.apply("aTrans", ParDo.of(new DoFn<String, String>(){ @ProcessElement public void processElement(ProcessContext c) { c.output(c.element().toUpperCase()); } }));(实践中建议给两个分支分别命名,例如"reverseTrans"、"upperTrans",便于在管道可视化与日志中区分。)
Python:两个Map分支
reversed = input | reverseString(...) toUpper = input | toUpperString(...)其中reverseString、toUpperString可以是自定义的beam.PTransform,也可以是beam.Map的 lambda,例如:
reversed = input | beam.Map(lambda x: x[::-1]) to_upper = input | beam.Map(lambda x: x.upper())Go:返回两个PCollection的分支函数
仓库中该分支单元的可运行示例 go-example/main.go 给出了完整的 Go 分支实现:applyTransform接收一个输入PCollection,内部调用reverseString与toUpperString两个函数,返回两个独立的PCollection:
func applyTransform(s beam.Scope, input beam.PCollection) (beam.PCollection, beam.PCollection) { reversed := reverseString(s, input) toUpper := toUpperString(s, input) return reversed, toUpper } func reverseString(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, reverseFn, input) } func toUpperString(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, strings.ToUpper, input) }reverseFn使用 rune 切片实现字符串反转(兼容中文等多字节字符),完整逻辑见 main.go。这个示例还展示了分支的经典调试方式:用debug.Printf分别打印两条分支的结果。
可直接运行的完整分支示例
除了文档中的片段,仓库为每个 SDK 都提供了标注beam-playground: name: branching、complexity: MEDIUM的完整可运行代码,可直接在 Tour of Beam 的 Playground 窗口中运行实验(对应元数据见 unit-info.yaml)。
Java 完整示例:乘以 5 与乘以 10 两个分支
Java 示例 Task.java 演示了对整数PCollection做两个数值分支:
PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline = Pipeline.create(options); PCollection<Integer> input = pipeline.apply(Create.of(1, 2, 3, 4, 5)); PCollection<Integer> mult5Results = applyMultiply5Transform(input); PCollection<Integer> mult10Results = applyMultiply10Transform(input); mult5Results.apply("Log multiplied by 5: ", ParDo.of(new LogOutput<Integer>("Multiplied by 5: "))); mult10Results.apply("Log multiplied by 10: ", ParDo.of(new LogOutput<Integer>("Multiplied by 10: "))); pipeline.run();其中两个分支变换都用MapElements.into(integers()).via(...)实现:
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)); }LogOutput<T>是一个自定义DoFn,通过LOG.info输出每个分支的元素,便于观察分支结果。
Python 完整示例
Python 示例 task.py 结构与之对应:
with beam.Pipeline() as p: input = p | beam.Create([1, 2, 3, 4, 5]) mult5_results = input | beam.Map(lambda num: num * 5) mult10_results = input | beam.Map(lambda num: num * 10) mult5_results | 'Log multiply 5' >> Output(prefix='Multiplied by 5: ') mult10_results | 'Log multiply 10' >> Output(prefix='Multiplied by 10: ')示例中还定义了一个可复用的Output(beam.PTransform)打印变换,展示如何把"日志输出"也封装成标准变换挂到分支上。
分支与"拆分"类变换的区别:Partition 与多输出 ParDo
Branching 常被与另外两种"一进多出"的机制混淆,这里做一次关键区分,帮助你按需选型。
Branching vs Partition(按分区函数拆分)
Partition是 Beam 中把单个PCollection按你提供的分区函数拆成固定数量(N 个)更小的集合的变换,返回PCollectionList(Java)/[]beam.PCollection(Go)/ tuple(Python),通过索引访问各分区。与 Branching 的核心差异在于:
- Branching对同一输入挂多个独立变换,各分支的处理逻辑可以完全不同(如一个过滤、一个反转、一个大写);
- Partition只做"拆分",所有分区共享同一个分区函数决定去向,输出数量在建图时必须确定(可运行时通过命令行参数传入,但不能在管道中途依据数据动态决定)。
Go SDK 中beam.Partition(s, n, fn, col) []beam.PCollection的实现位于 sdks/go/pkg/beam/partition.go,它返回分区切片、按下标访问,这与branching中"分支返回多个PCollection"的形态在 API 层面有相似之处,但语义不同:分支的每个输出对应一次完整的独立变换,而分区的每个输出只是同一个变换对不同元素的归类结果。
Branching vs 多输出 ParDo(Additional outputs)
多输出 ParDo 则是在单个DoFn内部通过MultiOutputReceiver+TupleTag(Java)、with_outputs()+pvalue.TaggedOutput(Python)或 emitter 函数(Go 的beam.ParDo2/beam.ParDo3/beam.ParDoN)把元素发射到多个带标签的输出。它与 Branching 的区别在于:
- Branching是"多个变换、共享同一输入",每个变换只有一个输出;
- 多输出 ParDo是"一个变换、多个输出",每个元素按运行时的条件被路由到不同的输出通道。
何时用哪个?如果各分支的处理逻辑差异巨大、需要独立命名与独立测试,优先用 Branching;如果只是按条件把元素分流到不同输出(且希望共享同一个DoFn的处理逻辑),优先用多输出。Go SDK 中多输出变换的源码定义可见 sdks/go/pkg/beam/pardo.go(ParDo2/ParDo3分别返回 2/3 个PCollection,ParDoN返回切片)。
在 Tour of Beam 课程中的定位与练习建议
Branching 单元隶属于 learning/tour-of-beam/learning-content/core-transforms 模块,是 Core Transforms 系列(map→additional-outputs→branching→combine→composite→flatten→partition→side-inputs)中的第三个单元,复杂度评级为 MEDIUM,覆盖 Java、Python、Go 三种 SDK(见 unit-info.yaml)。练习建议:
- 在 Playground 直接运行上述三个 SDK 的完整示例,观察两条分支各自的输出日志;
- 把条件过滤式分支改成数值分支(如"大于 100 / 小于 100"),体会过滤与转换两类分支的写法差异;
- 在分支上继续叠加变换,例如先分支再对每条分支各自
Map/Filter,理解分支之后还可以继续分支的树状组合能力; - 对比练习:把同一个"分流"需求分别用 Branching(多个
Filter)与多输出 ParDo(TupleTag/TaggedOutput)实现,感受两种模式的代码组织差异。
小结
Branching PCollections 是 Apache Beam 管道设计中"一输入、多输出"的基础范式,其成立的前提是变换逐元素处理且不消耗输入。掌握它意味着你能把一条数据处理流水线优雅地拆成多条并行分支:既可以在同一份数据上挂多个条件过滤(如按首字母分流),也可以对同一份数据做多种形态的转换(如反转、大写、乘以 5、乘以 10)。再结合Partition(固定数量拆分)与多输出ParDo(单变换多通道路由),你就可以针对不同场景选择最合适的扇出方案,构建出结构清晰、可独立调试的 Beam 管道。
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 核心变换实战:Branching PCollections 多路分支处理
Apache Beam 核心变换实战:Branching PCollections 多路分支处理 Branching(分支)是 Apache Beam 中最基础
Apache Beam Go SDK 实战:用 Branching 将一个 PCollection 分支为多个变换输出
Apache Beam Go SDK 实战:用 Branching 将一个 PCollection 分支为多个变换输出 本篇文章围绕 Apache Beam G
Apache Beam 分支(Branching)实战:在 Go SDK 中对同一 PCollection 应用多个变换
Apache Beam 分支(Branching)实战:在 Go SDK 中对同一 PCollection 应用多个变换 导读 本文围绕 Apache Beam
批处理流处理大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考