【问题标题】:Spark read Parquet files of different versionsSpark读取不同版本的Parquet文件
【发布时间】:2017-09-25 21:23:24
【问题描述】:
我使用 Version1 架构生成了一年多的镶木地板文件。并且随着最近的架构更改,较新的镶木地板文件具有 Version2 架构额外列。
所以当我从旧版本和新版本一起加载 parquet 文件并尝试过滤更改的列时,我得到了一个异常。
我希望 spark 读取新旧文件并在列不存在的情况下填充空值。是否有解决方法,在未找到列时 spark 填充空值?
【问题讨论】:
标签:
apache-spark
parquet
versions
【解决方案1】:
SparkSQL 本身支持 parquet 文件的模式合并。你可以在official documentation here阅读所有相关信息
与 ProtocolBuffer、Avro 和 Thrift 一样,Parquet 也支持模式
进化。用户可以从一个简单的模式开始,逐步添加
根据需要向架构添加更多列。这样,用户最终可能会
具有不同但相互兼容的多个 Parquet 文件
模式。 Parquet 数据源现在能够自动检测
这种情况并合并所有这些文件的模式。
由于模式合并是一项相对昂贵的操作,而不是
在大多数情况下,我们默认将其关闭,从
1.5.0。您可以通过
启用它
在读取 Parquet 时将数据源选项 mergeSchema 设置为 true
文件(如下例所示),或
将全局 SQL 选项 spark.sql.parquet.mergeSchema 设置为 true。
【解决方案2】:
有两种方法你可以试试。
1.喜欢这种方式可以使用地图变换,但不推荐,如spark.read.parquet("mypath").map(e => val field =if (e.isNullAt(e.fieldIndex("field"))) null else e.getAs[String]("field"))
2.使用mergeSchema选项的最佳方式,例如:
spark.read.option("mergeSchema", "true").parquet(xxx).as[MyClass]
参考:schema-merging
【解决方案3】:
假设您有一组要阅读的文件:
- 查询集合中每个文件的架构,生成 N 组文件,每组包含具有相似架构的文件。
- 使用与每组中的架构兼容的过滤器对每组文件进行操作。
- 联合过滤/操作每个集合的结果(假设您的输出对于每个文件架构的结果都相同)