【问题标题】:Convert a string variable in nested JSON to datetime using Spark Scala使用 Spark Scala 将嵌套 JSON 中的字符串变量转换为日期时间
【发布时间】:2017-08-05 10:27:18
【问题描述】:

我在 Spark 中有一个嵌套的 JSON 数据框,如下所示

root
 |-- data: struct (nullable = true)
 |    |-- average: long (nullable = true)
 |    |-- sum: long (nullable = true)
 |    |-- time: string (nullable = true)
 |-- password: string (nullable = true)
 |-- url: string (nullable = true)
 |-- username: string (nullable = true)

我需要将数据结构下的时间变量转换为时间戳数据类型。以下是我尝试过的代码,但没有给我想要的结果。

val jsonStr = """{
"url": "imap.yahoo.com",
"username": "myusername",
"password": "mypassword",
"data": {
"time":"2017-1-29 0-54-32",
"average": 234,
"sum": 123}}"""


  val json: JsValue = Json.parse(jsonStr)

  import sqlContext.implicits._
  val rdd = sc.parallelize(jsonStr::Nil);
  var df = sqlContext.read.json(rdd);
  df.printSchema()
  val dfRes = df.withColumn("data",makeTimeStamp(unix_timestamp(df("data.time"),"yyyy-MM-dd hh-mm-ss").cast("timestamp")))
  dfRes.printSchema();

case class Convert(time: java.sql.Timestamp)
val makeTimeStamp = udf((time: java.sql.Timestamp) => Convert(
  time))

我的代码结果:

root
 |-- data: struct (nullable = true)
 |    |-- time: timestamp (nullable = true)
 |-- password: string (nullable = true)
 |-- url: string (nullable = true)
 |-- username: string (nullable = true)

我的代码实际上是删除数据结构中的其他元素(平均值和总和),而不是仅仅将时间字符串转换为时间戳数据类型。对于 JSON 数据帧上的基本数据管理操作,我们是否需要在需要功能时编写 UDF,或者是否有可用于 JSON 数据管理的库。我目前正在使用 Play 框架来处理 Spark 中的 JSON 对象。提前致谢。

【问题讨论】:

    标签: json scala apache-spark


    【解决方案1】:

    你可以试试这个:

    val jsonStr = """{
    "url": "imap.yahoo.com",
    "username": "myusername",
    "password": "mypassword",
    "data": {
    "time":"2017-1-29 0-54-32",
    "average": 234,
    "sum": 123}}"""
    
    
    val json: JsValue = Json.parse(jsonStr)
    
    import sqlContext.implicits._
    val rdd = sc.parallelize(jsonStr::Nil);
    var df = sqlContext.read.json(rdd);
    df.printSchema()
    val dfRes = df.withColumn("data",makeTimeStamp(unix_timestamp(df("data.time"),"yyyy-MM-dd hh-mm-ss").cast("timestamp"), df("data.average"), df("data.sum")))
    
    case class Convert(time: java.sql.Timestamp, average: Long, sum: Long)
    val makeTimeStamp = udf((time: java.sql.Timestamp, average: Long, sum: Long) => Convert(time, average, sum))
    

    这将给出结果:

    root
    |-- url: string (nullable = true)
    |-- username: string (nullable = true)
    |-- password: string (nullable = true)
    |-- data: struct (nullable = true)
    |    |-- time: timestamp (nullable = true)
    |    |-- average: long (nullable = false)
    |    |-- sum: long (nullable = false)
    

    唯一改变的是Convert case class 和makeTimeStamp UDF。

    【讨论】:

    • 感谢您的回复。您的解决方案确实有助于保持数据树下的元素完好无损。但是我正在处理的真实场景中,数据结构中有太多元素。将它们全部添加到案例类不是一个可行的解决方案。数据结构内的元素数量也可能随时间变化。我正在寻找一种可以轻松扩展的解决方案。
    【解决方案2】:

    假设您可以预先指定 Spark 架构,自动字符串到时间戳类型强制转换应该负责转换。

    import org.apache.spark.sql.types._
    val dschema = (new StructType).add("url", StringType).add("username", StringType).add
               ("data", (new StructType).add("sum", LongType).add("time", TimestampType))
    val df = spark.read.schema(dschema).json("/your/json/on/hdfs")
    df.printSchema
    df.show
    

    This article 概述了更多处理不良数据的技术;值得一读您的用例。

    【讨论】:

      猜你喜欢
      • 2021-12-21
      • 1970-01-01
      • 1970-01-01
      • 2020-09-02
      • 2019-09-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多