资讯考证相关

Apache Beam Python Filter 变换详解:6 种过滤 PCollection 元素的实战方案

2026/10/12 5:27:47 安证通 考证咨询 特种作业
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),仅供参考
本文仅供参考,具体政策以官方公告为准 返回资讯列表 →
延伸阅读

更多相关内容

相关资讯、最新动态、本周本月更新,都在这里。

特种工证能提前退休吗?跨省换证+考前押题全攻略

特种工证能提前退休吗?跨省换证+考前押题全攻略

特种工证能提前退休吗?跨省换证+考前押题全攻略 之前在外省考的电工证,现在回三门峡想接活,能不能直接转过来?这是很多劳务班组负责人和持证工人最头疼的问题。很多人手里攥着证,却卡在“异地互认”的环节,甚至因为不了解 考前押题 后的最新题库变化,导致复审或换证时手忙脚乱。其实,关于…

查看 →
Java基础进阶:面向对象、集合框架、异常处理与泛型全梳理

Java基础进阶:面向对象、集合框架、异常处理与泛型全梳理

这是Java总结进阶之路系列的第二篇。写这篇的起因很简单:很多朋友学完基础语法之后会卡在一个不上不下的位置——变量、数组、循环、方法都会写,但一旦看到的代码开始出现类继承、集合框架、异常捕获这些内容,整个人就开始发懵。基础一解决的…

查看 →
长春电工证在哪复审?避开中介坑的3个官方渠道

长春电工证在哪复审?避开中介坑的3个官方渠道

长春电工证在哪复审?避开中介坑的3个官方渠道 不知道去哪报名怕被中介坑?别急,长春电工证复审其实有明确路径。核心在于搞清 在哪考 以及怎么避免花冤枉钱。很多老电工因为不懂流程,每年多花几百块甚至被忽悠去外地,其实本地就能办。 长春电工证复审具体去哪找官方入口?…

查看 →
嘉兴电工证去哪里办理?实操不过白跑,看这3步拿电子证书

嘉兴电工证去哪里办理?实操不过白跑,看这3步拿电子证书

嘉兴电工证去哪里办理?实操不过白跑,看这3步拿电子证书 很多兄弟在嘉兴准备考电工证,最纠结的不是报名费,而是 实操考试心里没底怕挂科 。毕竟理论背背书还能混过去,实操那一套闸刀、验电笔、绝缘测试,手一抖直接零分,重考还得等下个月,这时间成本太吓人。更让人头疼的是,很多人跑了好几个地方,不知道…

查看 →
AI技能包实战:从提示词到可复用能力单元

AI技能包实战:从提示词到可复用能力单元

技能(skills)这个词,前两年说出来大家想到的可能是简历上的Office三件套、会点Python、能剪视频。现在你要是还只在简历意义上理解它,就真的落伍了。这一两年AI圈把“技能”玩出了新花样——给模型装上“技能包”,让它…

查看 →
2023电工证培训:工地忙到飞起,靠考前押题拿证真不难

2023电工证培训:工地忙到飞起,靠考前押题拿证真不难

2023电工证培训:工地忙到飞起,靠考前押题拿证真不难 工地现场天天催进度,脑子转得比发电机还快,谁还有心思啃那厚厚一本《电工作业》教材?很多老铁跟我吐槽,白天盯现场、晚上对账,根本没时间复习,怕考试挂了白交钱。别慌,我是干这行十年的老顾问,专门给咱们这种“没空学习”的工地人支招。其实2023年的电…

查看 →
AI转头就忘?claude-mem给Claude配了个跨会话记忆库

AI转头就忘?claude-mem给Claude配了个跨会话记忆库

说个真实场景:昨天还在让Claude帮忙梳理技术方案,今天想开个新会话补问一句“方案里第三点结论是什么”,它却一脸茫然地反问“你是指哪份方案?”——你只能把背景从头再喂一遍。这个体验其实不是Claude变笨了,而是每个…

查看 →
三门峡高压电工证发放全解析:避开中介坑,拿全国通用证

三门峡高压电工证发放全解析:避开中介坑,拿全国通用证

三门峡高压电工证发放全解析:避开中介坑,拿全国通用证 在三门峡干电力的,谁没被中介的“包过”“快速拿证”忽悠过?心里没底,怕钱打水漂更怕证办下来是假的,这种焦虑太真实了。很多人不知道去哪报名,结果钱花了,证没考下来,或者考下来发现根本没法在国网或大型项目上注册。 别急,今天咱们就掰开揉碎了讲清楚…

查看 →
电工证7天出证企业靠谱吗?良心建议防坑

电工证7天出证企业靠谱吗?良心建议防坑

电工证7天出证企业靠谱吗?良心建议防坑 实操考试心里没底怕挂科?别慌,这行水太深,今天说点大实话。 很多人一听说“电工证7天出证企业”,第一反应是:这能信?是不是那种交钱就办事的野路子? 其实,这里有个巨大的认知误区。所谓的“7天出证”,指的绝不是从你报名到拿证只花7天,而是指…

查看 →
内部消息揭秘:考低压电工证有难度?郑州金水区水利人避坑指南

内部消息揭秘:考低压电工证有难度?郑州金水区水利人避坑指南

内部消息揭秘:考低压电工证有难度?郑州金水区水利人避坑指南 是不是每次看到朋友圈里有人晒电工证,心里痒痒的,想考一个备用,又怕被那些满大街发广告的中介坑得明明白白?这种“不知道去哪报名、怕交钱后没下文”的焦虑,我太懂了。在郑州金水区搞水利工程的,咱们干的是实打实的辛苦活,证书这东西,水太深,稍不留神…

查看 →
内部消息:涿州市电工证培训跨省转办全流程解析

内部消息:涿州市电工证培训跨省转办全流程解析

内部消息:涿州市电工证培训跨省转办全流程解析 很多在外地干了几年电工的老伙计,最近都在问同一个问题:之前考的那个电工证,现在回涿州干活,还能不能用?或者说,需不需要重新在涿州市找个地方重新培训一遍?这确实是咱们一线从业者最头疼的事,毕竟谁也不想白交那几百上千的学费,更不想耽误接活儿的时间。今天我就结…

查看 →
3招看懂中级电工证能干嘛用,避开官方报名入口陷阱

3招看懂中级电工证能干嘛用,避开官方报名入口陷阱

3招看懂中级电工证能干嘛用,避开官方报名入口陷阱 很多人怕考不过白交培训费,这种焦虑我太懂了。毕竟电工证是安监局(现应急管理局)发的硬通货,考不过不仅钱打水漂,还耽误找工作的时间。 别慌,今天咱们不聊虚的,直接拆解【中级电工证能干嘛用】,顺便把 官方报名入口…

查看 →
官方回应:怎么申领电工证电子证?异地转证全攻略

官方回应:怎么申领电工证电子证?异地转证全攻略

官方回应:怎么申领电工证电子证?异地转证全攻略 之前考的证在外省能不能转过来?这是无数在外打拼的电工兄弟最关心的事儿。别急,今天就把这事儿掰开了揉碎了讲清楚。 很多师傅拿着纸质本,担心换个城市工地就废了。其实, 国家安全生产考试网…

查看 →
箱面接地电工作业报名门槛低,考下来多少钱才合理

箱面接地电工作业报名门槛低,考下来多少钱才合理

箱面接地电工作业报名门槛低,考下来多少钱才合理 学历不高,怕报不上名?别自己瞎琢磨,直接看这里。 很多刚出校门的朋友,尤其是许昌这边刚毕业学工程的学弟学妹,手里攥着身份证,心里打鼓:我是不是得先拿个本科证才能考这个证?是不是得去报个几千块的培训班才能过?其实,关于 箱面接地电工作业…

查看 →
建湖高压电工证外省能转吗?官方回应来了

建湖高压电工证外省能转吗?官方回应来了

建湖高压电工证外省能转吗?官方回应来了 之前在外省考的高压电工证,回到江苏建湖能直接用吗?能不能直接转过来?这是很多跨省流动电工最头疼的问题。别急,关于【建湖高压电工证】异地互认的疑惑, 官方回应…

查看 →
相关服务

看完文章,下一步可以直接办

报考、备考、复审相关的服务入口,都在这里。

考试批次时间

近期各工种批次安排与报名截止提醒。

查看详情 →

报考条件查询

年龄、学历、体检条件逐项对照。

查看详情 →

材料免费预审

报名材料逐项核对,缺什么当场补齐。

查看详情 →

复审流程

复审时间、材料与流程一次说清。

查看详情 →
报名流程

从咨询到拿证,就四步

每一步都有明确产出,每一步都有人盯着。

01

意向沟通

说清岗位与目标,顾问推荐对应工种与报考方向。

02

材料预审

身份证、学历、体检逐项核对,缺什么当场补齐。

03

批次报名

锁定最近考试批次,考务信息逐一确认。

04

培训考试

题库辅导加实操要点,考完节点逐一跟进拿证。

免费咨询

想报考特种作业证?找顾问聊一聊

根据你的工作经历推荐工种,确认批次与材料,30 秒登记当天回访。