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

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),仅供参考

看完这篇,下一步怎么走

如果你正打算考证,先看报考条件判断自己符不符合,再按报名流程详解准备材料;不知道考哪个工种的,翻工种总目录;证快到期的,留意证书复审要求。拿不准的,直接打 18236992212。

关于这张证,你还要知道

证是全国通用的吗

应急管理部门发的特种作业操作证全国通用,跨省从业有效。换工作到外地,证不用重考,到期在当地办复审即可。

多久能考下来

正常情况从报名到拿证一个多月:材料预审三五天、等批次一到两周、辅导几天、考后等制证。具体看当月批次。

考不过怎么办

理论或实操单科没过,保留成绩约补考,不用全部重来。哪科弱我们辅导时重点补哪科。

这篇文章讲的是通用情况。你自己的条件符不符合、最近一批还能不能报,电话里一句就清楚:18236992212(微信同号),邮箱 809451989@qq.com。

FAQ

关于考证,电话里最常问的

零基础能考吗?

能。培训从零教,按当年大纲走。建议走"培训+考试"全包,别只买报名名额自己硬考。

多久能拿证?

正常一个多月。材料不卡壳、批次不延误的前提下,从报名到拿证一个多月;制证还要两三周。

证是全国通用的吗?

是。应急管理部门发的特种作业操作证全国通用,换工作到外地不用重考。

证过期了怎么办?

超期未复审的证失效,一般要重新考试。拿不准的把证号发过来查系统状态。

能包过吗?

我们不承诺包过。能做到的是按大纲辅导、训练覆盖考核点。

几个人一起报便宜吗?

企业团报走单独方案,看团报说明。个人三五个的也能凑一批。

HOW IT WORKS

从咨询到拿证,大概这么走

01

电话咨询

说你的工种、情况,我们判断条件、报费用、说批次。

02

材料预审

拍照发来,逐条核规格,缺的补、糊的重拍。

03

报名约考

按批次录系统、约考位,约好时间通知你。

04

考前辅导

按你时间排理论刷题和实操训练。

05

考试取证

到场、考试、等成绩,过了等制证。

看完这篇,下一步怎么办

如果你是来查考试通知的:对一下文章里的日期和截止时间,材料准备齐了打 18236992212 预约。

如果你是来看政策的:把你的工种、证号、到期时间说清楚,我们判断新规对你有没有影响。

如果你是来了解行业的:想考证的翻工种目录,想看流程的翻报名流程。

每篇文章详情页右侧栏会推荐相关、最新和最近一周/一日/一月的文章,不用来回翻列表。

文章里的图片和日期都是发布时的信息。考试安排以最新通知为准,政策条款以官方原文为准。这页只是帮你省时间,不是替代你打电话确认。

有任何拿不准的地方,直接拨 18236992212。接电话的人会按你的具体情况告诉你下一步,不用你对着文章猜。

这篇文章,你能用来干什么

如果你是来查考试通知的:对一下文章里的日期和截止时间,材料准备齐了打 18236992212 预约。别等通知快截止了才来。

如果你是来看政策的:把你的工种、证号、到期时间说清楚,我们判断新规对你有没有影响。别自己对着原文猜。

如果你是来了解行业的:想考证的翻工种目录挑方向,想看流程的翻报名流程。

每篇文章详情页右侧栏会推荐相关、最新和最近一周/一日/一月的文章,不用来回翻列表。

文章里的信息什么时候会变

通知公告的时效性最强。发布日期和截止日期都是那一批的安排,过了时间就失效了。

政策法规修订后,旧文章里的解读可能不适用了。我们会发新文章覆盖,以最新一篇为准。

行业动态是背景参考,不是即时信息。今天看到的趋势,下个月可能就变了。

安全常识长期有效,但考试题库会更新。考前以最新辅导资料为准。

拿不准文章里的信息还能不能用的,直接打电话问。

每篇文章详情页右侧栏会自动推荐相关文章、最新文章和最近一周/一日/一月的热门文章。不用来回翻列表,顺着推荐往下看就行。

看完这篇文章,如果你还是拿不准自己能不能报名、该考哪个工种、费用多少——别对着文章猜,打 18236992212。

文章是通用情况,每个人的条件不一样。同样是电工证,有人能直报,有人要先补学历,有人要先体检。电话里说你的具体情况,我们给你算准。

这个页面上的所有链接,都是按根路径写的。你点哪个都能直接跳过去,不用怕 404。

看完这篇文章,建议你下一步

想报名的:先对照报考条件,再按流程准备材料。

想复审的:查复审流程,看你证什么时候到期。

选工种的:翻工种目录,看哪个适合你。

企业团报的:看团报说明,或者直接打电话谈方案。

还拿不准的:打 18236992212,把你的情况说清楚,我们告诉你下一步。

VERIFY

文章里的信息,什么时候该核实

信息类型有效期怎么核实
考试批次通知截止日期前有效打 18236992212 问最近还能不能报
政策法规条款修订前有效看最新一篇解读,或打电话问
费用标准长期参考报名时按当时报价为准
考试地点当批次有效约考后收到通知
报考条件政策修订时变打电话说你的情况判断
复审要求长期参考拿证时我们会记着日期
行业动态数据背景参考不作为即时决策依据
安全常识长期有效考前以最新辅导资料为准

看完这篇,建议再看看

如果这篇是通知:再看看其他批次通知,对比时间和工种。

如果这篇是政策:再看看其他法规解读,了解全貌。

如果这篇是行业动态:再看看其他行业新闻,了解趋势。

如果这篇是安全常识:再看看其他安全知识,为考试做准备。

右侧栏还会推荐相关文章、最新文章和最近一周热门文章,顺着看就行。

这个页面上的所有链接都是按根路径写的。你点哪个都能直接跳过去,不用怕 404。

看完这篇文章,如果觉得有用,转给需要的工友。他们也在愁考证的事。

文章是通用情况,每个人的条件不一样。同样是电工证,有人能直报,有人要先补学历,有人要先体检。电话里说你的具体情况,我们给你算准。

拿不准的,打 18236992212,别对着文章猜。

文章是通用情况,每个人的条件不一样。电话里说你的具体情况,我们给你算准。

这个页面上的所有链接都是按根路径写的。你点哪个都能直接跳过去,不用怕 404。

看完这篇文章,如果觉得有用,转给需要的工友。他们也在愁考证的事。

文章是通用情况,每个人的条件不一样。电话里说你的具体情况,我们给你算准。

这个页面上的所有链接都是按根路径写的。你点哪个都能直接跳过去,不用怕 404。

看完这篇文章,如果觉得有用,转给需要的工友。他们也在愁考证的事。

拿不准的,打 18236992212,别对着文章猜。

这个页面上的所有链接都是按根路径写的。你点哪个都能直接跳过去,不用怕 404。

看完这篇文章,如果觉得有用,转给需要的工友。他们也在愁考证的事。

文章是通用情况,每个人的条件不一样。电话里说你的具体情况,我们给你算准。

文章讲的是通用情况,你的情况要单独问

符不符合条件、最近一批还能不能报、费用怎么算,打 18236992212 一句就清楚。