【问题标题】:DataFrame partitionBy on nested columns嵌套列上的 DataFrame partitionBy
【发布时间】:2016-11-14 03:27:42
【问题描述】:

我正在尝试在如下嵌套字段上调用 ​​partitionBy:

val rawJson = sqlContext.read.json(filename)
rawJson.write.partitionBy("data.dataDetails.name").parquet(filenameParquet)

我在运行它时收到以下错误。我确实看到“名称”列为以下架构中的字段。是否有不同的格式来指定嵌套的列名?

java.lang.RuntimeException: 在模式 StructType(StructField(name,StringType,true), StructField(time,StringType,true), StructField(data,StructType(StructField(dataDetails, StructType(StructField(name,StringType,true), StructField(id,StringType,true),true)),true))

这是我的 json 文件:

{  
  "name": "AssetName",
  "time": "2016-06-20T11:57:19.4941368-04:00",
  "data": {
    "type": "EventData",
    "dataDetails": {
      "name": "EventName"
      "id": "1234"

    }
  }
} 

【问题讨论】:

  • 我也遇到了同样的问题,你解决了吗?

标签: apache-spark apache-spark-sql spark-dataframe


【解决方案1】:

由于该功能在 Spark 2.3.1 中不可用,因此这里有一个解决方法。确保处理嵌套字段和根级别字段之间的名称冲突。

{"date":"20180808","value":{"group":"xxx","team":"yyy"}}
df.select("date","value.group","value.team")
      .write
      .partitionBy("date","group","team")
      .parquet(filenameParquet)

分区结束

date=20180808/group=xxx/team=yyy/part-xxx.parquet

【讨论】:

    【解决方案2】:

    这似乎是此处列出的已知问题:https://issues.apache.org/jira/browse/SPARK-18084

    我也遇到了这个问题,为了解决这个问题,我能够取消嵌套数据集上的列。我的数据集与您的数据集略有不同,但这里是策略...

    原始Json:

    {  
      "name": "AssetName",
      "time": "2016-06-20T11:57:19.4941368-04:00",
      "data": {
        "type": "EventData",
        "dataDetails": {
          "name": "EventName"
          "id": "1234"
    
        }
      }
    } 
    

    修改后的Json:

    {  
      "name": "AssetName",
      "time": "2016-06-20T11:57:19.4941368-04:00",
      "data_type": "EventData",
      "data_dataDetails_name" : "EventName",
      "data_dataDetails_id": "1234"
      }
    } 
    

    获取修改后 Json 的代码:

    def main(args: Array[String]) {
      ...
    
      val data = df.select(children("data", df) ++ $"name" ++ $"time"): _*)
    
      data.printSchema
    
      data.write.partitionBy("data_dataDetails_name").format("csv").save(...)
    }
    
    def children(colname: String, df: DataFrame) = {
      val parent = df.schema.fields.filter(_.name == colname).head
      val fields = parent.dataType match {
        case x: StructType => x.fields
        case _ => Array.empty[StructField]
      }
      fields.map(x => col(s"$colname.${x.name}").alias(s"$colname" + s"_" + s"${x.name}"))
    }
    

    【讨论】:

      猜你喜欢
      • 2018-11-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-07-24
      • 1970-01-01
      • 2017-12-26
      • 2015-12-20
      相关资源
      最近更新 更多