【发布时间】:2016-04-11 18:35:44
【问题描述】:
我正在尝试将我的数据框中的几个字段写入 JSON。我在数据框中的数据结构是
Key|col1|col2|col3|col4
key|a |b |c |d
Key|a1 |b1 |c1 |d1
现在我正在尝试将 col1 到 col4 字段转换为 JSON 并为 Json 字段命名
预期输出
[Key,{cols:[{col1:a,col2:b,col3:c,col4:d},{col1:a1,col2:b1,col3:c1,col4:d1}]
我为此写了一个udf。
val summary = udf(
(col1:String, col2:String, col3:String, col4:String) => "{\"cols\":[" + " {\"col1\":" + col1 + ",\"col2\":" + col2 + ",\"col3\":" + col3 + ",\"col4\":" + col4 + "}]}"
)
val result = input.withColumn("Summary",summary('col1,'col2,'col3,'col4))
val result1 = result.select('Key,'Summary)
result1.show(10)
这是我的结果
[Key,{cols:[{col1:a,col2:b,col3:c,col4:d}]}]
[Key,{cols:[{col1:a1,col2:b1,col3:c1,col4:d1}]}]
如您所见,它们没有分组。有没有办法使用 UDF 本身对这些行进行分组。我是 scala/Spark 的新手,无法找出正确的 udf。
【问题讨论】:
-
我认为您没有正确终止您的“预期输出”;我希望最后会有另一个“}]”来匹配开头的“[{”。
标签: json scala apache-spark apache-spark-sql