【问题标题】:How to use Partition Discovery in Spark SQL如何在 Spark SQL 中使用分区发现
【发布时间】:2019-01-29 10:45:36
【问题描述】:

鉴于 HDFS 中带有 parquet 文件的下一个结构:

    data
    ├── name=Steve
    └── name=Michael

在 SparkSQL 中,使用查询:

CREATE TABLE test USING parquet OPTIONS (path 'hdfs://namenode:8020/data')

分区未正常恢复,未检测到数据:

SELECT * FROM test LIMIT 1

+---+----+
|ID |name|
+---+----+
+---+----+

但是,还有一种替代方法,即在创建表时指定架构,然后使用 recover partitions

执行 alter table
CREATE TABLE test2(ID Int, name String) USING parquet OPTIONS (path 'hdfs://namenode:8020/data')

ALTER TABLE test2 RECOVER PARTITIONS

SELECT * FROM test2 LIMIT 1

+----+---------+
| ID |    name |
+----+---------+
|  1 |   Steve |
|  2 | Michael |
+----+---------+

Spark SQL 中是否有任何其他替代方法可以在创建表时仅使用一个查询而不指定架构来使用分区发现?

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    通过 spark sql 创建表后如:

    CREATE TABLE test USING parquet OPTIONS (path 'hdfs://namenode:8020/data')
    

    记得在使用前修复桌子:

    MSCK REPAIR TABLE test
    

    然后表的分区将被注册到 Metastore。否则该表将不会返回任何结果。如本在线文档所述:https://docs.databricks.com/user-guide/tables.html

    【讨论】:

      【解决方案2】:

      假设人口数据使用以下目录结构加载到分区表中,并使用两个额外的列,性别和国家作为分区列:

      path
      └── to
          └── table
              ├── gender=male
              │   ├── ...
              │   │
              │   ├── country=US
              │   │   └── data.parquet
              │   ├── country=CN
              │   │   └── data.parquet
              │   └── ...
              └── gender=female
                  ├── ...
                  │
                  ├── country=US
                  │   └── data.parquet
                  ├── country=CN
                  │   └── data.parquet
                  └── ...
      

      通过将 path/to/table 传递给 SparkSession.read.parquet 或 SparkSession.read.load,Spark SQL 将自动从路径中提取分区信息。现在返回的 DataFrame 的 schema 变为:

      root
      |-- name: string (nullable = true)
      |-- age: long (nullable = true)
      |-- gender: string (nullable = true)
      |-- country: string (nullable = true)
      

      注意:当您使用spark.read.parquet("path/to/table/*") 读取分区数据并运行df.show() 时,您会发现缺少两个额外的列(性别和国家/地区)。为确保将它们添加到您的数据框中,您必须将 basePath 选项添加为read.option("basePath", "path/to/table").parquet("path/to/table/*"),或者您可以将路径指定为spark.read.parquet("path/to/table")

      如果使用 spark.sql,则 Hive parquet - 提供设置正常,而不是在需要修复的 S3 上,然后 spark.sql 根据问题提供这样的分区感知。

      例如一个简单的例子(由于某种原因无法格式化):

      dfX.write.partitionBy("col2").format("parquet").saveAsTable("dfX_partitionBy_Table")

      但在你的问题中,如果这里不需要,那么答案是否定的。

      【讨论】:

      • 是不是说分区发现不能和SQL查询一起使用,必须使用dataframe API?
      • 上述方法适用于 spark.read... 但是, dfX.write.partitionBy("col2").format("parquet").saveAsTable("dfX_partitionBy_Table") 应该允许分区修剪发生。这是你要找的东西,那我会改编答案。
      • 问题是我需要使用HDFS中的分区发现只使用一个SQL查询
      • 但是如果你已经正确分区,那么它应该可以工作。我提供了 col c2 的 parquet table 示例分区。它在 Hive Metastore 中,所以应该可以工作。否则你需要更清楚地说明问题。
      • 如果 Hive 支持被禁用怎么办?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-11-10
      • 1970-01-01
      • 2016-03-06
      • 1970-01-01
      • 1970-01-01
      • 2020-10-11
      • 2016-10-14
      相关资源
      最近更新 更多