【问题标题】:Spark: Filters on Functions are not pushed downSpark:函数上的过滤器没有被下推
【发布时间】:2020-10-17 19:25:31
【问题描述】:

我正在运行 spark 作业 (scala) 将列转换为上列并过滤其值

val udfExampleDF = spark.read.parquet("a.parquet")
udfExampleDF.filter(upper(col("device_type"))==="PHONE").select(col("device_type")).show()

我看到这个过滤器没有被按下。这是物理计划的样子

== Physical Plan ==
CollectLimit 21
+- *(1) Filter (upper(device_type#138) = PHONE)
   +- *(1) FileScan parquet [device_type#138] Batched: true, Format: Parquet, Location: InMemoryFileIndex[a.parquet, PartitionFilters: [], PushedFilters: [], ReadSchema: struct<orig_device_type:string>

但是,如果我只过滤没有上层,过滤器会被推下。

知道为什么会这样。我假设这些过滤器在所有情况下都会被下推。感谢您的帮助。

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    由于谓词下推试图移除逻辑运算符并将它们推送到数据源,在我们的例子中从FilterFileScan parquet,它必须使用原始列值。

    理论上,如果其他文件格式支持使用更改的列值进行过滤,它可能会起作用(不确定PushDownPredicate 是否支持这种情况,即使存在这样的文件格式)。

    要解决您的问题,请通过在等式的另一侧设置动态值来解决此问题,尽管您有几个条件它会快得多(尝试将更频繁的值放在第一个条件中):

    udfExampleDF.filter(col("device_type")==="Phone" or col("device_type")==="PHONE" or col("device_type")==="phone").select(col("device_type")).show()
    

    您还可以在数据摄取中将列值统一为始终为 PHONE,即在编写 parquet 文件之前应用 upper 函数,然后像这样过滤:

    udfExampleDF.filter(col("device_type")==="PHONE").select(col("device_type")).show()
    

    【讨论】:

    • 谢谢。似乎没有办法直接使用函数并具有过滤器下推。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-04-06
    • 1970-01-01
    • 2023-03-13
    • 2022-01-07
    • 1970-01-01
    • 1970-01-01
    • 2016-01-25
    相关资源
    最近更新 更多