ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Beam Python 实战:用 ParDo 与 DoFn 实现通用并行处理(Katas 逐行拆解)

Apache Beam Python 实战:用 ParDo 与 DoFn 实现通用并行处理(Katas 逐行拆解) 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本指南以 Apache Beam 官方 Katas 训练项目中的ParDo关卡learning/katas/python/Core Transforms/Map/ParDo/task.md为核心讲解 Beam 中最核心的并行处理原语ParDo它的执行语义一进零出、一进一出、一进多出、DoFn.process的写法与生命周期并带你在 Python SDK 中亲手实现一个将输入元素乘以 10的完整任务配合仓库内测试用例完成验证。读完本文你将能独立编写、运行并测试自己的 ParDo 变换为理解 Beam 的批流统一编程模型打下基础。一、任务背景Katas 是什么learning/katas/python是 Beam 仓库内置的一套按难度分级的 Python 编码训练项目Katas覆盖 Common Transforms、Core Transforms、IO、Windowing 等主题。其中Core Transforms/Map目录下依次排列着Map、FlatMap、ParDo、ParDo OneToMany等关卡每个关卡由三部分组成task.md任务描述与提示即本指南关联的核心文档task.py需要你补全/实现的源码骨架该关卡的参考实现tests/test_task.py隐藏的单元测试用于自动校验你的实现是否正确。task-info.yaml中标注了本关卡为complexity: BASIC基础难度、categories: Core Transforms说明它是入门并行变换的第一课。每个任务还带有 beam-playground 元数据name: MapPardo意味着这类题目可直接在 Beam Playground 在线环境运行。二、ParDo通用并行处理的核心原语任务文档开门见山地给出了 ParDo 的权威定义ParDo is a Beam transform for generic parallel processing. The ParDo processing paradigm is similar to the Map phase of a Map/Shuffle/Reduce-style algorithm: a ParDo transform considers each element in the input PCollection, performs some processing function (your user code) on that element, and emits zero, one, or multiple elements to an output PCollection.翻译并拆解这一定义可以得到 ParDo 的三个关键事实逐元素处理ParDo 作用于输入PCollection中的每一个元素其处理模式类似于 Map/Shuffle/Reduce 算法中的Map 阶段——这是理解 Beam 数据流模型的重要类比用户代码即处理函数具体的处理逻辑由你编写即DoFn中的用户代码Beam 负责在分布式环境中调度执行输出数量不受限针对每个输入元素ParDo 可以输出零个、一个或多个元素——这正是它与严格一对一映射的Map的本质区别也是实现过滤零输出、拆分多输出等逻辑的基础。从源码看Python SDK 中ParDo类定义于 sdks/python/apache_beam/transforms/core.pyclass ParDo(PTransformWithSideInputs)它本身是一个PTransform的子类在执行时会把你的DoFn实例分发到各个工作单元Bundle中并行调用。三、Kata 任务拆解把每个元素乘以 10任务文档给出的 Kata 要求非常明确Kata:Please write a simple ParDo that maps the input element by multiplying it by 10.即编写一个简单的 ParDo把输入元素乘以 10 后输出。结合配套提示hint完成它需要掌握两个 API 要点重写DoFn.process方法使用ParDo配合DoFn一起使用。3.1 什么是 DoFnDoFnDo Function是承载你业务逻辑的类Beam 规定用户代码必须包装在DoFn中才能交给 ParDo 执行。在 Python SDK 中DoFn基类同样位于 sdks/python/apache_beam/transforms/core.pyclass DoFn(WithTypeHints, HasDisplayData, urns.RunnerApiFn)约第 597 行。它的核心思想是你只管描述对单个元素做什么Beam 负责把它并行化到成千上万个 worker 上。3.2 参考实现逐行解读仓库中 task.py 给出了完整答案全貌如下import apache_beam as beam class MultiplyByTenDoFn(beam.DoFn): def process(self, element): yield element * 10 with beam.Pipeline() as p: (p | beam.Create([1, 2, 3, 4, 5]) | beam.ParDo(MultiplyByTenDoFn()) | beam.LogElements())逐部分拆解class MultiplyByTenDoFn(beam.DoFn)自定义的DoFn子类命名上遵循动词化 意图的惯例MultiplyByTen便于他人从名称读懂处理逻辑def process(self, element)核心处理方法对传入的单个element执行用户逻辑。注意process的返回方式是迭代器/生成器——使用yield逐个产出结果元素这是任务文档 hint 中强调的关键写法yield element * 10对输入元素做乘法并产出由于每个输入只产出一个输出这是一进一出one-to-one的典型用法beam.Create([1, 2, 3, 4, 5])把内存中的 Python 列表包装成一个输入PCollection在分布式 runner 上元素会被分发到多个 worker 并行处理beam.ParDo(MultiplyByTenDoFn())把DoFn实例交给ParDo变换执行——注意传入的是实例而不是类beam.LogElements()把输出 PCollection 中的每个元素打印到日志/控制台便于在本地验证结果with beam.Pipeline() as p:创建流水线上下文管理器with块结束时自动调用run()并等待执行完成。该流水线执行后的预期输出为10 20 30 40 503.3 为什么用 yield 而不是 return 列表process方法既可以直接yield生成器也可以返回一个Iterable。任务文档在姊妹关卡ParDo OneToMany中给出了明确提示You can return an Iterable for multiple elements or call yield for each element to return a generator.。用yield的工程价值在于对超大数据集而言yield是惰性的元素逐个产出不必一次性构造完整结果列表内存友好便于在一个process中组合多个输出逻辑先 yield A 再 yield B与 Beam 的流式处理模型天然契合元素可随产随发。四、从源码理解 DoFn 的完整生命周期仅仅实现process就能通过本关但理解DoFn的完整生命周期能让你写出更健壮的生产级代码。在 sdks/python/apache_beam/transforms/core.py 的DoFn基类中定义了以下可重写方法约第 699–770 行方法调用时机典型用途setup(self)每个 worker 上处理任何 bundle 之前仅调用一次建立数据库连接、加载模型、初始化昂贵资源start_bundle(self)每个 bundle一批元素处理之前批量初始化、打开批量写入缓冲区process(self, element, *args, **kwargs)对 PCollection 中的每个元素调用核心业务逻辑用yield产出零到多个输出process_batch(self, batch, *args, **kwargs)按批处理元素批量模式下的替代实现利用向量化库如 NumPy做批量计算finish_bundle(self)每个 bundle 处理完成之后冲刷缓冲区、提交批量事务teardown(self)整个 worker 生命周期结束前关闭连接、释放资源对照本关任务MultiplyByTenDoFn是无状态、无外部资源的纯计算因此只需实现process如果任务涉及读取文件每行 → 逐行处理 → 关闭文件则应分别在setup中打开、在teardown中关闭避免对每个元素重复开关资源。此外DoFn还支持通过process(elementDoFn.ElementParam, timestampDoFn.TimestampParam)这类参数注解声明需要注入的上下文如事件时间戳、PaneInfo这在流式窗口计算中非常实用属于进阶内容。五、用仓库自带测试验证你的实现Katas 之所以是训练营而非示例就在于每个任务都配有隐藏的自动化测试。本关的测试位于 tests/test_task.py核心逻辑如下import unittest from test_helper import get_file_output class TestCase(unittest.TestCase): def test_output(self): output get_file_output(pathtask.py) answers [10, 20, 30, 40, 50] for ans in answers: self.assertIn(ans, output, Incorrect output. Multiply each element by 10.)测试的运行原理是通过 test_helper.py 中的get_file_output用subprocess以当前 Python 解释器实际执行task.py捕获其标准输出并按行拆分然后断言10、20、30、40、50全部出现在输出中。这意味着测试真实运行你的流水线不是静态检查代码文本——任何语法错误、导入错误都会导致测试失败结果顺序不敏感用assertIn而非按序比较因为分布式执行本身不保证元素顺序这体现了 Beam元素无序、窗口有序的语义若你误写为element 10或element / 10测试会抛出 Incorrect output. Multiply each element by 10. 的失败信息。本地验证方式在仓库根目录执行cd learning/katas/python python Core\ Transforms/Map/ParDo/task.py python -m unittest discover -s Core Transforms/Map/ParDo/tests -p test_*.py输出中应能看到10、20、30、40、50五个数字且单元测试全部通过Ran 1 test ... OK。六、ParDo 与周边变换的对比理解了 ParDo 之后把它和同目录下的姊妹关卡放在一起对比能更清晰地定位它在 Beam 变换体系中的位置6.1 ParDo vs Map一进一出Map是 ParDo 的轻量语法糖。同目录 Map/task.py 展示了用 lambda 完成乘以 5的写法(p | beam.Create([10, 20, 30, 40, 50]) | beam.Map(lambda num: num * 5) | beam.LogElements())对比本关的 ParDo 写法二者对每个输入都恰好输出一个元素one-to-one功能等价但Map省去了定义DoFn子类的样板代码适合简单映射当逻辑复杂、需要复用或需要生命周期钩子时应回归 ParDo。从源码看beam.Map底层也是通过ParDoFlatMapDoFn等内部DoFn实现的只是把 lambda 包装成了DoFn。6.2 ParDo vs FlatMap / ParDo OneToMany一进多出姊妹关卡ParDo OneToMany要求把每个句子按空格切分成单词此时每个输入会产生多个输出。写法上只需在process中多次yieldclass WordSplittingDoFn(beam.DoFn): def process(self, element): for word in element.split( ): yield word这正是任务文档开篇所述emits zero, one, or multiple elements的完整体现过滤逻辑 零输出拆分逻辑 多输出同一个process框架天然覆盖三种场景这是Map无法做到的。七、实战建议与常见误区基于本关的实现与测试给出几点实战建议命名要表达语义DoFn类名用动词 宾语如MultiplyByTenDoFn让流水线读起来像业务文档保持process无状态除setup初始化的共享资源外不要在process内维护跨元素状态如计数器分布式环境下多个 worker 各自持有副本结果不可预期需要跨元素聚合时请使用 Beam 的 State/Timer 或 Combine输出必须可迭代忘记yield/返回可迭代对象会导致运行时错误单输出场景请确保写yield element * 10而非return element * 10后者会把整个标量当作一次产出失败或语义错误善用beam.Map做简单映射一行 lambda 能表达的映射用Map涉及资源管理、多输出或复用逻辑再用 ParDo以测试为准绳Katas 的隐藏测试即是最好的行为规范——先跑通测试再谈优化。八、总结本关虽然标注为 BASIC 难度却浓缩了 Beam 最核心的编程心智ParDo 把对每个元素做什么用户代码与如何并行做运行时调度彻底解耦。通过乘以 10这个小任务你掌握了DoFn.process的 yield 写法、ParDo的一进多出语义、DoFn生命周期钩子以及仓库测试的自动化验证方式。下一步可以继续挑战同目录的ParDo OneToMany一词切多词与 Windowing 系列关卡把 ParDo 放入窗口与触发器的大背景下理解流式处理。相关代码均可在本仓库中直接查阅任务描述见 task.md参考实现见 task.py测试见 tests/test_task.pyDoFn/ParDo源码见 sdks/python/apache_beam/transforms/core.py。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Katas 实战用 ParDo DoFn 实现数据过滤Filter using ParDoApache Beam Java Katas 实战用 ParDo DoFn 实现数据过滤Filter using ParDo 本文是 Apache B大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战用 ParDo 实现通用并行处理变换Apache Beam Kotlin Kata 实战用 ParDo 实现通用并行处理变换 本指南以 Apache Beam 仓库中 Kotlin 版 Kata大数据批处理流处理数据工程Apache Beam Go SDK 实战用 ParDo 与 DoFn 实现 Filter 过滤变换Apache Beam Go SDK 实战用 ParDo 与 DoFn 实现 Filter 过滤变换 本教程基于 Apache Beam 官方 Katas 课大数据批处理流处理数据工程上一篇拓展Rust编程边界Diesel——安全且可扩展的ORM和查询构建器下一篇开源项目 acwj 使用教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表