批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载本篇技术指南以 Apache Beam 官方 Katas 课程中Common Transforms / Filter一课task.md为核心骨架系统讲解 Beam Python SDK 中Filter变换的用法如何用beam.Filter配合 lambda 表达式过滤奇数、保留偶数并结合仓库中的core.py源码揭示其底层实现原理与测试验证方法。读完本文你将掌握Filter的两种等价实现方式beam.Filter与ParDo、参数传递技巧、类型提示type hints的推导逻辑以及如何独立运行与验证 Beam Python 数据过滤管道。一、Kata 任务概览过滤奇数保留偶数Apache Beam 的官方 Katas 是一系列面向初学者的编码练习位于仓库 learning/katas 目录下按语言Java、Kotlin、Python、Go与主题Common Transforms、Core Transforms、Windowing、Triggers 等组织。本文对应的任务位于 learning/katas/python/Common Transforms/Filter/Filter/task.md属于Common Transforms常用变换课程中Filter一课课程结构由 lesson-info.yaml 定义包含 ParDo 与 Filter 两个子任务。任务描述原文Kata:Implement a filter function that filters out the odd numbers by using Filter.Hint:Use Filter with a lambda.翻译过来即实现一个过滤函数利用Filter变换把 1~10 中的奇数1、3、5、7、9过滤掉只保留偶数2、4、6、8、10。官方提示给出两条关键线索使用apache_beam.transforms.core.Filter即beam.Filter配合lambda匿名函数编写过滤谓词。task.md还点出了本课的背景主旨The Beam SDKs provide language-specific ways to simplify how you provide your DoFn implementation.——即 Beam 各语言 SDK 提供了一系列简化 DoFn 实现的专用变换如Map、FlatMap、Filter让你不必每次都手动编写完整的DoFn类。二、参考解法一行 lambda 完成过滤本任务的参考解法位于 learning/katas/python/Common Transforms/Filter/Filter/task.py完整代码如下import apache_beam as beam with beam.Pipeline() as p: (p | beam.Create(range(1, 11)) | beam.Filter(lambda num: num % 2 0) | beam.LogElements())这段代码虽然简短却完整演示了 Beam Python 管道Pipeline的标准骨架逐行拆解如下代码片段作用import apache_beam as beam导入 Beam Python SDK统一以beam命名空间访问Pipeline、Create、Filter、LogElements等核心构件with beam.Pipeline() as p:以上下文管理器方式创建管道p。with块结束时会自动运行管道触发__exit__执行 run无需手动调用run()p \| beam.Create(range(1, 11))beam.Create把 Python 的range(1, 11)即整数 1~10转换为一个内存型PCollection作为管道的输入源\| beam.Filter(lambda num: num % 2 0)核心过滤步骤Filter对每个元素调用谓词函数保留使谓词返回True的元素丢弃返回False的元素。num % 2 0判定偶数的条件成立因此只保留偶数\| beam.LogElements()把PCollection中的每个元素打印到日志/标准输出方便本地验证过滤结果运行该管道控制台将输出2 4 6 8 10奇数 1、3、5、7、9 全部被过滤掉。三、Filter 的底层实现本质是一个包装过的 FlatMapbeam.Filter并非一个独立的底层算子而是一个基于FlatMap的语法糖封装。查看其定义所在文件 sdks/python/apache_beam/transforms/core.py源码实现如下def Filter(fn, *args, **kwargs): # pylint: disableinvalid-name :func:Filter is a :func:FlatMap with its callable filtering out elements. Filter accepts a function that keeps elements that return True, and filters out the remaining elements. if not callable(fn): raise TypeError( Filter can be used only with callable objects. Received %r instead. % (fn)) wrapper lambda x, *args, **kwargs: [x] if fn(x, *args, **kwargs) else [] label Filter(%s) % ptransform.label_from_callable(fn) # TODO: What about callable classes? if hasattr(fn, __name__): wrapper.__name__ fn.__name__ # Get type hints from this instance or the callable. Do not use output type # hints from the callable (which should be bool if set). fn_type_hints typehints.decorators.IOTypeHints.from_callable(fn) ... return ptransform(pardo.DoFn(*args, **kwargs), label, wrapper)从源码可以提炼出三个关键机制谓词语义Predicate SemanticsFilter接受的fn必须是可调用对象callable其第一个参数是待判断的元素。源码用包装函数wrapper lambda x, *args, **kwargs: [x] if fn(x, *args, **kwargs) else []实现语义——谓词返回True时包装函数输出包含该元素的单元素列表[x]返回False时输出空列表[]。这正是FlatMap的一个元素可以产生 0 到多个输出语义因此Filter天然支持对每个输入元素保留或不保留。严格的入参校验Filter第一行便检查callable(fn)若非可调用对象例如误传一个DoFn实例——DoFn实例只支持用于ParDo会抛出TypeError错误信息为 Filter can be used only with callable objects. Received %r instead.。这一点在 core.py 的 docstring 与实现中均有明确说明。自动生成标签与类型推导Filter会自动生成形如Filter(fn_name)的变换标签便于在管道图中识别同时通过typehints.decorators.IOTypeHints.from_callable(fn)提取被包装函数的类型提示并将输出类型设置为与输入类型一致with_output_types(typehints.Iterable[...])包装保证过滤不改变元素类型这一直观语义。从源码结构看Filter最终调用ptransform(pardo.DoFn(*args, **kwargs), label, wrapper)把包装函数包进一个DoFn交给ParDo执行因此它在运行时与ParDo走的是同一条执行链路只是编程模型上更简洁。Filter 的参数签名Filter的函数签名来自 core.py为Filter(fn, *args, **kwargs)各参数含义如下参数类型说明fnCallable[..., bool]过滤谓词。第一个参数是待判断的元素必须返回布尔值或可转换为布尔值必须是可调用对象否则抛TypeError*args-传递给谓词函数的额外位置参数**kwargs-传递给谓词函数的关键字参数*args/**kwargs的存在意味着谓词可以携带额外参数。例如# 保留大于 threshold 的元素 threshold 5 filtered pcoll | beam.Filter(lambda num, t: num t, threshold) # 使用关键字参数 filtered pcoll | beam.Filter(lambda num, t0: num t, tthreshold)四、等价实现用 ParDo 手写过滤逻辑Filter虽然简洁但它的底层执行单元仍然是DoFn。Katas 课程在同一课程目录下安排了姊妹任务 learning/katas/python/Common Transforms/Filter/ParDo/task.mdFilter using ParDo用于展示使用原始ParDo实现相同过滤逻辑的写法——任务要求过滤掉偶数、保留奇数参考解法在 ParDo/task.pyimport apache_beam as beam class FilterOutEvenNumber(beam.DoFn): def process(self, element): if element % 2 1: yield element with beam.Pipeline() as p: (p | beam.Create(range(1, 11)) | beam.ParDo(FilterOutEvenNumber()) | beam.LogElements())这段代码的运行输出为奇数1 3 5 7 9对比两种写法可以清晰看到Filter的简化 DoFn 实现价值对比维度beam.Filter lambdabeam.ParDo 自定义 DoFn谓词逻辑lambda num: num % 2 0一行搞定需定义class FilterOutEvenNumber(beam.DoFn)并覆写process方法保留元素谓词返回True即保留在process中对该元素执行yield element才会输出丢弃元素谓词返回False自动丢弃process中不yield该元素即可也可return或yield from []适用场景简单、单条件过滤需要setup/teardown、start_bundle/finish_bundle生命周期钩子或复杂多输出逻辑值得说明的是core.py 的注释明确指出Map、FlatMap、Filter这类简化变换适用于不需要start_bundle/finish_bundle或setup/teardown生命周期方法的场景一旦你需要这些生命周期控制就应当退回到完整的ParDoDoFn写法。五、测试验证如何确认过滤结果正确Katas 的每个任务都配有自动化测试用于在教学环境PyCharm Education / EduTools中校验你的解法。本任务的测试文件位于 learning/katas/python/Common Transforms/Filter/Filter/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 [2, 4, 6, 8, 10] for num in answers: self.assertIn(num, output, Incorrect output. Filter out the odd numbers.)测试逻辑非常直白可以拆解为两层test_not_empty校验task.py文件非空防止交空文件使用 test_helper.py 中的test_is_not_empty()实现test_output通过get_file_output(pathtask.py)实际执行task.py内部用subprocess.Popen([sys.executable, path], ...)运行并把标准输出逐行读取为字符串列表然后断言偶数2、4、6、8、10全部出现在输出中。如果缺少任何一个偶数测试失败并提示 Incorrect output. Filter out the odd numbers.。姊妹任务 ParDo 版的测试ParDo/tests/test_task.py结构完全相同只是期望答案为奇数1、3、5、7、9。这套测试本身也是学习素材它展示了 Beam 管道在本地环境默认 Direct Runner下直接运行、并把LogElements输出到标准输出的行为是验证任何小型 Beam Python 管道的最简单方式。六、本地运行环境准备Filter练习的代码不依赖任何外部运行器使用 Beam 内置的 Direct Runner 即可本地执行。参考仓库中 learning/katas/python/README.md 的指引标准做法是使用PyCharm Education或安装了 EduTools 插件的 PyCharm以 Create New Project 方式打开 learning/katas/python 目录为项目选择 Python 解释器例如 virtualenv并创建虚拟环境等待索引完成后在 Project 工具窗口切换到 Course 视图即可逐个完成任务并运行测试项目依赖见 learning/katas/python/requirements.txt确保环境中已安装apache_beamPython 包。如果不使用 IDE也可以直接命令行运行参考解法验证结果cd learning/katas/python/Common Transforms/Filter/Filter python task.py只要环境中已安装apache_beam即可在标准输出看到过滤后的偶数序列与测试断言完全一致。七、进阶练习把 Filter 用在更真实的场景掌握了beam.Filter后可以把同样的技巧迁移到真实数据处理场景。下面给出三个与 Katas 任务同构的进阶示例可直接替换range(1, 11)的输入示例 1过滤字符串列表import apache_beam as beam with beam.Pipeline() as p: (p | beam.Create([apple, banana, avocado, cherry]) | beam.Filter(lambda word: word.startswith(a)) | beam.LogElements()) # 输出apple、avocado示例 2过滤自定义对象结合属性访问import apache_beam as beam class Order: def __init__(self, amount): self.amount amount with beam.Pipeline() as p: (p | beam.Create([Order(80), Order(150), Order(220)]) | beam.Filter(lambda order: order.amount 100) | beam.LogElements()) # 输出amount 为 150、220 的订单示例 3谓词携带外部参数演示*args/**kwargsimport apache_beam as beam with beam.Pipeline() as p: (p | beam.Create(range(1, 21)) | beam.Filter(lambda num, divisor: num % divisor 0, 3) | beam.LogElements()) # 输出3、6、9、12、15、18即 1~20 中能被 3 整除的数这些例子与 Katas 任务的核心知识点完全一致——Filter的谓词语义、lambda 用法、以及额外的参数传递能力可以平滑迁移到日志清洗、异常值剔除、条件标记等常见过滤需求。八、小结与自检清单本文围绕 learning/katas/python/Common Transforms/Filter/Filter/task.md 展开核心结论归纳如下beam.Filter(fn, *args, **kwargs)是 Beam Python SDK 提供的过滤变换保留谓词返回True的元素丢弃其余元素其底层实现在 sdks/python/apache_beam/transforms/core.py本质是基于FlatMap包装的ParDo语法糖并带有可调用性校验、自动标签与类型提示推导简单过滤场景优先用beam.Filter lambda需要生命周期钩子或复杂逻辑时退回到ParDo 自定义DoFn参考 ParDo/task.py正确性可通过 tests/test_task.py 自动化测试验证测试会实际执行task.py并断言偶数序列2, 4, 6, 8, 10全部出现在输出中。你可以用下面这份自检清单确认自己已掌握本课内容能用beam.Filter lambda 过滤出偶数能说出Filter保留/丢弃元素的判定规则谓词返回True保留能解释Filter与FlatMap、ParDo的关系能写出ParDoDoFn.process版本的等价过滤逻辑知道谓词函数如何接收额外参数*args/**kwargs能在本地直接运行task.py并对照测试断言验证输出。完成以上练习后建议继续学习 Katas 课程中 Core Transforms 的其它常用变换如Map、FlatMap、GroupByKey、CombinePerKey见 learning/katas/python 目录它们与Filter共同构成 Beam Python 管道编写的基础工具箱。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Python Filter 变换实战从 Kata 任务到源码级原理Apache Beam Python Filter 变换实战从 Kata 任务到源码级原理 导读 本文围绕 Apache Beam Python SDK 中Apache Beam Python SDK Filter 变换实战从 Kata 过滤奇数任务到源码级原理剖析Apache Beam Python SDK Filter 变换实战从 Kata 过滤奇数任务到源码级原理剖析 Apache Beam 的 Python SD大数据批处理流处理数据工程使用 Apache Beam Kotlin SDK 的 Min 聚合变换从 Aggregation Kata 到源码级原理剖析使用 Apache Beam Kotlin SDK 的 Min 聚合变换从 Aggregation Kata 到源码级原理剖析 Apache Beam 的 M批处理流处理大数据上一篇HunterPie完整教程5分钟掌握《怪物猎人世界》最强游戏覆盖层工具下一篇Windows 11任务栏终极自定义指南用Taskbar11释放你的桌面生产力创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?