Apache Beam Java SDK 的 Filter 变换:按谓词与自然顺序过滤 PCollection 的完整指南
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载导读Filter是 Apache Beam Java SDK 中最常用的元素级Elementwise变换之一它根据给定的谓词Predicate从PCollection中剔除不满足条件的元素也可基于元素类型的自然排序natural ordering直接按数值不等式进行过滤。本文以 Apache Beam 官方文档 filter.md 为核心结合 SDK 源码与测试系统讲解Filter.by、Filter.greaterThan、Filter.lessThanEq等全部内置过滤器以及其底层通过ParDo实现的执行原理帮助你写出可运行、可验证的过滤逻辑。Filter 变换能做什么Filter是一个PTransformPCollectionT, PCollectionT其输入与输出的元素类型完全一致功能是筛选而非转换Given a predicate, filter out all elements that dont satisfy that predicate. May also be used to filter based on an inequality with a given value based on the natural ordering of the element.即在给定一个谓词的情况下过滤掉所有不满足该谓词的元素也可以基于元素类型的自然排序按与某个给定值的不等式关系进行过滤。在数据管道中Filter的典型应用场景包括清洗数据剔除空值、非法格式或长度异常的记录分流业务只保留满足特定业务规则的订单、日志或事件数值筛选按阈值大于、小于、等于圈定数据范围。由于Filter不改变元素类型、只改变集合的成员它非常适合与其他元素级变换如 MapElements、FlatMapElements、ParDo串联成完整的处理链路。方式一用Filter.by传入自定义谓词Filter.by(predicate)接受一个谓词函数返回新的PTransform只保留谓词返回true的元素。谓词类型在 Filter.java 中有两个重载public static T, PredicateT extends ProcessFunctionT, Boolean FilterT by( PredicateT predicate) // Binary compatibility adapter public static T, PredicateT extends SerializableFunctionT, Boolean FilterT by( PredicateT predicate)ProcessFunctionT, Boolean允许apply方法抛出受检异常SerializableFunctionT, Boolean二进制兼容适配层内部会转成ProcessFunction调用。使用匿名类实现谓词官方文档给出的第一个示例只保留长度大于 3 的字符串。PCollectionString allStrings Create.of(Hello, world, hi); PCollectionString longStrings allStrings .apply(Filter.by(new SerializableFunctionString, Boolean() { Override public Boolean apply(String input) { return input.length() 3; } }));结果是包含Hello和world的PCollectionhi长度 2被过滤掉。使用 Lambda 表达式简化谓词在 Java 8 中可以直接用 Lambda 替代匿名类代码更紧凑PCollectionInteger numbers pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollectionInteger evenNumbers numbers.apply(Filter.by(n - n % 2 0));这也是仓库中练习任务的标准写法。例如 学习路径 Filter 练习 中就是用Filter.by(number - number % 2 0)过滤出偶数Tour of Beam 的 Filter 示例 同样使用Filter.by(number - number % 2 0)保留 2、4、6、8、10。使用方法引用Method ReferenceFilter.by也支持方法引用形式。在 FilterTest.java 的testFilterByMethodReferenceWithLambda中可以看到PCollectionInteger output p.apply(Create.of(1, 2, 3, 4, 5, 6, 7)).apply(Filter.by(new EvenFilter()::isEven));方式二基于自然排序的内置不等式过滤器如果元素类型实现了Comparable可以直接使用Filter提供的五个内置不等式过滤器无需编写任何谓词方法保留条件内部实现compareToDisplayData 描述Filter.greaterThan(v)x vinput.compareTo(value) 0x vFilter.greaterThanEq(v)x vinput.compareTo(value) 0x ≥ vFilter.lessThan(v)x vinput.compareTo(value) 0x vFilter.lessThanEq(v)x vinput.compareTo(value) 0x ≤ vFilter.equal(v)x vinput.compareTo(value) 0x v这些方法的方法签名均为public static T extends ComparableT FilterT greaterThan(final T value)即要求T必须实现ComparableT比较基于元素的自然排序natural ordering因此适用于数值、字符串、日期等可比较类型。官方文档的第二个示例PCollectionLong numbers Create.of(1L, 2L, 3L, 4L, 5L); PCollectionLong bigNumbers numbers.apply(Filter.greaterThan(3)); PCollectionLong smallNumbers numbers.apply(Filter.lessThanEq(3));bigNumbers保留4L、5LsmallNumbers保留1L、2L、3L。其余变体还包括Filter.greaterThanEq、Filter.lessThan和Filter.equal。下面给出覆盖全部五个内置过滤器的完整可运行示例// 输入[1, 2, 3, 4, 5, 6, 7, 8, 9, 10] PCollectionInteger input pipeline.apply(Create.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)); PCollectionInteger greaterThanEqNumbers input.apply(Filter.greaterThanEq(3)); // 结果[3, 4, 5, 6, 7, 8, 9, 10] PCollectionInteger greaterThanNumbers input.apply(Filter.greaterThan(4)); // 结果[5, 6, 7, 8, 9, 10] PCollectionInteger lessThanNumbers input.apply(Filter.lessThan(10)); // 结果[1, 2, 3, 4, 5, 6, 7, 8, 9] PCollectionInteger lessThanEqNumbers input.apply(Filter.lessThanEq(7)); // 结果[1, 2, 3, 4, 5, 6, 7] PCollectionInteger equalNumbers input.apply(Filter.equal(9)); // 结果[9]该示例与 Tour of Beam Filter 单元描述 中的Example 2: Filtering with a built-in methods完全一致可直接照搬验证。源码原理Filter 底层就是一个 ParDo从源码结构看Filter并非独立实现元素级处理而是直接复用了ParDo的执行机制。Filter.java 的expand方法是其核心Override 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将谓词包装进一个DoFn通过ParDo.of(...)应用每个输入元素进入processElement后只有predicate.apply(element)返回true才调用r.output(element)输出从而实现过滤因为输出元素与输入元素类型完全相同expand末尾通过.setCoder(input.getCoder())显式复用输入PCollection的 Coder避免重新推断编码器所有内置不等式过滤器greaterThan、lessThanEq等本质上是Filter.by的语法糖——它们内部把compareTo比较表达式封装成一个ProcessFunction谓词例如public static T extends ComparableT FilterT greaterThan(final T value) { return by((ProcessFunctionT, Boolean) input - input.compareTo(value) 0) .described(String.format(x %s, value)); }.described(...)会生成一个新的Filter实例不修改原变换并把描述信息写入DisplayData便于在监控界面中查看当前过滤器语义。populateDisplayData中builder.add(DisplayData.item(predicate, predicateDescription).withLabel(Filter Predicate));相应的FilterTest.java 的testDisplayData验证了每种内置过滤器的显示描述assertThat(DisplayData.from(Filter.lessThan(123)), hasDisplayItem(predicate, x 123)); assertThat(DisplayData.from(Filter.lessThanEq(234)), hasDisplayItem(predicate, x ≤ 234)); assertThat(DisplayData.from(Filter.greaterThan(345)), hasDisplayItem(predicate, x 345)); assertThat(DisplayData.from(Filter.greaterThanEq(456)), hasDisplayItem(predicate, x ≥ 456)); assertThat(DisplayData.from(Filter.equal(567)), hasDisplayItem(predicate, x 567));用测试验证 Filter 的边界行为FilterTest.java 覆盖了Filter的完整行为边界是理解语义的最佳参考恒真谓词Filter.by(new TrivialFn(true))保留全部元素identity 语义恒假谓词Filter.by(new TrivialFn(false))输出空PCollectionPAssert.that(output).empty()普通谓词Filter.by(new EvenFn())在[1..7]上只保留[2, 4, 6]ProcessFunction 谓词Filter.by(new EvenProcessFn())验证可抛异常的ProcessFunction路径内置不等式lessThan(4)→[1,2,3]greaterThan(4)→[5,6,7]lessThanEq(4)→[1,2,3,4]greaterThanEq(4)→[4,5,6,7]equal(4)→[4]Lambda 谓词Filter.by(i - i % 2 0)与匿名类结果一致。其中testFilterParDoOutputTypeDescriptorRawWithLambda还记录了一个值得注意的细节当 Lambda 导致 raw type 时输出类型的 TypeDescriptor 无法用于 Coder 推断会抛出CannotProvideCoderException。因此在实际生产中建议为Filter.by的 Lambda 提供显式类型信息或依赖expand中.setCoder(input.getCoder())复用输入 Coder 来规避该问题。实战练习与进阶用法练习一Katas 中的 Filter 任务在仓库的 Katas 学习路径 中练习要求从 110 中过滤出偶数标准答案就是static PCollectionInteger applyTransform(PCollectionInteger input) { return input.apply(Filter.by(number - number % 2 0)); }同一目录下的 Filter/ParDo 任务 则展示了一个对照实验手工编写DoFn保留奇数number % 2 1帮助理解Filter与原生ParDo的等价关系。练习二组合多个 Filter 实现复杂条件Tour of Beam 的 Filter 说明 给出了两个进阶方向链式组合多个简单 Filter例如过滤出以字母 a 开头不分大小写且长度大于 3的单词可以串联两次Filter.by单个 Filter 内实现复杂逻辑用一个谓词同时判断多个条件减少算子数量。注意Filter 与 ParDo 的取舍官方文档在 Related transforms 部分明确指出FlatMapElements行为与Map相同但每个输入可能产生零个或多个输出ParDo最通用的元素级映射操作还支持多输出集合与侧输入side-inputs。因此当过滤逻辑需要依赖侧输入、需要多路输出或需在过滤的同时改写元素时应优先考虑ParDo当只是保留满足条件的一类元素时Filter是语义最清晰、代码最简洁的选择。小结Filter是PCollectionT→PCollectionT的筛选变换不改变元素类型Filter.by支持SerializableFunction、ProcessFunction、Lambda 与方法引用四种谓词写法内置greaterThan/greaterThanEq/lessThan/lessThanEq/equal五个基于Comparable自然排序的不等式过滤器并有对应的 DisplayData 描述便于监控源码层面Filter底层复用ParDo.of(DoFn)通过if (predicate.apply(element)) r.output(element)实现逐元素过滤并复用输入 Coder行为边界恒真、恒假、边界值包含与否可参考 FilterTest.java 中的完整测试用例。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Java SDK Filter 转换实战用 Filter.by 谓词过滤 PCollection 元素Apache Beam Java SDK Filter 转换实战用 Filter.by 谓词过滤 PCollection 元素 Apache Beam 的 J批处理流处理大数据Apache Beam Java Filter 转换详解基于谓词与自然序不等式的元素过滤Apache Beam Java Filter 转换详解基于谓词与自然序不等式的元素过滤 Apache Beam 的 Filter 是 PCollection大数据批处理流处理数据工程Apache Beam Filter 变换实战三语言按条件过滤 PCollection 的完整指南Apache Beam Filter 变换实战三语言按条件过滤 PCollection 的完整指南 Filter 是 Apache Beam 中最常用的基础变批处理流处理大数据上一篇gh_mirrors/caf/caffe2超参数优化高效调参方法与工具下一篇Easings.net社区生态第三方插件与扩展资源创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考