【问题标题】: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。您可以通过

    启用它
    1. 在读取 Parquet 时将数据源选项 mergeSchema 设置为 true 文件(如下例所示),或

    2. 将全局 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】:

        假设您有一组要阅读的文件:

        1. 查询集合中每个文件的架构,生成 N 组文件,每组包含具有相似架构的文件。
        2. 使用与每组中的架构兼容的过滤器对每组文件进行操作。
        3. 联合过滤/操作每个集合的结果(假设您的输出对于每个文件架构的结果都相同)

        【讨论】:

          猜你喜欢
          • 2015-12-19
          • 2016-12-15
          • 2020-10-28
          • 2015-08-05
          • 1970-01-01
          • 1970-01-01
          • 2016-02-06
          • 2017-04-24
          • 2018-10-04
          相关资源
          最近更新 更多