【问题标题】:Partition column disappears in result set dataframe Spark分区列在结果集数据框 Spark 中消失
【发布时间】:2019-05-30 18:21:24
【问题描述】:

我尝试通过时间戳列update_database_time 拆分 Spark 数据帧,并将其写入具有定义的 Avro 架构的 HDFS。但是,在调用重新分区方法后,我得到了这个异常:

Caused by: org.apache.spark.sql.avro.IncompatibleSchemaException: Cannot convert Catalyst type StructType(StructField(random_pk,DecimalType(38,0),true), StructField(random_string,StringType,true), StructField(code,StringType,true), StructField(random_bool,BooleanType,true), StructField(random_int,IntegerType,true), StructField(random_float,DoubleType,true), StructField(random_double,DoubleType,true), StructField(random_enum,StringType,true), StructField(random_date,DateType,true), StructField(random_decimal,DecimalType(4,2),true), StructField(update_database_time_tz,TimestampType,true), StructField(random_money,DecimalType(19,4),true)) to Avro type {"type":"record","name":"TestData","namespace":"DWH","fields":[{"name":"random_pk","type":["null",{"type":"bytes","logicalType":"decimal","precision":38,"scale":0}]},{"name":"random_string","type":["string","null"]},{"name":"code","type":["string","null"]},{"name":"random_bool","type":["boolean","null"]},{"name":"random_int","type":["int","null"]},{"name":"random_float","type":["double","null"]},{"name":"random_double","type":["double","null"]},{"name":"random_enum","type":["null",{"type":"enum","name":"enumType","symbols":["VAL_1","VAL_2","VAL_3"]}]},{"name":"random_date","type":["null",{"type":"int","logicalType":"date"}]},{"name":"random_decimal","type":["null",{"type":"bytes","logicalType":"decimal","precision":4,"scale":2}]},{"name":"update_database_time","type":["null",{"type":"long","logicalType":"timestamp-millis"}]},{"name":"update_database_time_tz","type":["null",{"type":"long","logicalType":"timestamp-millis"}]},{"name":"random_money","type":["null",{"type":"bytes","logicalType":"decimal","precision":19,"scale":4}]}]}.

我假设用于分区的列在结果中消失了。如何重新定义操作以使其不会发生?

这是我使用的代码:

    dataDF.write
      .partitionBy("update_database_time")
      .format("avro")
      .option(
        "avroSchema",
        SchemaRegistry.getSchema(
          schemaRegistryConfig.url,
          schemaRegistryConfig.dataSchemaSubject,
          schemaRegistryConfig.dataSchemaVersion))
  .save(s"${hdfsURL}${pathToSave}")

【问题讨论】:

  • 通常分区列不应该是要保存的架构的一部分。在 hdfs 中,保存功能将使用该信息创建文件夹。如果你在 hive 中读取这个 avro 文件,它应该根据文件夹结构创建 update_database_time 列的表示。
  • 问题是我需要用分区数据重命名目录。因此,当我要在 Hive 中读回它时,目录名称将没有列名:update_database_time_2019-04-12 而我将使用 2019-04-12 作为目录名称。对于这个问题,我应该使用其他东西而不是分区依据吗?
  • 如果您启用了 hivesupport,您可以直接将带有分区文件夹的 DF 写入 hive 元存储。为什么要先写入 hdfs 然后重命名文件夹?
  • 问题是管道的构建方式是我需要先写入 HDFS,然后才将数据推送到 Hive 元存储
  • 什么是推送到 hive 的过程以及通过该过程推送的其他内容的一般文件结构是什么?如果您不关心使用这些字段拆分数据,还有其他(通常是更好的)方法可以对数据进行分区。通常,partitionby 用于已包含在数据中的字段,因为它为您的数据创建逻辑上可理解的分区,这些分区易于在磁盘上导航。但你不必那样做。

标签: scala apache-spark apache-spark-sql avro


【解决方案1】:

根据您提供的例外情况,该错误似乎源于获取的 AVRO 架构和 Spark 架构之间的不兼容架构。快速浏览一下,最令人担忧的部分可能是这些:

  1. (可能催化剂不知道如何将字符串转换为枚举类型)

Spark 架构:

StructField(random_enum,StringType,true)

AVRO 架构:

{
      "name": "random_enum",
      "type": [
        "null",
        {
          "type": "enum",
          "name": "enumType",
          "symbols": [
            "VAL_1",
            "VAL_2",
            "VAL_3"
          ]
        }
      ]
    }
  1. update_databse_time_tz 在数据帧的架构中只出现一次,但在 AVRO 架构中出现两次)

Spark 架构:

StructField(update_database_time_tz,TimestampType,true)

AVRO 架构:

{
      "name": "update_database_time",
      "type": [
        "null",
        {
          "type": "long",
          "logicalType": "timestamp-millis"
        }
      ]
    },
    {
      "name": "update_database_time_tz",
      "type": [
        "null",
        {
          "type": "long",
          "logicalType": "timestamp-millis"
        }
      ]
    }

我建议先整合架构并消除该异常,然后再解决其他可能的分区问题。

编辑:关于 2 号,我错过了 AVRO 架构中有不同名称的面孔,这导致数据帧中缺少列 update_database_time 的问题。

【讨论】:

  • 首先,我不确定random_enum问题,我会调查一下。但是,update_database_time 和 update_database_time_tz(tz 代表时区)是不同的字段。所以第二个出现在数据帧和 Avro 模式中,第一个只出现在 Avro 模式中,而不出现在数据帧中。
  • 哦,你完全正确,我错过了名称差异,抱歉:) 我会更新帖子
猜你喜欢
  • 2018-07-15
  • 2016-05-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-01-12
  • 2018-03-23
相关资源
最近更新 更多