【问题标题】:Spark does not push filter (PushedFilters array is empty)Spark 不推送过滤器(PushedFilters 数组为空)
【发布时间】:2023-03-14 03:12:02
【问题描述】:

简介

我注意到我们项目中的所有推送过滤器都不起作用。它解释了为什么执行时间会受到影响,因为它读取了数百万次读取,而应该将其减少到数千次。为了调试问题,我编写了一个读取 CSV 文件、过滤内容(PushDown Filter)并返回结果的小测试。

它不适用于 CSV,因此我尝试读取 parquet 文件。它们都不起作用。

数据

people.csv 文件具有以下结构:

first_name,last_name,city  // header
FirstName1,LastName1,Bern // 1st row
FirstName2,LastName2,Sion // 2nd row
FirstName3,LastName3,Bulle // 3rd row

注意:parquet 文件具有相同的结构

读取 CSV 文件

为了重现这个问题,我编写了一个最小的代码来读取一个 csv 文件并且应该只返回过滤后的数据。

读取csv文件并打印实物计划:

Dataset<Row> ds = sparkSession.read().option("header", "true").csv(BASE_PATH+"people.csv");
ds.where(col("city").equalTo("Bern")).show();
ds.explain(true);

实物图:

+----------+---------+----+
|名|姓|城市|
+---------+---------+----+
|名1|姓1|伯尔尼|
+---------+---------+----+

== 解析逻辑计划 == 关系[first_name#10,last_name#11,city#12] csv

== 分析的逻辑计划 == first_name: string, last_name: string, city: string Relation[first_name#10,last_name#11,city#12] csv

== 优化逻辑计划 == 关系[first_name#10,last_name#11,city#12] csv

== 物理计划 == *(1) FileScan csv [first_name#10,last_name#11,city#12] 批处理:false,格式:CSV,位置: InMemoryFileIndex[文件:people.csv], PartitionFilters:[],PushedFilters:[],ReadSchema: 结构体

我已经用 parquet 文件进行了测试,不幸的是结果是一样的。

我们可以注意到的是:

  • PushedFilters 为空,我希望过滤器包含谓词。
  • 返回的结果还是正确的。

我的问题是:为什么这个 PushedFilters 是空的?

注:

  • Spark 版本:2.4.3
  • 文件系统:ext4(和集群上的HDFS,都没有工作)

【问题讨论】:

  • 您确定它不适用于镶木地板文件吗?你用的是什么版本的火花?文件存储在哪里? (hdfs?s3? ...)
  • 嗨奥利!我都尝试过(CSV 和镶木地板文件),但都没有工作。对于火花版本,我在问题末尾添加了

标签: java csv apache-spark parquet


【解决方案1】:

您正在对第一个数据集调用解释,即只有读取的数据集。试试类似的东西(对不起,我只有 Scala 环境可用):

val ds: DataFrame = spark.read.option("header", "true").csv("input.csv")
val f = ds.filter(col("city").equalTo("Bern"))

f.explain(true)

f.show()

此外,由于this,在使用类型化数据集 API 时要小心。不过不应该是你的情况。

【讨论】:

    【解决方案2】:

    只是为了记录,这里是解决方案(感谢 LizardKing):

    结果

    • 之前:PushedFilters: []
    • 之后:PushedFilters: [IsNotNull(city), EqualTo(city,Bern)]

    代码

    Dataset<Row> ds = sparkSession.read().option("header", "true").csv(BASE_PATH+"people.csv");
    Dataset<Row> dsFiltered = ds.where(col("city").equalTo("Bern"));
    dsFiltered.explain(true);
    

    物理计划

    物理计划看起来好多了:

    == Parsed Logical Plan ==
    'Filter ('city = Bern)
    +- Relation[first_name#10,last_name#11,city#12] csv
    
    == Analyzed Logical Plan ==
    first_name: string, last_name: string, city: string
    Filter (city#12 = Bern)
    +- Relation[first_name#10,last_name#11,city#12] csv
    
    == Optimized Logical Plan ==
    Filter (isnotnull(city#12) && (city#12 = Bern))
    +- Relation[first_name#10,last_name#11,city#12] csv
    
    == Physical Plan ==
    *(1) Project [first_name#10, last_name#11, city#12]
    +- *(1) Filter (isnotnull(city#12) && (city#12 = Bern))
       +- *(1) FileScan csv [first_name#10,last_name#11,city#12] Batched: false, Format: CSV, Location: InMemoryFileIndex[file:./people.csv], PartitionFilters: [], PushedFilters: [IsNotNull(city), EqualTo(city,Bern)], ReadSchema: struct<first_name:string,last_name:string,city:string>
    

    【讨论】:

      猜你喜欢
      • 2021-07-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-09-12
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-06-13
      相关资源
      最近更新 更多