【发布时间】: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