)
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载本指南基于 Apache Beam 官方学习项目 Tour of Beam 中 core-transforms/branching 单元 展开系统讲解 Beam 中「分支 PCollectionBranching 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作用于同一输入PCollectionString input ...; PCollectionString aCollection input.apply(aTrans, ParDo.of(new DoFnString, String(){ ProcessElement public void processElement(ProcessContext c) { if(c.element().startsWith(A)){ c.output(c.element()); } } })); PCollectionString bCollection input.apply(bTrans, ParDo.of(new DoFnString, String(){ ProcessElement public void processElement(ProcessContext c) { if(c.element().startsWith(B)){ c.output(c.element()); } } }));注意两处input.apply(...)都作用于同一个input对象只是各自携带了不同的变换名称aTrans、bTrans用于在管道图上区分这两个分支节点。Python两行beam.Filterstarts_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反转与大写并行PCollectionString input pipeline.apply(Create.of(Apache, Beam, is, an, open, source, unified, programming, model, To, define, and, execute, data, processing, pipelines, Go, SDK)); PCollectionString reverseCollection input.apply(aTrans, ParDo.of(new DoFnString, String(){ ProcessElement public void processElement(ProcessContext c) { c.output(new StringBuilder(c.element()).reverse().toString()); } })); PCollectionString upperCollection input.apply(aTrans, ParDo.of(new DoFnString, 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两个函数返回两个独立的PCollectionfunc 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); PCollectionInteger input pipeline.apply(Create.of(1, 2, 3, 4, 5)); PCollectionInteger mult5Results applyMultiply5Transform(input); PCollectionInteger mult10Results applyMultiply10Transform(input); mult5Results.apply(Log multiplied by 5: , ParDo.of(new LogOutputInteger(Multiplied by 5: ))); mult10Results.apply(Log multiplied by 10: , ParDo.of(new LogOutputInteger(Multiplied by 10: ))); pipeline.run();其中两个分支变换都用MapElements.into(integers()).via(...)实现static PCollectionInteger applyMultiply5Transform(PCollectionInteger input) { return input.apply(Multiply by 5, MapElements.into(integers()).via(num - num * 5)); } static PCollectionInteger applyMultiply10Transform(PCollectionInteger input) { return input.apply(Multiply by 10, MapElements.into(integers()).via(num - num * 10)); }LogOutputT是一个自定义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(prefixMultiplied by 5: ) mult10_results | Log multiply 10 Output(prefixMultiplied by 10: )示例中还定义了一个可复用的Output(beam.PTransform)打印变换展示如何把日志输出也封装成标准变换挂到分支上。分支与拆分类变换的区别Partition 与多输出 ParDoBranching 常被与另外两种一进多出的机制混淆这里做一次关键区分帮助你按需选型。Branching vs Partition按分区函数拆分Partition是 Beam 中把单个PCollection按你提供的分区函数拆成固定数量N 个更小的集合的变换返回PCollectionListJava/[]beam.PCollectionGo/ tuplePython通过索引访问各分区。与 Branching 的核心差异在于Branching对同一输入挂多个独立变换各分支的处理逻辑可以完全不同如一个过滤、一个反转、一个大写Partition只做拆分所有分区共享同一个分区函数决定去向输出数量在建图时必须确定可运行时通过命令行参数传入但不能在管道中途依据数据动态决定。Go SDK 中beam.Partition(s, n, fn, col) []beam.PCollection的实现位于 sdks/go/pkg/beam/partition.go它返回分区切片、按下标访问这与branching中分支返回多个PCollection的形态在 API 层面有相似之处但语义不同分支的每个输出对应一次完整的独立变换而分区的每个输出只是同一个变换对不同元素的归类结果。Branching vs 多输出 ParDoAdditional outputs多输出 ParDo 则是在单个DoFn内部通过MultiOutputReceiverTupleTagJava、with_outputs()pvalue.TaggedOutputPython或 emitter 函数Go 的beam.ParDo2/beam.ParDo3/beam.ParDoN把元素发射到多个带标签的输出。它与 Branching 的区别在于Branching是多个变换、共享同一输入每个变换只有一个输出多输出 ParDo是一个变换、多个输出每个元素按运行时的条件被路由到不同的输出通道。何时用哪个如果各分支的处理逻辑差异巨大、需要独立命名与独立测试优先用 Branching如果只是按条件把元素分流到不同输出且希望共享同一个DoFn的处理逻辑优先用多输出。Go SDK 中多输出变换的源码定义可见 sdks/go/pkg/beam/pardo.goParDo2/ParDo3分别返回 2/3 个PCollectionParDoN返回切片。在 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与多输出 ParDoTupleTag/TaggedOutput实现感受两种模式的代码组织差异。小结Branching PCollections 是 Apache Beam 管道设计中一输入、多输出的基础范式其成立的前提是变换逐元素处理且不消耗输入。掌握它意味着你能把一条数据处理流水线优雅地拆成多条并行分支既可以在同一份数据上挂多个条件过滤如按首字母分流也可以对同一份数据做多种形态的转换如反转、大写、乘以 5、乘以 10。再结合Partition固定数量拆分与多输出ParDo单变换多通道路由你就可以针对不同场景选择最合适的扇出方案构建出结构清晰、可独立调试的 Beam 管道。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐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 GApache Beam 分支Branching实战在 Go SDK 中对同一 PCollection 应用多个变换Apache Beam 分支Branching实战在 Go SDK 中对同一 PCollection 应用多个变换 导读 本文围绕 Apache Beam批处理流处理大数据上一篇Temporal TypeScript SDK核心功能解析Workflow、Activity与Worker实战指南下一篇Formily实时验证输入反馈与错误提示设计创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考