免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

Apache Beam Python Kata:用 ParDo 实现 OneToMany 分词(一进多出)

Apache Beam Python Kata:用 ParDo 实现 OneToMany 分词(一进多出) 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的ParDo是面向通用并行处理的核心变换其处理范式与 Map/Shuffle/Reduce 类算法中的 Map 阶段类似对输入 PCollection 中的每个元素执行用户自定义的处理函数并向外输出零个、一个或多个元素。本文围绕 Beam Katas 系列中 ParDo OneToMany 这一实战练习讲解如何在 Python SDK 中编写一个DoFn把每个输入句子按空白符切分为多个单词同时结合仓库中的参考实现、单元测试与底层源码厘清ParDo、process方法的两种产出方式以及它与Map、FlatMap等轻量抽象之间的关系。读完本文你将能够独立完成该 Kata并掌握用ParDo做一进多出OneToMany元素变换的标准写法。练习背景从 ParDo 到 ParDo OneToMany在 learning/katas/python/Core Transforms/Map/lesson-info.yaml 中可以看到Map 这一课由四个练习按难度递进组成ParDo→ParDo OneToMany→Map→FlatMap。本练习ParDo OneToMany位于该课程第二关它的前导练习ParDo要求实现一个把输入元素乘以 10 的DoFn一对一映射参考实现见 ParDo/task.py而本关则将映射关系扩展到一对多。练习题目本身定义在 task.md 中Kata请编写一个ParDo将每个输入句子映射为由空白符 分词得到的单词。也就是说输入 PCollection 中的元素是句子字符串例如Hello Beam期望的输出是这些句子拆分后的每个单词例如Hello、Beam且一个句子会对应产出多个单词元素。这正是ParDo区别于一对一Map的核心场景process方法一次可以向外发出多个元素。DoFn 与 process 方法一进多出的两种写法在 Beam Python SDK 中使用ParDo需要配合一个beam.DoFn子类。核心要点是重写process方法其参数element就是输入 PCollection 中的单个元素。要对外产出多个结果官方推荐两种等价写法返回一个 Iterable如列表、元组框架会把返回的容器展开将其中的每一个元素分别发射到输出 PCollection使用yield逐个产出让process成为生成器函数每次yield的元素都会被框架作为独立输出元素发射。本练习的参考答案位于 task.py它采用了第二种写法import apache_beam as beam class BreakIntoWordsDoFn(beam.DoFn): def process(self, element): for w in element.split(): yield w with beam.Pipeline() as p: (p | beam.Create([Hello Beam, It is awesome]) | beam.ParDo(BreakIntoWordsDoFn()) | beam.LogElements())逐行拆解这段代码class BreakIntoWordsDoFn(beam.DoFn)任何传给beam.ParDo的用户处理逻辑都必须继承beam.DoFndef process(self, element)重写process方法element是输入 PCollection 中的一条数据此处为一个句子字符串element.split()Python 字符串的split()不带参数时会按任意连续空白字符空格、制表符、换行等切分并返回单词列表。由于题目明确要求按空白符分词这里既可以直接写element.split( )也可以依赖默认行为for w in element.split(): yield w逐单词产出实现一个句子 → 多个单词的 OneToMany 变换。作为对照若换成return element.split()返回 Iterable结果完全一致beam.Create([...])构造输入 PCollection包含两个句子beam.ParDo(BreakIntoWordsDoFn())把DoFn实例应用到每个元素上这是本练习的主角beam.LogElements()将输出 PCollection 的每个元素打印到日志便于在终端观察运行结果。运行这段程序后输出应为五个单词Hello、Beam、It、is、awesome两条句子的分词结果被拍平到同一条输出 PCollection 中。通过测试验证实现是否正确每个 Kata 都配有一套隐藏的单测位于 tests/test_task.pyimport unittest from test_helper import test_is_not_empty, get_file_output class TestCase(unittest.TestCase): def test_not_empty(self): self.assertTrue(test_is_not_empty(), The output is empty) def test_output(self): output get_file_output(pathtask.py) answers [Hello, Beam, It, is, awesome] for word in answers: self.assertIn(word, output, Incorrect output. Break each sentence into words.)测试包含两个用例test_not_empty通过 test_helper.py 中的test_is_not_empty()检查task.py是否非空防止提交空文件test_output调用get_file_output(pathtask.py)用当前 Python 解释器以子进程方式实际运行task.py捕获其标准输出并按行拆分然后逐一断言Hello、Beam、It、is、awesome这五个单词都出现在输出中。从test_helper.py的实现可以看出这套测试机制并不做静态代码检查而是真实执行你的 pipeline并校验运行结果因此只要process方法能够正确把句子拆成单词无论你采用yield还是返回 Iterable测试都会通过。与 Map、FlatMap 的对比何时用 ParDo 的一进多出在 lesson-info.yaml 编排的四关练习中ParDo OneToMany之后还有Map与FlatMap两个练习它们恰好是解决同一类问题的更轻量写法Map一对一beam.Map(fn)只接受一个元素进、一个元素出的普通函数等价于process只yield一个元素的ParDo。参考实现见 Map/task.pyFlatMap一对多、函数式beam.FlatMap(fn)接受一个返回 Iterable 的普通函数并把返回的 iterable 扁平化后逐个发射。它正是用函数代替DoFn的ParDoOneToMany。参考实现见 FlatMap/task.py例如(p | beam.Create([Apache Beam, Unified Batch and Streaming]) | beam.FlatMap(lambda sentence: sentence.split()) | beam.LogElements())ParDoDoFn最通用当处理逻辑需要维护状态、接收额外参数、访问上下文如计时器、窗口信息或产出一对多结果时DoFn的完整生命周期setup/start_bundle/process/finish_bundle/teardown是最灵活的选择。关于FlatMap的底层语义可以在 sdks/python/apache_beam/transforms/core.py 的源码注释中找到权威说明FlatMap与ParDo类似只是它接受一个 callable 来指定变换该 callable 必须为输入 PCollection 的每个元素返回一个 iterable这些 iterable 中的元素会被扁平化进输出 PCollection。可以看到本练习句子的split()结果是一个列表Iterable这一事实正是FlatMap能够一行 Lambda 解决同样问题的前提——而ParDo版本则通过yield显式地把列表展开两者输出完全等价。因此本练习ParDo OneToMany的价值在于让你理解ParDo的通用能力一个输入元素可以产出任意多个输出元素。掌握了这一点再学习Map与FlatMap时就能清楚地认识到它们只是ParDo在特定场景下的语法糖。练习的运行环境与工程元数据本 Kata 属于 Beam Katas 交互式课程的一部分有两种运行方式PyCharm Education或安装 EduTools 插件的 PyCharm按 learning/katas/python/README.md 的指引以该目录为根创建项目、选定 Python 解释器然后在 Course 视图中逐课完成练习并运行内置测试独立执行直接运行task.py依赖apache_beam已安装观察beam.LogElements()打印出的分词结果。此外每个任务还带有面向 Beam Playground 的元数据声明见 task-info.yaml# beam-playground: # name: MapParDoOneToMany # description: Task from katas is a ParDo that maps each input sentence into # words splitter by whitespace ( ). # multifile: false # context_line: 40 # categories: # - Core Transforms # complexity: BASIC # tags: # - transforms # - strings其中name: MapParDoOneToMany是该任务在 Playground 中的唯一标识categories: Core Transforms说明其归属complexity: BASIC标记为基础难度context_line: 40指向task.py中第 40 行with beam.Pipeline() as p:所在处用于在 IDE / Playground 中定位代码上下文。这些元数据说明本练习是 Beam 官方为初学者设计的入门级变换训练重点就是一个DoFn的process如何产出多个元素。小结通过本练习你可以掌握 Beam Python SDK 中ParDo最核心的 OneToMany 用法继承beam.DoFn并重写process(element)用yield逐元素产出或直接返回一个 Iterable实现一个输入、多个输出用beam.Create构造输入、beam.ParDo应用逻辑、beam.LogElements观察输出理解Map一对一函数式、FlatMap一对多函数式本质上都是ParDo的轻量特化。句子分词是文本处理类 pipeline 的经典起步动作——在 WordCount 这类经典示例中把每行文本拆成单词正是通过这里练习的 OneToMany 变换完成的。完成本关后建议继续完成课程中的Map与FlatMap练习从对照中体会 Beam 变换抽象由通用到轻量的设计脉络。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Kotlin Katas 实战用 ParDo 实现 OneToMany 一对多映射Apache Beam Kotlin Katas 实战用 ParDo 实现 OneToMany 一对多映射 Apache Beam 的 ParDo 是最核心的大数据批处理流处理数据工程Apache Beam Java Katas用 ParDo 实现 OneToMany一对多映射将句子拆分为单词Apache Beam Java Katas用 ParDo 实现 OneToMany一对多映射将句子拆分为单词 导读 Apache Beam 的 Par大数据批处理流处理数据工程Apache Beam Go SDK 实战用 ParDo 实现 One-to-Many 一对多变换句子分词 Kata 详解Apache Beam Go SDK 实战用 ParDo 实现 One to Many 一对多变换句子分词 Kata 详解 Apache Beam 的 P大数据批处理流处理数据工程上一篇三步解锁中文心理咨询数据集从零构建你的AI心理助手下一篇Citra 3DS模拟器画质优化从模糊到清晰的5步配置指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表