【问题标题】:Use spark SQL to read non-existing column with Parquet format使用 Spark SQL 读取 Parquet 格式的不存在列
【发布时间】:2018-03-05 23:22:38
【问题描述】:

我有两个月的镶木地板文件 2017_01.parquet 和 2017_08.parquet,这些架构是:

2017_01.parquet:

root
|-- value: struct (nullable = true)
|    |-- version: struct (nullable = true)
|    |    |-- major: integer (nullable = true)
|    |    |-- minor: integer (nullable = true)
|    |-- guid: string (nullable = true)

2017_08.parquet:

root
|-- value: struct (nullable = true)
|    |-- version: struct (nullable = true)
|    |    |-- major: integer (nullable = true)
|    |    |-- minor: integer (nullable = true)
|    |    |-- vnum: integer (nullable = true)
|    |-- guid: string (nullable = true)

还有我的代码

SQL = """
SELECT value.version.major,
       value.version.minor,
       value.version.vnum
FROM OUT_TABLE 
LIMIT 10"""

parquetFile = spark.read.parquet("/mydata/2017_08.parquet")
parquetFile.createOrReplaceTempView("OUT_TABLE")
out_osce = spark.sql(SQL)
out_osce.show()

当我加载 2017_08.parquet show 时:

+-----+-----+----+
|major|minor|vnum|
+-----+-----+----+
| 0001| 4610|1315|
| 0002| 4610|6206|
| 0003| 4610|6125|

但如果我加载 2017_01.parquet 之类的 parquetFile = spark.read.parquet("/mydata/2017_01.parquet")

SQL 显示错误:

pyspark.sql.utils.AnalysisException: u'No such struct field vnum in major, minor; line 4 pos 11'

我知道原因是2017_01.parquet没有vnum列,我有两个slove解决方案,一个是使用mergeSchema另一个是在读取parquet文件时使用schema,但是这些方式也有很大的问题。

第一个解决方案需要读取2017_08.parquet,如果我不需要08的数据就会有问题,如果运气不好vnum是一个选项列而08没有这个列仍然会出错

第二种解决方案是读取时给出schema,如spark.read.schema(schema).parquet("/mydata/2017_01.parquet"),这种方式需要先写入schema,但如果文件是一个非常复杂的嵌套表,用户可能无法写入schema,并且schema会更新。

我想问任何人有第三种解决方案,然后只阅读 2017_01.parquet 并输出如下:

+-----+-----+----+
|major|minor|vnum|
+-----+-----+----+
| 0001| 4600|null|
| 0002| 4600|null|
| 0003| 4600|null|

【问题讨论】:

  • 感谢编辑建议@himanshuIIITian

标签: python apache-spark pyspark apache-spark-sql parquet


【解决方案1】:

阅读时可以简单地使用 case 语句或合并:

parquetFile = spark.read.parquet("") \
              .withColumn("vnum", coalesce("vnum"))

来自文档:

coalesce(e: Column*): 列

返回不为空的第一列,如果所有输入都为空,则返回 空。

如果您的 Parquet 文件有此字段,则会使用该字段。如果没有,将使用空值,并且新列将在您的架构中

【讨论】:

  • 嗨 T. Gawęda 当我使用 coalesce("vnum") 时它总是显示错误,所以我尝试使用 lit(""),但问题是 .withColumn 新列添加顶层,如果我有嵌套模式,则无法将列添加到正确的级别
  • 这不是一个稳定的解决方案,因为我理解只有当它出现在其他行中时才添加的列,但如果根本没有列,这不起作用。
【解决方案2】:

我可以通过在创建选择时检查 DF 的列列表来解决类似的问题。 在我的情况下,以下就足够了:

 parquetFile = spark.read.parquet("").withColumn("vnum", coalesce(if
 parquetFile.columns.contains("vnum") $"vnum" else lit(null)))

在您的情况下,使用嵌套架构,您可以使用以下内容:

// Define the full struct type schema to check if nested field exists.

val structToCheck = new StructField("value", new StructType().add("version",new StructType().add("major",StringType).add("minor",StringType).add("vnum",StringType)))

val SQL = """ SELECT value.version.major,
       value.version.minor,""" +
       if (parquetFile.schema.contains(structToCheck)) 
           "value.version.vnum" 
       else 
           "'' as vnum" +
       "FROM OUT_TABLE  LIMIT 10"

您还可以进行一些更具体的搜索,获取 value.version 结构并检查其元素。

【讨论】:

  • 嗨 Vapira,这是示例模式的一个很好的解决方案,但在实际情况下,我的模式非常庞大,我不知道它什么时候添加新列(数据提供者是其他人),所以使用检查功能可能有其他问题。谢谢
【解决方案3】:

您可以创建一个索引表来存储每个parquet文件的列情况,例如:

parquetfile1 列 1,列 2 parquetfile2 第 1 列,第 3 列 .....

读取parquet文件时,先读取索引数据过滤部分文件,然后执行查询操作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-01-27
    • 1970-01-01
    • 1970-01-01
    • 2020-05-17
    • 2018-10-04
    • 1970-01-01
    • 2019-03-16
    • 2020-02-10
    相关资源
    最近更新 更多