【问题标题】:Retain raw JSON as column in Spark DataFrame on read/load?在读取/加载时保留原始 JSON 作为 Spark DataFrame 中的列?
【发布时间】:2018-10-17 11:08:35
【问题描述】:

在将我的数据读入 Spark DataFrame 时,我一直在寻找一种将原始 (JSON) 数据添加为列的方法。我有一种方法可以通过加入来做到这一点,但我希望有一种方法可以使用 Spark 2.2.x+ 在单个操作中做到这一点。

例如数据:

{"team":"Golden Knights","colors":"gold,red,black","origin":"Las Vegas"}
{"team":"Sharks","origin": "San Jose", "eliminated":"true"}
{"team":"Wild","colors":"red,green,gold","origin":"Minnesota"}

执行时:

val logs = sc.textFile("/Users/vgk/data/tiny.json") // example data file
spark.read.json(logs).show

可以预见的是:

+--------------+----------+--------------------+--------------+
|        colors|eliminated|              origin|          team|
+--------------+----------+--------------------+--------------+
|gold,red,black|      null|           Las Vegas|Golden Knights|
|          null|      true|            San Jose|        Sharks|
|red,green,gold|      null|           Minnesota|          Wild|
|red,white,blue|     false|District of Columbia|      Capitals|
+--------------+----------+--------------------+--------------+

我希望在初始加载时拥有上述内容,但将原始 JSON 数据作为附加列。例如(截断的原始值):

+--------------+-------------------------------+--------------+--------------------+
|        colors|eliminated|              origin|          team|               value|
+--------------+----------+--------------------+--------------+--------------------+
|red,white,blue|     false|District of Columbia|      Capitals|{"colors":"red,wh...|
|gold,red,black|      null|           Las Vegas|Golden Knights|{"colors":"gold,r...|
|          null|      true|            San Jose|        Sharks|{"eliminated":"tr...|
|red,green,gold|      null|           Minnesota|          Wild|{"colors":"red,gr...|
+--------------+----------+--------------------+--------------+--------------------+

非理想解决方案涉及连接:

val logs = sc.textFile("/Users/vgk/data/tiny.json")
val df = spark.read.json(logs).withColumn("uniqueID",monotonically_increasing_id)
val rawdf = df.toJSON.withColumn("uniqueID",monotonically_increasing_id)
df.join(rawdf, "uniqueID")

这会产生与上面相同的数据框,但添加了 uniqueID 列。此外,json 是从 DF 呈现的,不一定是“原始”数据。实际上它们是相等的,但对于我的用例,实际的原始数据更可取。

是否有人知道将原始 JSON 数据捕获为加载时的附加列的解决方案?

【问题讨论】:

  • 另一个不理想的解决方案是从Row 到它的 JSON 表示进行内联转换。这可能更容易、更可靠。
  • 我应该补充一点,数据是异构的,不符合持久模式。此外,下游处理更愿意传递未经修改的原始 json 值,而不是基于数据生成的块。
  • 是的。您可以使用 rdd api 并映射到具有所有列的新行对象以及具有整行的 json 表示的新文本字段。使用类似杰克逊的东西应该这样做。

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


【解决方案1】:

如果您有收到的数据的架构,那么您可以使用from_jsonschema 来获取所有字段并保持raw 字段原样

val logs = spark.sparkContext.textFile(path) // example data file

val schema = StructType(
  StructField("team", StringType, true)::
  StructField("colors", StringType, true)::
  StructField("eliminated", StringType, true)::
  StructField("origin", StringType, true)::Nil
)

logs.toDF("values")
    .withColumn("json", from_json($"values", schema))
    .select("values", "json.*")

    .show(false)

输出:

+------------------------------------------------------------------------+--------------+--------------+----------+---------+
|values                                                                  |team          |colors        |eliminated|origin   |
+------------------------------------------------------------------------+--------------+--------------+----------+---------+
|{"team":"Golden Knights","colors":"gold,red,black","origin":"Las Vegas"}|Golden Knights|gold,red,black|null      |Las Vegas|
|{"team":"Sharks","origin": "San Jose", "eliminated":"true"}             |Sharks        |null          |true      |San Jose |
|{"team":"Wild","colors":"red,green,gold","origin":"Minnesota"}          |Wild          |red,green,gold|null      |Minnesota|
+------------------------------------------------------------------------+--------------+--------------+----------+---------+

希望他的帮助!

【讨论】:

  • 很遗憾,源数据是日志流,未定义架构。此外,可以随时在流的日志行中添加或删除新属性。
【解决方案2】:

您可以简单地将to_json 内置函数.withColumn函数结合使用

val logs = sc.textFile("/Users/vgk/data/tiny.json")
val df = spark.read.json(logs)
import org.apache.spark.sql.functions._
df.withColumn("value", to_json(struct(df.columns.map(col): _*))).show(false)

或者更好,不要用sparkContexttextFile读成rdd,直接用sparkSession读成json文件

val df = spark.read.json("/Users/vgk/data/tiny.json")

import org.apache.spark.sql.functions._
df.withColumn("value", to_json(struct(df.columns.map(col): _*))).show(false)

你应该得到

+--------------+----------+---------+--------------+------------------------------------------------------------------------+
|colors        |eliminated|origin   |team          |value                                                                   |
+--------------+----------+---------+--------------+------------------------------------------------------------------------+
|gold,red,black|null      |Las Vegas|Golden Knights|{"colors":"gold,red,black","origin":"Las Vegas","team":"Golden Knights"}|
|null          |true      |San Jose |Sharks        |{"eliminated":"true","origin":"San Jose","team":"Sharks"}               |
|red,green,gold|null      |Minnesota|Wild          |{"colors":"red,green,gold","origin":"Minnesota","team":"Wild"}          |
+--------------+----------+---------+--------------+------------------------------------------------------------------------+

【讨论】:

  • 这似乎确实有效,尽管保留原始原始数据(相对于通过某种机制重新生成数据)的既定目标仍然是一个障碍。我怀疑这种机制会比 uniqueID & join 技术更高效——一些测试应该可以验证。如果是这样的话,我怀疑,那么我们仍然向前迈进了。
  • 我不明白你的评论@reverend的意思
  • 换一种方式,不需要通过toJSONto_json 将json 重新组合在一起,只需在读取步骤期间将原始数据转储到列中即可。使用 toJSON/to_json 可以有效地将原始 json 解析到 DataFrame 中,这似乎不必要地增加了将原始数据保存在列中所需的计算能力。这是一个不寻常的用例,所以尽量优化。
  • 然后将其读取为 text json dataframe 。那应该有一个包含 json 字符串的列。
  • 是的,但是如何扭转解决方案呢?将RDD[String]Dataset[String] 转换为JSON 会产生一个新的Dataset/DataFrame,然后需要将其连接回原始Dataset[String]。同样,没有一致模式的异构数据不允许 json 函数在这种情况下有用。
【解决方案3】:

使用 rdd 映射器读取每一行,并操作字符串以将原始行添加到 json 字符串中,然后将该 rdd 解析到数据帧 json 读取器中。

def addRawToJson(line):
    line = line.strip()
    rawJson = line.replace('\\', '\\\\').replace('"', '\\"')
    linePlusRaw = f'{line[0:len(line)-1]}, "{RAW_JSON_FIELD_NAME}":"{rawJson}"' + '}'
    return linePlusRaw
    
rawAugmentedJsonRdd = sc.textFile('add file path here').map(addRawToJson)
df = spark.read.json(rawAugmentedJsonRdd)

这样取原来的json而不是重新构建,不需要两次读取数据合并,也不需要你提前知道schema。

请注意,我的答案是在 python 中使用 pyspark,但应该很容易更改为使用 scala。

另请注意,该方法假定简单的单行 json 输入,并且在直接操作字符串之前不测试有效的 json,这对于我的用例来说是可以接受的。

【讨论】:

    猜你喜欢
    • 2018-11-13
    • 1970-01-01
    • 2019-07-05
    • 2020-12-21
    • 2018-08-19
    • 1970-01-01
    • 1970-01-01
    • 2019-02-17
    • 1970-01-01
    相关资源
    最近更新 更多