☰
Apache Beam Kotlin Kata 实战:用 Filter.by 实现数据过滤(Common Transforms)
2026/9/28 21:09:12 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

导读

本文围绕 Apache Beam 官方 Kotlin Kata 练习集中的 Filter 一课,讲解如何利用 Beam SDK 内置的Filter变换,以声明式的方式从PCollection中过滤满足条件的元素。文章将以“过滤出偶数”这一经典任务为主线,完整给出练习骨架、解题实现、单元测试与运行方式,并深入Filter变换的底层源码,说明它与手写ParDo + DoFn的等价关系以及Filter.by / lessThan / greaterThan / equal等常用 API 的适用场景。读完本文,你将掌握在 Kotlin 中编写 Beam 过滤逻辑的标准范式,并能够对照测试与源码验证自己的实现。

练习背景:Beam Katas 中的 Filter 一课

本练习位于仓库 learning/katas/kotlin/Common Transforms/Filter 目录下,属于 Beam Katas 交互式练习课程的 "Common Transforms" 单元。该单元包含两课:

  • ParDo:使用ParDo+DoFn手动实现过滤;
  • Filter:使用 Beam SDK 提供的Filter变换完成同样的任务。

课程内容定义在 lesson-info.yaml 中,两个任务依次为ParDo与Filter。这种"先手写 DoFn,再使用内置变换"的安排,是为了让学习者体会 Beam SDK 提供的语言级简化:当你的过滤逻辑足够简单时,不必手写完整的DoFn,直接使用Filter.by(...)即可,这正是练习文档 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 usingFilter.

即:实现一个过滤函数,使用Filter变换从输入数字中过滤掉奇数,只保留偶数。练习给出的提示是使用Filter.by(...)。

练习骨架解读

任务目录 learning/katas/kotlin/Common Transforms/Filter/Filter 下包含以下关键文件:

文件作用
Task.kt练习主文件,包含待补全的applyTransform函数
TaskTest.kt隐藏的单元测试,用于验证你的实现
task-info.yaml练习配置,其中placeholders指定了待补全的代码位置
task.md任务描述文档

从 task-info.yaml 可以看出,练习在Task.kt中预留了一个TODO()占位符(位于源码 offset 1654、长度 79 字节处),测试文件对学习者隐藏(visible: false),你需要自己实现过滤逻辑并通过隐藏测试。

练习主文件的骨架如下:

object Task { @JvmStatic fun main(args: Array<String>) { val options = PipelineOptionsFactory.fromArgs(*args).create() val pipeline = Pipeline.create(options) val numbers = pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)) val output = applyTransform(numbers) output.apply(Log.ofElements()) pipeline.run() } @JvmStatic fun applyTransform(input: PCollection<Int>): PCollection<Int> { // TODO(): 在此实现过滤逻辑 } }

骨架的流程非常清晰,构成了一个最小可运行的 Beam 管道:

  1. PipelineOptionsFactory.fromArgs(*args).create()从命令行参数解析管道配置;
  2. Pipeline.create(options)创建管道实例;
  3. pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10))用Create变换将内存中的 1~10 构造成输入PCollection<Int>;
  4. applyTransform(numbers)是你要补全的过滤步骤;
  5. output.apply(Log.ofElements())将结果逐条打印到日志。Log是 Katas 项目提供的调试辅助工具(见 Log.kt),其内部通过ParDo+DoFn将每个元素输出为日志行,并原样传递(out.output(element)),不会改变数据内容;
  6. pipeline.run()触发管道执行。

你只需要关注applyTransform这一个函数——输入是PCollection<Int>,输出也必须是PCollection<Int>,元素类型保持不变。

标准解法:Filter.by 一行完成过滤

根据 task.md 的提示,使用Filter.by(...)即可。完整的实现如下(即仓库中 Task.kt 给出的标准答案):

import org.apache.beam.sdk.transforms.Filter import org.apache.beam.sdk.transforms.SerializableFunction object Task { // ... main 函数同骨架 ... @JvmStatic fun applyTransform(input: PCollection<Int>): PCollection<Int> { return input.apply(Filter.by(SerializableFunction { number: Int -> number % 2 == 0 })) } }

要点说明:

  • Filter.by(predicate)接收一个"谓词"(predicate):对每个元素返回true则保留,返回false则丢弃。这里number % 2 == 0对偶数返回true,因此奇数被过滤掉、偶数被保留,恰好满足任务"过滤掉奇数"的要求;
  • 在 Kotlin 中,SerializableFunction { number: Int -> number % 2 == 0 }使用 SAM 转换直接以 lambda 形式提供谓词,并保证其可序列化——这是 Beam 将处理逻辑分发到分布式执行环境(如 Dataflow、Flink、Spark)时对函数对象的基本要求;
  • 谓词逻辑放在applyTransform中,管道骨架与业务逻辑分离,便于单元测试直接针对applyTransform做断言。

运行该管道时,Log.ofElements()会依次打印出2、4、6、8、10五个偶数。

单元测试:如何验证过滤正确性

练习的隐藏测试 TaskTest.kt 是验证实现正确性的权威依据:

class TaskTest { @get:Rule @Transient val testPipeline: TestPipeline = TestPipeline.create() @Test fun common_transforms_filter_filter() { val values = Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10) val numbers = testPipeline.apply(values) val results = applyTransform(numbers) PAssert.that(results).containsInAnyOrder(2, 4, 6, 8, 10) testPipeline.run().waitUntilFinish() } }

测试使用了 Beam 官方测试框架的核心组件,这也是在真实 Beam 项目中验证管道逻辑的标准写法:

  • TestPipeline(来自org.apache.beam.sdk.testing.TestPipeline):作为 JUnit 的@Rule注入,负责自动管理管道的创建、执行与资源释放;
  • PAssert.that(results).containsInAnyOrder(2, 4, 6, 8, 10):对输出PCollection做断言。containsInAnyOrder校验结果集合与期望集合内容一致、与元素顺序无关——这在分布式并行处理中至关重要,因为元素在并行执行时的到达顺序无法保证;
  • testPipeline.run().waitUntilFinish():运行管道并阻塞等待执行完成,确保断言真正被执行。

如果你把谓词写反(例如number % 2 == 1),测试会立即失败,因为输出变成了1, 3, 5, 7, 9,与期望的2, 4, 6, 8, 10不符。所以测试不仅指导你"怎么写",也精确约束了"过滤方向"。

底层原理:Filter 就是内置的 ParDo

了解了如何使用之后,我们再深入一层:Filter变换在源码层面是如何实现的?答案在 sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Filter.java 中非常直观。

Filter<T>本身是一个PTransform<PCollection<T>, PCollection<T>>(第 30 行),它的核心expand方法(第 201~216 行)实现如下:

@Override public PCollection<T> expand(PCollection<T> input) { return input .apply( ParDo.of( new DoFn<T, T>() { @ProcessElement public void processElement(@Element T element, OutputReceiver<T> r) throws Exception { if (predicate.apply(element)) { r.output(element); } } })) .setCoder(input.getCoder()); }

这段源码揭示了一个重要事实:Filter内部就是包裹了一个ParDo+DoFn——DoFn对每个元素调用你提供的predicate,返回true才调用r.output(element)放行,否则丢弃;最后通过.setCoder(input.getCoder())显式沿用输入PCollection的编码器(coder),保证输出元素类型与输入一致。

这与你在"ParDo"一课中手写的过滤逻辑(见 ParDo 任务的 Task.kt)在语义上是完全等价的——ParDo 版本在DoFn.processElement中手动判断number % 2 == 1后仅输出符合条件的元素,而Filter只是把这个"判断后选择性输出"的模式封装成了可复用的谓词形式。理解了这一点,你就明白了练习安排两课的意义:先用 ParDo 理解执行模型,再用 Filter 体验高层抽象带来的简洁性。

另外值得注意by方法的两个重载(第 49~58 行):

  • by(ProcessFunction<T, Boolean> predicate):主入口,接收ProcessFunction;
  • by(SerializableFunction<T, Boolean> predicate):为兼容既有代码而保留的二进制兼容适配器,内部把SerializableFunction转成ProcessFunction再调用主版本。

Kotlin 练习中使用的SerializableFunction { ... }走的就是第二个重载,在 Java/Kotlin 互操作层面依然顺畅可用。

更多过滤方式:比较型 API 与命名谓词

Filter不止提供by一种用法。从 Filter.java 的源码可以看到,它还基于元素的**自然排序(natural ordering)**提供了一组比较型工厂方法,适用于元素类型实现了Comparable的场景:

静态方法保留条件(源码实现)典型示例
Filter.by(predicate)predicate.apply(element) == truenumber % 2 == 0
Filter.lessThan(value)element < value(compareTo < 0,第 79~82 行)过滤出小于 10 的数字
Filter.lessThanEq(value)element <= value(compareTo <= 0,第 127~130 行)过滤出不超过 10 的数字
Filter.greaterThan(value)element > value(compareTo > 0,第 103~106 行)过滤出大于 1000 的数字
Filter.greaterThanEq(value)element >= value(compareTo >= 0,第 151~154 行)过滤出大于等于 0 的数字
Filter.equal(value)element == value(compareTo == 0,第 174~177 行)只保留等于 1000 的数字

这些比较型方法在源码内部同样通过by(...)构造谓词(例如lessThan等价于by(input -> input.compareTo(value) < 0)),并调用described(...)为变换附加了人类可读的描述(如"x < 10"),这些描述会通过DisplayData暴露到监控界面中,方便调试与运维查看管道结构。

从源码结构看,无论使用哪种工厂方法,最终构造出的都是同一个内部状态——predicate(谓词)与predicateDescription(描述),并在expand中统一按"谓词为真即输出"的方式执行(第 181~191 行、第 201~216 行)。

在实际业务中,过滤条件往往不止是"偶数"这样简单。两个实用建议:

  1. 命名谓词可读性更强:当条件复杂时,可以像源码 javadoc 中的示例那样,把谓词封装成有名字的类或函数,例如定义一个MatchIfWordLengthGT(6)风格的对象,让管道意图一目了然;
  2. 复杂逻辑仍然可以回退到 DoFn:如果过滤逻辑涉及有状态计算、定时器或需要访问窗口信息,Filter.by的谓词形式就无法覆盖了,此时应像 ParDo 任务的 Task.kt 那样手写DoFn。Filter是"简单即用、复杂可退"的最佳实践示例。

运行与练习方式

Kotlin Katas 练习位于 learning/katas/kotlin 目录,它是一套独立可运行的 Gradle 项目(含gradlew、settings.gradle与course-info.yaml)。练习的基本运行方式如下:

# 在 learning/katas/kotlin 目录下运行某个任务的 main 函数 ./gradlew run

实际使用时,通常是在支持 Edu 插件的 IDE(如 IntelliJ IDEA 的 JetBrains Academy 插件)中打开该目录,按课程顺序逐题完成:先读 task.md 理解任务,在 Task.kt 中替换TODO()占位符,然后运行隐藏的 TaskTest.kt 验证是否通过。

需要说明的前提条件:运行这些练习需要 JDK 与 Gradle 环境就绪,main函数默认使用本地执行(管道通过PipelineOptionsFactory创建,未指定 Runner 时由 Beam 默认机制选择),因此无需任何外部集群即可完成本课的验证。

小结

通过本课你可以掌握一个核心结论:在 Apache Beam 中,简单的按条件筛选应优先使用Filter变换而非手写DoFn。具体来说:

  • Filter.by(predicate)接收"谓词",true保留、false丢弃,Kotlin 下用SerializableFunction { ... }一行即可完成;
  • Filter的底层实现就是ParDo+DoFn(见 Filter.java 第 201~216 行),与手写 DoFn 语义等价,但代码更简洁、意图更清晰;
  • 除by外,lessThan、lessThanEq、greaterThan、greaterThanEq、equal提供了基于自然排序的常见过滤捷径;
  • 用PAssert.that(...).containsInAnyOrder(...)配合TestPipeline验证过滤结果,是 Beam 社区的标准测试范式(参见 TaskTest.kt)。

将这一课延伸到真实项目:无论是清洗日志、剔除无效记录还是按阈值筛选指标,Filter都是 Beam 管道中最常用也最高频的基础变换之一。

  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询