【问题标题】:Spark Parquet Statistics(min/max) integrationSpark Parquet 统计(最小/最大)集成
【发布时间】:2017-06-01 16:46:26
【问题描述】:

我一直在研究 Spark 如何在 Parquet 中存储统计信息(最小值/最大值),以及它如何使用这些信息进行查询优化。 我有几个问题。 第一个设置:Spark 2.1.0,下面设置一个1000行的Dataframe,一个long类型,一个string类型列。 不过,它们按不同的列排序。

scala> spark.sql("select id, cast(id as string) text from range(1000)").sort("id").write.parquet("/secret/spark21-sortById")
scala> spark.sql("select id, cast(id as string) text from range(1000)").sort("Text").write.parquet("/secret/spark21-sortByText")

我在 parquet-tools 中添加了一些代码来打印统计数据并检查生成的 parquet 文件:

hadoop jar parquet-tools-1.9.1-SNAPSHOT.jar meta /secret/spark21-sortById/part-00000-39f7ac12-6038-46ee-b5c3-d7a5a06e4425.snappy.parquet 
file:        file:/secret/spark21-sortById/part-00000-39f7ac12-6038-46ee-b5c3-d7a5a06e4425.snappy.parquet 
creator:     parquet-mr version 1.8.1 (build 4aba4dae7bb0d4edbcf7923ae1339f28fd3f7fcf) 
extra:       org.apache.spark.sql.parquet.row.metadata = {"type":"struct","fields":[{"name":"id","type":"long","nullable":false,"metadata":{}},{"name":"text","type":"string","nullable":false,"metadata":{}}]} 

file schema: spark_schema 
--------------------------------------------------------------------------------
id:          REQUIRED INT64 R:0 D:0
text:        REQUIRED BINARY O:UTF8 R:0 D:0

row group 1: RC:5 TS:133 OFFSET:4 
--------------------------------------------------------------------------------
id:           INT64 SNAPPY DO:0 FPO:4 SZ:71/81/1.14 VC:5 ENC:PLAIN,BIT_PACKED STA:[min: 0, max: 4, num_nulls: 0]
text:         BINARY SNAPPY DO:0 FPO:75 SZ:53/52/0.98 VC:5 ENC:PLAIN,BIT_PACKED

hadoop jar parquet-tools-1.9.1-SNAPSHOT.jar meta /secret/spark21-sortByText/part-00000-3d7eac74-5ca0-44a0-b8a6-d67cc38a2bde.snappy.parquet 
file:        file:/secret/spark21-sortByText/part-00000-3d7eac74-5ca0-44a0-b8a6-d67cc38a2bde.snappy.parquet 
creator:     parquet-mr version 1.8.1 (build 4aba4dae7bb0d4edbcf7923ae1339f28fd3f7fcf) 
extra:       org.apache.spark.sql.parquet.row.metadata = {"type":"struct","fields":[{"name":"id","type":"long","nullable":false,"metadata":{}},{"name":"text","type":"string","nullable":false,"metadata":{}}]} 

file schema: spark_schema 
--------------------------------------------------------------------------------
id:          REQUIRED INT64 R:0 D:0
text:        REQUIRED BINARY O:UTF8 R:0 D:0

row group 1: RC:5 TS:140 OFFSET:4 
--------------------------------------------------------------------------------
id:           INT64 SNAPPY DO:0 FPO:4 SZ:71/81/1.14 VC:5 ENC:PLAIN,BIT_PACKED STA:[min: 0, max: 101, num_nulls: 0]
text:         BINARY SNAPPY DO:0 FPO:75 SZ:60/59/0.98 VC:5 ENC:PLAIN,BIT_PACKED

所以问题是为什么 Spark,特别是 2.1.0,只为数字列生成最小值/最大值,而不是字符串(BINARY)字段,即使字符串字段包含在排序中?也许我错过了配置?

第二个问题,我如何确认 Spark 正在使用 min/max?

scala> sc.setLogLevel("INFO")
scala> spark.sql("select * from parquet.`/secret/spark21-sortById` where id=4").show

我有很多这样的行:

17/01/17 09:23:35 INFO FilterCompat: Filtering using predicate: and(noteq(id, null), eq(id, 4))
17/01/17 09:23:35 INFO FileScanRDD: Reading File path: file:///secret/spark21-sortById/part-00000-39f7ac12-6038-46ee-b5c3-d7a5a06e4425.snappy.parquet, range: 0-558, partition values: [empty row]
...
17/01/17 09:23:35 INFO FilterCompat: Filtering using predicate: and(noteq(id, null), eq(id, 4))
17/01/17 09:23:35 INFO FileScanRDD: Reading File path: file:///secret/spark21-sortById/part-00193-39f7ac12-6038-46ee-b5c3-d7a5a06e4425.snappy.parquet, range: 0-574, partition values: [empty row]
...

问题是 Spark 似乎正在扫描每个文件,即使从 min/max 来看,Spark 也应该能够确定只有 part-00000 具有相关数据。或者也许我读错了,Spark 正在跳过文件?也许Spark只能使用分区值来跳过数据?

【问题讨论】:

    标签: apache-spark parquet


    【解决方案1】:

    PARQUET-686 进行了更改,以便在似乎合适时有意忽略二进制字段的统计信息。您可以通过将parquet.strings.signed-min-max.enabled 设置为true 来覆盖此行为。

    设置该配置后,您可以使用 parquet-tools 读取二进制字段中的最小值/最大值。

    更多详情my another stackoverflow question

    【讨论】:

    • 您好,我为我的 spark 作业启用了“spark.parquet.strings.signed-min-max.enabled”属性。但是,我仍然看到没有为字符串列显示任何统计信息:“VLE:PLAIN DICTIONARY ST:[no stats for this column]”
    【解决方案2】:

    这已在 Spark-2.4.0 版本中解决。在这里,他们将镶木地板版本从 1.8.2 升级到 1.10.0。

    [SPARK-23972] 将 Parquet 从 1.8.2 更新到 1.10.0

    对于这些所有列类型,无论它们是 Int/String/Decimal 都将包含 min/max 统计信息。

    【讨论】:

      【解决方案3】:

      对于第一个问题,我认为这是一个定义问题(字符串的最小值/最大值是多少?词法排序?)但据我所知,无论如何,spark 的镶木地板目前只索引数字。

      至于第二个问题,我相信如果你看得更深,你会发现 spark 本身并没有加载文件。相反,它正在读取元数据,因此它知道是否要读取一个块。所以基本上它是将谓词推到文件(块)级别。

      【讨论】:

        猜你喜欢
        • 2021-06-19
        • 2016-09-02
        • 2019-12-26
        • 2015-11-30
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-11-10
        • 2018-07-08
        相关资源
        最近更新 更多