
【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Filter 是 Apache Beam Python SDK 中用于按条件筛选PCollection元素的元素级elementwise变换给定一个返回布尔值的谓词函数它会保留所有满足条件的元素并丢弃其余元素。本文以仓库中的官方文档 filter.md 为主线结合其配套的 6 个可运行示例与底层实现源码系统讲解函数过滤、lambda 过滤、多参数过滤以及基于侧输入side inputs的三种过滤方式读完即可在真实管道中熟练运用。一、Filter 是什么Filter是 Apache Beam 中对PCollection做元素筛选的标准变换它的语义非常朴素接受一个谓词函数保留返回True的元素过滤掉其余元素。官方文档还指出它也可以基于元素自身的比较排序comparison ordering与给定值进行不等式过滤例如筛选出所有大于某阈值的数据。从实现层面看Filter并不是一个独立的DoFn而是构建在FlatMap之上的语法糖。查看源码 core.py 可以看到它的核心逻辑def Filter(fn, *args, **kwargs): # pylint: disableinvalid-name 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) ... pardo FlatMap(wrapper, *args, **kwargs) pardo.label label return pardo这段实现揭示了几点关键信息Filter的返回值必须是一个可调用对象callable否则会抛出TypeError。特别地把DoFn实例直接传给Filter会报错因为DoFn只支持用在ParDo上。内部通过一个包装 lambda 实现谓词返回True时产出[x]保留元素返回False时产出[]丢弃元素这与FlatMap的每个输入可以产出零个或多个输出语义完全吻合。变换的标签label会自动命名为Filter(函数名)例如beam.Filter(is_perennial)在流水线图上显示为Filter(is_perennial)。Filter会代理被包装函数的类型提示type hints输入类型取自已包装函数输出类型被修正为与输入类型相同确保 Beam 的类型推断系统如编码器选择正常工作。后续所有示例都基于同一组蔬菜水果produce数据每一条记录包含icon图标、name名称和duration生长周期三个字段其中duration的取值包括annual一年生、biennial两年生和perennial多年生。二、示例数据与环境准备6 个示例均可在本地直接运行也可以在 Apache Beam Playground 中在线体验。在仓库中这些示例以可测试片段的形式存放在 sdks/python/apache_beam/examples/snippets/transforms/elementwise/ 目录下每个.py文件头部都带有beam-playground元数据注解name、description、complexity、tags 等并定义了[START ...]/[END ...]标记块供文档系统提取。以filter_function.py为例基础环境只需引入apache_beamimport apache_beam as beam def filter_function(testNone): # [START filter_function] import apache_beam as beam def is_perennial(plant): return plant[duration] perennial with beam.Pipeline() as pipeline: perennials ( pipeline | Gardening plants beam.Create([ { icon: , name: Strawberry, duration: perennial }, { icon: , name: Carrot, duration: biennial }, { icon: , name: Eggplant, duration: perennial }, { icon: , name: Tomato, duration: annual }, { icon: , name: Potato, duration: perennial }, ]) | Filter perennials beam.Filter(is_perennial) | beam.Map(print)) # [END filter_function] if test: test(perennials) if __name__ __main__: filter_function()运行方式很简单直接以python filter_function.py执行各示例文件末尾都有if __name__ __main__:入口也可以在测试框架中通过传入test回调做断言验证。示例预期的输出是三行多年生植物{icon: , name: Strawberry, duration: perennial} {icon: , name: Eggplant, duration: perennial} {icon: , name: Potato, duration: perennial}这一断言逻辑与仓库测试文件 filter_test.py 中的check_perennials完全一致。三、六种 Filter 用法详解1. 使用具名函数过滤最直观的写法是定义一个具名谓词函数。is_perennial接收一个元素这里是字典返回plant[duration] perennial的布尔结果。Filter对每个元素调用该函数仅保留返回True的元素。当过滤逻辑较复杂、需要在多处复用、或希望被测试单独覆盖时优先使用具名函数。2. 使用 lambda 函数过滤对于简单的判断可以直接内联 lambda省去单独定义函数perennials ( pipeline | Gardening plants beam.Create([...]) | Filter perennials beam.Filter(lambda plant: plant[duration] perennial) | beam.Map(print))完整代码见 filter_lambda.py输出与示例 1 相同。lambda 与具名函数在语义上完全等价选择哪一种主要看可读性与复用需求。3. 传递多个参数进行过滤Filter支持向谓词函数传递额外的位置参数和关键字参数。这些参数会在调用函数时附加到元素之后参数形式与beam.Filter(has_duration, perennial)完全对应。例如def has_duration(plant, duration): return plant[duration] duration perennials ( pipeline | Gardening plants beam.Create([...]) | Filter perennials beam.Filter(has_duration, perennial) | beam.Map(print))其中perennial以位置参数的形式作为第二个实参传给has_duration。关键字参数同样支持例如beam.Filter(has_duration, durationperennial)。这种模式很适合把过滤阈值/目标值作为外部可配置参数传入让同一个谓词函数服务于不同的过滤条件。完整代码见 filter_multiple_arguments.py。注意额外参数必须是可在 worker 间序列化的值如字符串、数字、简单结构因为管道定义会被分发给分布式执行环境。4. 以单例Singleton形式使用侧输入过滤当过滤条件来自另一个PCollection且该PCollection只有一个值时例如某个上游计算的平均值可以用beam.pvalue.AsSingleton(pcollection)把它包装成单例侧输入在调用谓词时按值访问perennial pipeline | Perennial beam.Create([perennial]) perennials ( pipeline | Gardening plants beam.Create([...]) | Filter perennials beam.Filter( lambda plant, duration: plant[duration] duration, durationbeam.pvalue.AsSingleton(perennial), ) | beam.Map(print))这里durationbeam.pvalue.AsSingleton(perennial)以关键字参数形式传入侧输入AsSingleton会把PCollection中唯一的值解包出来作为duration的实参。完整代码见 filter_side_inputs_singleton.py。从源码看AsSingleton定义于 pvalue.py是AsSideInput的子类内部通过_view_options支持可选的default_value兜底。它解决了动态过滤条件来自运行时计算结果的典型需求避免把结果落盘再读回。5. 以迭代器Iterator形式使用侧输入过滤当作为过滤条件的PCollection包含多个值时应使用beam.pvalue.AsIter(pcollection)以迭代器方式传入。迭代器按需惰性访问元素因此可以遍历大到无法一次性装入内存的PCollection这是它相比AsList的核心优势valid_durations pipeline | Valid durations beam.Create([ annual, biennial, perennial, ]) valid_plants ( pipeline | Gardening plants beam.Create([ {icon: , name: Strawberry, duration: perennial}, {icon: , name: Carrot, duration: biennial}, {icon: , name: Eggplant, duration: perennial}, {icon: , name: Tomato, duration: annual}, # 注意这里特意写成大写的 PERENNIAL以演示过滤失效 {icon: , name: Potato, duration: PERENNIAL}, ]) | Filter valid plants beam.Filter( lambda plant, valid_durations: plant[duration] in valid_durations, valid_durationsbeam.pvalue.AsIter(valid_durations), ) | beam.Map(print))值得注意的细节本示例的蔬菜数据中Potato的duration被刻意写成大写PERENNIAL因此它不会命中valid_durations中的perennial最终输出只有 4 条合法植物对应测试文件 filter_test.py 中的check_valid_plants期望。这提醒我们字符串过滤是精确匹配大小写敏感若需要大小写不敏感应提前做归一化如统一.lower()。完整代码见 filter_side_inputs_iter.py。官方文档特别提示你也可以用beam.pvalue.AsList(pcollection)把侧输入整体转成列表但这要求该PCollection的所有元素都能装进内存。6. 以字典Dictionary形式使用侧输入过滤如果侧输入PCollection足够小、可以整体载入内存且每个元素都是(key, value)键值对就可以用beam.pvalue.AsDict(pcollection)以字典方式访问——直接用 key 做 O(1) 查询keep_duration pipeline | Duration filters beam.Create([ (annual, False), (biennial, False), (perennial, True), ]) perennials ( pipeline | Gardening plants beam.Create([...]) | Filter plants by duration beam.Filter( lambda plant, keep_duration: keep_duration[plant[duration]], keep_durationbeam.pvalue.AsDict(keep_duration), ) | beam.Map(print))这里keep_duration[plant[duration]]直接按植物周期取值perennial对应True则保留annual/biennial对应False则丢弃实现了动态过滤策略表的效果。完整代码见 filter_side_inputs_dict.py。字典方式的适用前提所有元素必须能装入内存且每个元素必须是(key, value)二元组。如果PCollection太大装不下官方文档明确建议改用beam.pvalue.AsIter(pcollection)。四、侧输入三种形态的选型速查结合官方文档与源码可以归纳出三种侧输入形态的适用场景侧输入形态用法适用场景内存要求单例beam.pvalue.AsSingleton(pcoll)过滤条件只有一个值如均值、单值配置单值迭代器beam.pvalue.AsIter(pcoll)过滤条件为多值集合元素逐个惰性访问无需整体载入内存字典beam.pvalue.AsDict(pcoll)过滤条件为(key, value)映射按键查询必须全部装入内存另外beam.pvalue.AsList(pcoll)也能以列表形式传入侧输入但正如文档强调的它要求所有元素同时驻留内存因此不适合超大PCollection。实践中多值场景优先AsIter需要按键查询且数据量可控时用AsDict。五、如何验证与测试 Filter仓库为每个示例都配套了单元测试见 filter_test.py。测试通过mock.patch将beam.Pipeline替换为TestPipeline、将示例中的print替换为收集函数然后对每个示例函数传入check_perennials或check_valid_plants回调做输出断言mock.patch(apache_beam.Pipeline, TestPipeline) class FilterTest(unittest.TestCase): def test_filter_function(self): filter_function.filter_function(check_perennials) def test_filter_lambda(self): filter_lambda.filter_lambda(check_perennials) def test_filter_multiple_arguments(self): filter_multiple_arguments.filter_multiple_arguments(check_perennials) def test_filter_side_inputs_singleton(self): filter_side_inputs_singleton.filter_side_inputs_singleton(check_perennials) def test_filter_side_inputs_iter(self): filter_side_inputs_iter.filter_side_inputs_iter(check_valid_plants) def test_filter_side_inputs_dict(self): filter_side_inputs_dict.filter_side_inputs_dict(check_perennials)在 examples/snippets 目录下执行对应的测试命令即可验证这 6 种写法。如果你想在自己的管道中复用这套模式可以仿照示例把业务函数写成接收test回调的形式便于在TestPipeline中做输出断言。六、Filter 与相关变换的关系官方文档在末尾列出了两个紧密相关的元素级变换FlatMap行为与Map相同但每个输入元素可以产生零个或多个输出。从前文源码可知Filter本质就是FlatMap(wrapper, ...)的特例——谓词为真时产出一个元素为假时产出零个元素。ParDo最通用的元素级映射变换支持多输出集合TaggedOutput、侧输入等更复杂的能力。当过滤逻辑伴随额外副作用、需要按窗口或键做精细控制时可直接用ParDo配合条件分支实现。选型建议纯保留/丢弃判断优先用Filter代码最简洁需要每个输入产生多条输出时用FlatMap需要多路输出、按时间戳处理或更细粒度的 DoFn 生命周期控制时升级到ParDo。七、小结本文围绕官方文档 filter.md 完整讲解了 Apache Beam PythonFilter变换的 6 种实战写法具名函数、lambda、多参数、单例侧输入、迭代器侧输入、字典侧输入。同时结合 core.py 的源码揭示了Filter基于FlatMap的实现原理、callable 校验与类型提示代理机制并用 filter_test.py 给出了可复用的验证范式。掌握这些写法后你可以在批处理与流式管道中灵活实现数据清洗、异常值剔除、基于动态配置的过滤等常见需求。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Kotlin Kata 实战使用 Filter 变换过滤 PCollection 中的元素Apache Beam Kotlin Kata 实战使用 Filter 变换过滤 PCollection 中的元素 导读 本文围绕 Apache Beam 官批处理流处理大数据Apache Beam Java SDK Filter 转换实战用 Filter.by 谓词过滤 PCollection 元素Apache Beam Java SDK Filter 转换实战用 Filter.by 谓词过滤 PCollection 元素 Apache Beam 的 J批处理流处理大数据Apache Beam Go SDK 实战使用 filter 包Include/Exclude过滤 PCollection 元素Apache Beam Go SDK 实战使用 filter 包Include/Exclude过滤 PCollection 元素 导读 本文聚焦 Apac上一篇3分钟掌握BBDown高效命令行B站视频下载解决方案下一篇PaddleX 产线全景指南CPU/GPU 基础产线与特色产线配置解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考