Apache Beam Java Filter 变换实战:用 Filter.by 一行过滤 PCollection 中的奇数

📅 发布时间:2026/10/9 2:41:22
Apache Beam Java Filter 变换实战:用 Filter.by 一行过滤 PCollection 中的奇数
【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Apache Beam 的 Java SDK 提供了一组语言级的高层变换PTransform用来简化自定义DoFn的编写Filter就是其中最常用的一员。本文以learning/katas训练仓库中的「Filter Kata」为主线讲解如何用Filter.by(...)实现「过滤掉奇数、只保留偶数」的过滤逻辑并结合 SDK 源码剖析Filter的底层实现ParDoDoFn及其完整 API 家族lessThan、greaterThanEq、equal等。读完本文你将掌握 Beam Java 中最简洁、最规范的元素过滤写法并理解它与手写ParDo的等价关系与取舍。从 Kata 出发理解任务目标Kata编程训练是 Apache Beam 官方提供的渐进式学习材料位于仓库的 learning/katas/java 目录下。其中「Common Transforms / Filter / Filter」这一课的任务描述task.md原文如下The Beam SDKs provide language-specific ways to simplify how you provide your DoFn implementation.Kata:Implement a filter function that filters out the odd numbers by using Filter.这句话点明了本课的核心知识点Beam SDK 为各种语言提供了简化DoFn实现的高层方式而本课的 Kata 任务就是——使用Filter变换写一个过滤函数把输入集合中的奇数全部过滤掉。任务还给出了明确的实现提示UseFilter.by(...)也就是说标准解法不是手写一个完整的DoFn去逐元素判断而是直接调用Filter.by(predicate)传入一个谓词predicate让 SDK 帮你完成过滤。先看任务骨架你要补全的那一行本课的任务骨架位于 Filter/Task.java其中留有一个TODO()占位符对应配置见 task-info.yaml占位符位于applyTransform方法内。完整代码结构如下public class Task { public static void main(String[] args) { PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline Pipeline.create(options); PCollectionInteger numbers pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollectionInteger output applyTransform(numbers); output.apply(Log.ofElements()); pipeline.run(); } static PCollectionInteger applyTransform(PCollectionInteger input) { // TODO(): 在这里实现过滤逻辑 return input.apply(...); } }骨架中已经替你完成了三件事你需要关心的只有applyTransform构造管线PipelineOptionsFactory.fromArgs(args).create()从命令行参数创建PipelineOptions再交给Pipeline.create(options)创建Pipeline实例创建输入集合Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)生成一个包含 1~10 共十个整数的PCollectionInteger输出结果output.apply(Log.ofElements())把过滤后的每个元素打印到日志Log工具类位于 learning/katas/java/util/src/org/apache/beam/learning/katas/util/Log.java内部用ParDo SLF4JLogger实现逐元素输出。标准解法一行Filter.by过滤奇数按照任务提示在applyTransform中补全过滤逻辑只保留偶数、过滤掉奇数static PCollectionInteger applyTransform(PCollectionInteger input) { return input.apply(Filter.by(number - number % 2 0)); }Filter.by(...)接收一个谓词ProcessFunctionT, Boolean此处为 lambdanumber - number % 2 0对输入PCollection中的每个元素求值谓词返回true的元素保留偶数返回false的元素被丢弃奇数PCollectionInteger输入、输出类型一致FilterT本身是PTransformPCollectionT, PCollectionT。将骨架中的TODO()替换为上面这行后运行日志会依次输出2 4 6 8 10——1、3、5、7、9 这五个奇数全部被过滤掉了。用测试验证PAssert 断言结果Kata 的每个任务都配套了隐藏的单元测试visible: false本课的测试位于 Filter/TaskTest.java基于 Beam 官方测试组件TestPipelinePAssert编写public class TaskTest { Rule public final transient TestPipeline testPipeline TestPipeline.create(); Test public void filter() { Create.ValuesInteger values Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10); PCollectionInteger numbers testPipeline.apply(values); PCollectionInteger results Task.applyTransform(numbers); PAssert.that(results) .containsInAnyOrder(2, 4, 6, 8, 10); testPipeline.run().waitUntilFinish(); } }这里验证了三点测试直接调用你实现的Task.applyTransform(...)而不是复制一遍主方法中的逻辑因此被测对象是 Kata 真正要求你写的过滤函数PAssert.that(results).containsInAnyOrder(2, 4, 6, 8, 10)断言过滤结果精确地、以任意顺序包含2, 4, 6, 8, 10——既不能多如残留奇数也不能少如丢掉某个偶数testPipeline.run().waitUntilFinish()实际执行管线保证断言建立在真实运行结果之上而非静态检查。Filter API 家族不止by一种过滤方式Filter变换的完整实现位于 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Filter.java。除了 Kata 中用到的by(predicate)SDK 还提供了基于元素**自然顺序natural ordering**比较的一组快捷方法全部要求元素类型实现Comparable静态方法语义基于元素的自然顺序谓词等价形式by(predicate)保留满足任意谓词的元素predicate.apply(element)为truelessThan(v)保留小于v的元素x.compareTo(v) 0lessThanEq(v)保留小于等于v的元素x.compareTo(v) 0greaterThan(v)保留大于v的元素x.compareTo(v) 0greaterThanEq(v)保留大于等于v的元素x.compareTo(v) 0equal(v)保留等于v的元素x.compareTo(v) 0例如官方源码注释中给出的示例从PCollectionString中筛出长度大于 6 的单词、从数字集合中筛出小于 10 或大于 1000 的数PCollectionString longWords wordList.apply(Filter.by(new MatchIfWordLengthGT(6))); PCollectionInteger smallNumbers listOfNumbers.apply(Filter.lessThan(10)); PCollectionInteger largeNumbers listOfNumbers.apply(Filter.greaterThan(1000)); PCollectionInteger largeOrEqualNumbers listOfNumbers.apply(Filter.greaterThanEq(1000)); PCollectionInteger equalNumbers listOfNumbers.apply(Filter.equal(1000));从源码看这些比较类方法内部都是通过by(...) lambda 实现的例如lessThan的实现是public static T extends ComparableT FilterT lessThan(final T value) { return by((ProcessFunctionT, Boolean) input - input.compareTo(value) 0) .described(String.format(x %s, value)); }它们只是在谓词外面多包了一层.described(...)用一句人类可读的描述如x 10来填充DisplayData方便在监控和诊断界面中辨认这条过滤规则。源码透视Filter 的本质是 ParDo DoFn理解Filter底层原理的关键在于其expand方法Filter.javaOverride public PCollectionT expand(PCollectionT input) { return input .apply( ParDo.of( new DoFnT, T() { ProcessElement public void processElement(Element T element, OutputReceiverT r) throws Exception { if (predicate.apply(element)) { r.output(element); } } })) .setCoder(input.getCoder()); }这段实现揭示了两个关键事实Filter并非魔法而是一个内置的ParDo它对每个元素执行predicate.apply(element)只有结果为true时才调用r.output(element)向下游发射否则直接丢弃。这正是「过滤」在 Beam 中的标准语义——一个DoFn可以有零到多个输出不输出即过滤。输入输出的编码Coder保持一致过滤不改变元素类型因此expand末尾通过.setCoder(input.getCoder())复用输入PCollection的编码器避免重复推断也保证了流水线的高效。这也是为什么任务描述开篇会说「Beam SDKs provide language-specific ways to simplify how you provide your DoFn implementation」——Filter这类高层变换本质上是把你的意图编译成一个隐藏的DoFn让你免去手动声明DoFn类、ProcessElement注解、OutputReceiver等样板代码。对照实验用原生 ParDo 手写过滤同一个「过滤」需求你也可以完全手写ParDo。本课所在的「Common Transforms / Filter」课程还包含一节姊妹课 Filter using ParDo其任务恰好相反——过滤掉偶数、保留奇数且要求使用原生DoFn实现其参考实现位于 learning/katas/java/Common%20Transforms/Filter/ParDo/src/.../pardo/Task.javastatic PCollectionInteger applyTransform(PCollectionInteger input) { return input.apply(ParDo.of( new DoFnInteger, Integer() { ProcessElement public void processElement(Element Integer number, OutputReceiverInteger out) { if (number % 2 1) { out.output(number); } } })); }把两节课并排对比就能直观看到Filter带来的抽象收益语义层面等价两段代码都是「满足条件就output不满足就跳过」元素处理完全一一对应代码量悬殊手写版本需要完整的匿名DoFn类、ProcessElement注解、Element参数和OutputReceiver而Filter.by(number - number % 2 0)只有一行适用场景不同当过滤逻辑足够简单、可以表达为纯谓词时优先用Filter当过滤同时需要附带额外副作用如计数、窗口信息、状态处理时才需要退回到完整ParDo。如何运行与完成这个 KataKata 采用「先改代码、再跑测试、通过即完成」的训练闭环打开骨架编辑 learning/katas/java/Common%20Transforms/Filter/Filter/src/org/apache/beam/learning/katas/commontransforms/filter/filter/Task.java把applyTransform中的TODO()替换为return input.apply(Filter.by(number - number % 2 0));运行主方法执行Task.main(...)观察日志输出2, 4, 6, 8, 10确认奇数已被过滤默认运行在 Direct Runner 上即本地单机模式运行单元测试执行 TaskTest.java 中的filter()测试PAssert断言通过即表示解法正确进阶验证可以尝试把谓词改成number - number % 2 1保留奇数观察断言containsInAnyOrder(2, 4, 6, 8, 10)失败从而体会PAssert对输出集合的精确约束。整个 Java Kata 项目是独立的 Gradle 工程settings.gradle位于 learning/katas/java/settings.gradle可直接构建运行不依赖仓库其他模块。小结本文围绕「Filter Kata」讲解了 Beam Java SDK 中元素过滤的标准做法用一行Filter.by(predicate)替代手写DoFn通过PAssert精确验证输出并深入源码确认了Filter的底层就是「谓词驱动的ParDoDoFn」元素不满足条件就不发射。同时对比了同课程的 ParDo 手写版本明确了两者的等价性与适用边界。掌握Filter.by以及lessThan、greaterThanEq、equal这一组基于自然顺序的快捷过滤方法你就掌握了 Beam 数据清洗与预处理中最常用的一把利器。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战用 Filter 变换过滤奇数Filter.by 谓词编程详解Apache Beam Java Kata 实战用 Filter 变换过滤奇数Filter.by 谓词编程详解 本指南围绕 Apache Beam 官方大数据批处理流处理数据工程Apache Beam 的 Filter 变换多语言 PCollection 数据过滤完整指南Apache Beam 的 Filter 变换多语言 PCollection 数据过滤完整指南 导读 Filter 是 Apache Beam 中最常用的基础大数据批处理流处理数据工程Apache Beam Go SDK 实战使用 filter 包Include/Exclude过滤 PCollection 元素Apache Beam Go SDK 实战使用 filter 包Include/Exclude过滤 PCollection 元素 导读 本文聚焦 Apac上一篇learn-harness-engineering Project 06完全なエージェント harness を構築する — ランタイム観測性・クリーン状態・ablation 検証の総合実践ガイド下一篇opencodex 桌面端上下文窗口与推理强度控制Cycle 2 A-门审计合成实战解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考