Apache Beam Python Filter 变换详解:6 种过滤 PCollection 元素的实战方案
【免费下载链接】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),仅供参考

相关新闻

我把 OpenRouter 换成了 ofox:0 加价 + 双协议 + Claude Code 直连实测(2026)

我把 OpenRouter 换成了 ofox:0 加价 + 双协议 + Claude Code 直连实测(2026)

上个月我们团队的 API 月账单突破 $400,我拉了一下 OpenRouter 的消费明细,发现光 5.5% 的平台手续费就吃掉了 $22。一开始是拒绝折腾的——OpenRouter 用了快一年,模型切换确实方便。但 $22 一个月、一年 $264,这钱花得有点冤。T…

2026/10/12 5:28:13 阅读更多 →
广州违法辞退赔付律师怎么找 2N赔付口径与律所适配参考

广州违法辞退赔付律师怎么找 2N赔付口径与律所适配参考

在广州主张违法辞退赔付,先要弄清 2N、N、半个月工资这几个口径分别对应什么情形:2N 是违法解除劳动合同的赔付金标准,N 是经济补偿,不满六个月对应的是半个月工资,三者不能混用。搞清口径后再选律师,综合公…

2026/10/12 5:28:13 阅读更多 →
一句话造Agent:PenguinHarness 0.2.1声明式编排实战

一句话造Agent:PenguinHarness 0.2.1声明式编排实战

1. 从“手搓 Agent”到“一句话造 Agent”的认知转变1.1 为什么“手搓 Agent”正在成为过去式如果你在过去一年里尝试过搭建一个能自主完成任务的智能体,大概率经历过这样的场景:打开编辑器,先写一个循环,再定义工具调用格式&…

2026/10/12 5:28:13 阅读更多 →

最新新闻

有害气体控制洁净工程的底层逻辑:从过滤到吸附,从压差到监测

有害气体控制洁净工程的底层逻辑:从过滤到吸附,从压差到监测

空气里最危险的不是脏,而是失控:有害气体控制洁净工程的底层逻辑干了这么多年洁净工程,我越来越觉得“洁净”这个词会误导人。很多人一听到洁净室,想到的就是无尘、高等级过滤、白大褂和干干净净的地板,下意识把“颗粒…

2026/10/12 6:03:33 阅读更多 →
Pygame乒乓球游戏开发实战:从游戏循环到碰撞检测的完整指南

Pygame乒乓球游戏开发实战:从游戏循环到碰撞检测的完整指南

简介:游戏循环是几乎所有实时游戏的心跳,它决定了每一帧里输入、更新与渲染的执行顺序。碰撞检测则负责回答“物体是否重叠”这个基本问题,而引擎中那些微妙的物理手感,往往源于对碰撞响应和状态管理的精细控制。乒乓球游戏恰好是…

2026/10/12 6:03:33 阅读更多 →
栈的压入、弹出序列判定算法详解:辅助栈模拟与 Java 实现(YCBlogs 剑指 Offer 系列)

栈的压入、弹出序列判定算法详解:辅助栈模拟与 Java 实现(YCBlogs 剑指 Offer 系列)

教程技术博客文档 【免费下载链接】YCBlogs 技术博客笔记大汇总,包括Java基础,线程,并发,数据结构;Android技术博客等等;常用设计模式;常见的算法;网络协议知识点;部分fl…

2026/10/12 6:03:33 阅读更多 →
C# TCP服务器与客户端双向通信:骨架搭建与避坑指南

C# TCP服务器与客户端双向通信:骨架搭建与避坑指南

简介:这是一份面向C#网络编程初学者的TCP通信示例工程,目标是用一个程序实现TCP客户端与服务器之间的互发消息,并支持在客户端界面点击按钮弹出服务器界面。资源围绕System.Net命名空间下的TcpListener与TcpClient展开,覆盖端口绑…

2026/10/12 6:03:32 阅读更多 →
CodeIgniter 4.7.4 安全与稳定性更新详解:四个安全公告与十余项缺陷修复

CodeIgniter 4.7.4 安全与稳定性更新详解:四个安全公告与十余项缺陷修复

后端Web框架 【免费下载链接】CodeIgniter4 Open Source PHP Framework (originally from EllisLab) 项目地址: https://gitcode.com/gh_mirrors/co/CodeIgniter4 点击查看 免费下载 CodeIgniter 4.7.4(2026 年 7 月 7 日发布)是一次以安全加…

2026/10/12 6:03:32 阅读更多 →
Tortoise-ORM 与 Sanic 集成实战:register_tortoise 生命周期管理全解析

Tortoise-ORM 与 Sanic 集成实战:register_tortoise 生命周期管理全解析

数据库后端 【免费下载链接】tortoise-orm Familiar asyncio ORM for python, built with relations in mind 项目地址: https://gitcode.com/gh_mirrors/to/tortoise-orm 点击查看 免费下载 本文以 Tortoise-ORM 仓库中 Sanic 集成示例 为主线,系统讲解…

2026/10/12 6:02:32 阅读更多 →

日新闻

复古胶片颗粒感噪点合成器:Canvas ImageData 像素高斯杂色注入算法

复古胶片颗粒感噪点合成器:Canvas ImageData 像素高斯杂色注入算法

在数码相机、高清显示屏与现代矢量图形技术高度发达的今天,画面可以做到绝对的锐利、平滑与无瑕。然而,当一张秋日手账插画或拍立得照片过于“平整无瑕”时,往往会散发出一种冰冷生硬的“数码塑料感(Digital Plasticity&#xff0…

2026/10/12 0:00:59 阅读更多 →
活字印刷古籍线装排版:Canvas 竖排文字与栏线自适应算法

活字印刷古籍线装排版:Canvas 竖排文字与栏线自适应算法

在现代网页与移动端设计中,横排(Horizontal Layout)早已经成为了绝对的主流。然而,当我们翻开泛黄的线装古籍、宋版木刻诗集,或是欣赏一张茶道雅集的手写便签时,那种**自上而下纵向书写、自右向左逐列铺展&…

2026/10/12 0:00:59 阅读更多 →
周日晚间的“精神松绑减震器”:无压力情绪倾倒箱与温和轻声陪伴

周日晚间的“精神松绑减震器”:无压力情绪倾倒箱与温和轻声陪伴

每到周日的晚上八点到十点,很多人心里都会悄悄亮起一盏警示灯。 在心理学上,这种现象有一个专门的称谓——“周日夜晚焦虑症(Sunday Scaries)”。明天又是周一,闹钟又要重新在七点响彻卧房;脑海里仿佛有一个…

2026/10/12 0:00:59 阅读更多 →

周新闻

流感时间序列预测实战:ARIMA/LSTM全流程拆解与避坑指南

流感时间序列预测实战:ARIMA/LSTM全流程拆解与避坑指南

简介:基于 ARIMA、LSTM、Transformer 等模型的流感时间序列预测 Python 源码,面向计算机相关专业课程设计与期末大作业学生,以及项目实战学习者。内容覆盖预处理、平稳性检验、定阶、残差分析、多模型对比预测的完整时序建模流程,…

2026/10/12 0:16:30 阅读更多 →
影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别

影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别

影刀RPA新手教程:键盘模拟输入实战——输入文本与模拟按键的区别 做影刀RPA自动化,十个新手有八个栽在"往输入框里填东西"这件事上:要么填不进去,要么填了一半,要么直接把原来内容追加在后面。这背后的根因&…

2026/10/12 0:16:38 阅读更多 →
影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容

影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容

影刀RPA新手教程:阅文起点小说数据采集实战——书籍信息与章节内容 1. 认识影刀:什么场景该用RPA采小说数据 起点中文网的页面结构相对稳定——分类榜单、书籍详情、章节内容三块独立页面,跳转链路清晰。这种场景非常适合影刀自动化&#x…

2026/10/12 0:16:43 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/11 10:45:37 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/11 14:36:53 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/11 14:36:54 阅读更多 →