【问题标题】:Re-using A Schema from JSON within a Spark DataFrame using Scala使用 Scala 在 Spark DataFrame 中重用来自 JSON 的模式
【发布时间】:2016-08-12 09:59:05
【问题描述】:

我有一些这样的 JSON 数据:

{"gid":"111","createHour":"2014-10-20 01:00:00.0","revisions":[{"revId":"2","modDate":"2014-11-20 01:40:37.0"},{"revId":"4","modDate":"2014-11-20 01:40:40.0"}],"comments":[],"replies":[]}
{"gid":"222","createHour":"2014-12-20 01:00:00.0","revisions":[{"revId":"2","modDate":"2014-11-20 01:39:31.0"},{"revId":"4","modDate":"2014-11-20 01:39:34.0"}],"comments":[],"replies":[]}
{"gid":"333","createHour":"2015-01-21 00:00:00.0","revisions":[{"revId":"25","modDate":"2014-11-21 00:34:53.0"},{"revId":"110","modDate":"2014-11-21 00:47:10.0"}],"comments":[{"comId":"4432","content":"How are you?"}],"replies":[{"repId":"4441","content":"I am good."}]}
{"gid":"444","createHour":"2015-09-20 23:00:00.0","revisions":[{"revId":"2","modDate":"2014-11-20 23:23:47.0"}],"comments":[],"replies":[]}
{"gid":"555","createHour":"2016-01-21 01:00:00.0","revisions":[{"revId":"135","modDate":"2014-11-21 01:01:58.0"}],"comments":[],"replies":[]}
{"gid":"666","createHour":"2016-04-23 19:00:00.0","revisions":[{"revId":"136","modDate":"2014-11-23 19:50:51.0"}],"comments":[],"replies":[]}

我可以读到:

val df = sqlContext.read.json("./data/full.json")

我可以使用df.printSchema 打印架构

root
 |-- comments: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- comId: string (nullable = true)
 |    |    |-- content: string (nullable = true)
 |-- createHour: string (nullable = true)
 |-- gid: string (nullable = true)
 |-- replies: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- content: string (nullable = true)
 |    |    |-- repId: string (nullable = true)
 |-- revisions: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- modDate: string (nullable = true)
 |    |    |-- revId: string (nullable = true)

我可以显示数据df.show(10,false)

+---------------------+---------------------+---+-------------------+---------------------------------------------------------+
|comments             |createHour           |gid|replies            |revisions                                                |
+---------------------+---------------------+---+-------------------+---------------------------------------------------------+
|[]                   |2014-10-20 01:00:00.0|111|[]                 |[[2014-11-20 01:40:37.0,2], [2014-11-20 01:40:40.0,4]]   |
|[]                   |2014-12-20 01:00:00.0|222|[]                 |[[2014-11-20 01:39:31.0,2], [2014-11-20 01:39:34.0,4]]   |
|[[4432,How are you?]]|2015-01-21 00:00:00.0|333|[[I am good.,4441]]|[[2014-11-21 00:34:53.0,25], [2014-11-21 00:47:10.0,110]]|
|[]                   |2015-09-20 23:00:00.0|444|[]                 |[[2014-11-20 23:23:47.0,2]]                              |
|[]                   |2016-01-21 01:00:00.0|555|[]                 |[[2014-11-21 01:01:58.0,135]]                            |
|[]                   |2016-04-23 19:00:00.0|666|[]                 |[[2014-11-23 19:50:51.0,136]]                            |
+---------------------+---------------------+---+-------------------+---------------------------------------------------------+

我可以打印/读取架构val dfSc = df.schema

StructType(StructField(comments,ArrayType(StructType(StructField(comId,StringType,true), StructField(content,StringType,true)),true),true), StructField(createHour,StringType,true), StructField(gid,StringType,true), StructField(replies,ArrayType(StructType(StructField(content,StringType,true), StructField(repId,StringType,true)),true),true), StructField(revisions,ArrayType(StructType(StructField(modDate,StringType,true), StructField(revId,StringType,true)),true),true))

我可以更好地打印出来:

println(df.schema.fields.mkString(",\n"))
StructField(comments,ArrayType(StructType(StructField(comId,StringType,true), StructField(content,StringType,true)),true),true),
StructField(createHour,StringType,true),
StructField(gid,StringType,true),
StructField(replies,ArrayType(StructType(StructField(content,StringType,true), StructField(repId,StringType,true)),true),true),
StructField(revisions,ArrayType(StructType(StructField(modDate,StringType,true), StructField(revId,StringType,true)),true),true)

现在,如果我在没有commentsreplies 行的情况下读取同一个文件,而val df2 = sqlContext.read. json("./data/partialRevOnly.json") 只是删除这些行,我会在printSchema 中得到类似的结果:

root
 |-- comments: array (nullable = true)
 |    |-- element: string (containsNull = true)
 |-- createHour: string (nullable = true)
 |-- gid: string (nullable = true)
 |-- replies: array (nullable = true)
 |    |-- element: string (containsNull = true)
 |-- revisions: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- modDate: string (nullable = true)
 |    |    |-- revId: string (nullable = true)

我不喜欢这样,所以我使用:

val df3 = sqlContext.read.
  schema(dfSc).
  json("./data/partialRevOnly.json")

原始架构是dfSc。所以现在我得到了我之前删除数据的模式:

root
 |-- comments: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- comId: string (nullable = true)
 |    |    |-- content: string (nullable = true)
 |-- createHour: string (nullable = true)
 |-- gid: string (nullable = true)
 |-- replies: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- content: string (nullable = true)
 |    |    |-- repId: string (nullable = true)
 |-- revisions: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- modDate: string (nullable = true)
 |    |    |-- revId: string (nullable = true)

这太完美了……差不多。我想将此模式分配给与此类似的变量:

val textSc =  StructField(comments,ArrayType(StructType(StructField(comId,StringType,true), StructField(content,StringType,true)),true),true),
    StructField(createHour,StringType,true),
    StructField(gid,StringType,true),
    StructField(replies,ArrayType(StructType(StructField(content,StringType,true), StructField(repId,StringType,true)),true),true),
    StructField(revisions,ArrayType(StructType(StructField(modDate,StringType,true), StructField(revId,StringType,true)),true),true)

好的 - 由于双引号和“其他一些结构性”的东西,这不起作用,所以试试这个(有错误):

import org.apache.spark.sql.types._

val textSc = StructType(Array(
    StructField("comments",ArrayType(StructType(StructField("comId",StringType,true), StructField("content",StringType,true)),true),true),
    StructField("createHour",StringType,true),
    StructField("gid",StringType,true),
    StructField("replies",ArrayType(StructType(StructField("content",StringType,true), StructField("repId",StringType,true)),true),true),
    StructField("revisions",ArrayType(StructType(StructField("modDate",StringType,true), StructField("revId",StringType,true)),true),true)
))

Name: Compile Error
Message: <console>:78: error: overloaded method value apply with alternatives:
  (fields: Array[org.apache.spark.sql.types.StructField])org.apache.spark.sql.types.StructType <and>
  (fields: java.util.List[org.apache.spark.sql.types.StructField])org.apache.spark.sql.types.StructType <and>
  (fields: Seq[org.apache.spark.sql.types.StructField])org.apache.spark.sql.types.StructType
 cannot be applied to (org.apache.spark.sql.types.StructField, org.apache.spark.sql.types.StructField)
           StructField("comments",ArrayType(StructType(StructField("comId",StringType,true), StructField("content",StringType,true)),true),true),

... 如果没有这个错误(我想不出一个快速的解决方法),我想使用textSc 代替dfSc 来读取带有强制模式的 JSON 数据。

我找不到“一对一匹配”的方式来获取(通过 println 或 ...)具有可接受语法(类似于上面)的模式。我想一些编码可以通过大小写匹配来消除双引号。但是,我仍然不清楚需要什么规则才能从测试夹具中获取确切的架构,我可以简单地在我的重复生产(相对于测试夹具)代码中重复使用。有没有办法让这个模式完全像我编码的那样打印?

注意:这包括双引号和所有正确的 StructField/Types 等等,以便与代码兼容。

作为侧边栏,我曾考虑保存一个完整的黄金 JSON 文件以在 Spark 作业开始时使用,但我希望最终在适用的结构位置使用日期字段和其他更简洁的类型而不是字符串.

如何从我的测试工具中获取 dataFrame 信息(使用带有 cmets 和回复的完整 JSON 输入行)到可以将架构作为源代码放入生产代码 Scala Spark 作业的程度?

注意:最好的答案是一些编码方法,但是解释一下,这样我就可以通过编码来跋涉、沉闷、辛劳、涉水、耕耘和跋涉,这也很有帮助。 :)

【问题讨论】:

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


【解决方案1】:

我最近遇到了这个问题。我使用的是 Spark 2.0.2,所以我不知道这个解决方案是否适用于早期版本。

import scala.util.Try
import org.apache.spark.sql.Dataset
import org.apache.spark.sql.catalyst.parser.LegacyTypeStringParser
import org.apache.spark.sql.types.{DataType, StructType}

/** Produce a Schema string from a Dataset */
def serializeSchema(ds: Dataset[_]): String = ds.schema.json

/** Produce a StructType schema object from a JSON string */
def deserializeSchema(json: String): StructType = {
    Try(DataType.fromJson(json)).getOrElse(LegacyTypeStringParser.parse(json)) match {
        case t: StructType => t
        case _ => throw new RuntimeException(s"Failed parsing StructType: $json")
    }
}

请注意,我刚刚从 Spark StructType 对象的私有函数中复制了“反序列化”函数。不知道跨版本支持的好不好。

【讨论】:

  • 我已经在 Spark 2.1.0 上试过了,效果很好。只需导入这些:import org.apache.spark.sql.Datasetimport org.apache.spark.sql.catalyst.parser.LegacyTypeStringParserimport scala.util.Tryval caseclassstring = """StructType(Array(StructField(comments,ArrayType(StructType(List(StructField(comId,DateType,true) ... """ deserializeSchema(caseclassstring)
【解决方案2】:

嗯,错误消息应该告诉你这里你必须知道的一切——StructType 需要一个字段序列作为参数。因此,在您的情况下,架构应如下所示:

StructType(Seq(
  StructField("comments", ArrayType(StructType(Seq(       // <- Seq[StructField]
    StructField("comId", StringType, true),
    StructField("content", StringType, true))), true), true), 
  StructField("createHour", StringType, true),
  StructField("gid", StringType, true),
  StructField("replies", ArrayType(StructType(Seq(        // <- Seq[StructField]
    StructField("content", StringType, true),
    StructField("repId", StringType, true))), true), true),
  StructField("revisions", ArrayType(StructType(Seq(      // <- Seq[StructField]
    StructField("modDate", StringType, true),
    StructField("revId", StringType, true))),true), true)))

【讨论】:

  • 好的 - 我承认错误信息对我来说有点神秘。我已经测试了你的建议并且它有效。我仍在寻找一种编程方式来构建模式。在您的帮助下,我现在知道如何手工制作模式......这是一大步——尤其是当需要调整模式以使 date 等列比 strings 更具代表性时。感谢您的帮助。
  • 这个错误肯定很神秘。 Scala 中的类型和其他类可以看起来像函数,它们定义了一个apply 方法。附加 .apply 是可选的。例如StructType.apply()StructType() 相同。这解释了错误“值应用方法”。 “重载”是指具有多个咒语的 StructType 应用方法(在技术上也是构造函数)。它可以容纳SeqArrayList。其中任何一个都可以在这里工作。 Seq 是一个不错的选择,它更通用,不需要 List 或 Array 的额外特征。
猜你喜欢
  • 1970-01-01
  • 2018-01-02
  • 2020-10-02
  • 2020-07-24
  • 1970-01-01
  • 2022-12-18
  • 2019-09-26
  • 2019-10-22
  • 1970-01-01
相关资源
最近更新 更多