【问题标题】:aggregating jsonarray into Map<key, list> in spark in spark2.x在 spark2.x 中将 jsonarray 聚合到 Spark 中的 Map<key, list>
【发布时间】:2018-06-18 18:36:24
【问题描述】:

我对 Spark 很陌生。我有一个输入 json 文件,我正在阅读它

val df = spark.read.json("/Users/user/Desktop/resource.json");

resource.json 的内容如下:

{"path":"path1","key":"key1","region":"region1"}
{"path":"path112","key":"key1","region":"region1"}
{"path":"path22","key":"key2","region":"region1"}

有什么方法可以处理这个数据框并将结果聚合为

Map<key, List<data>>

其中 data 是每个存在 key 的 json 对象。

例如:预期结果是

Map<key1 =[{"path":"path1","key":"key1","region":"region1"}, {"path":"path112","key":"key1","region":"region1"}] ,
key2 = [{"path":"path22","key":"key2","region":"region1"}]>

任何进一步进行的参考/文档/链接都会有很大帮助。

谢谢。

【问题讨论】:

  • 你可以在scala中使用groupBy。

标签: scala apache-spark bigdata


【解决方案1】:

你可以这样做:

import org.json4s._
import org.json4s.jackson.Serialization.read
case class cC(path: String, key: String, region: String)

val df = spark.read.json("/Users/user/Desktop/resource.json");

scala> df.show
+----+-------+-------+
| key|   path| region|
+----+-------+-------+
|key1|  path1|region1|
|key1|path112|region1|
|key2| path22|region1|
+----+-------+-------+
//Please note that original json structure is gone. Use .toJSON to get json back and extract key from json and create RDD[(String, String)] RDD[(key, json)]

val rdd = df.toJSON.rdd.map(m => {
implicit val formats = DefaultFormats
val parsedObj = read[cC](m)
(parsedObj.key, m)
})

scala> rdd.collect.groupBy(_._1).map(m => (m._1,m._2.map(_._2).toList))
res39: scala.collection.immutable.Map[String,List[String]] = Map(key2 -> List({"key":"key2","path":"path22","region":"region1"}), key1 -> List({"key":"key1","path":"path1","region":"region1"}, {"key":"key1","path":"path112","region":"region1"}))

【讨论】:

  • 这正是我一直在寻找的。谢谢@hadooper。
  • 如果某个答案已经解决了您的问题,请接受它 - 请参阅当有人回答我的问题时我该怎么办? stackoverflow.com/help/someone-answers
【解决方案2】:

您可以将groupBycollect_list 结合使用,这是一个聚合函数,可将所有匹配值收集到每个键的列表中。

请注意,原始 JSON 字符串已经“消失”(Spark 将它们解析为单独的列),所以如果您真的想要所有 记录 的列表(包括它们的所有列,包括键) ,您可以使用struct 函数将列合并为一列:

import org.apache.spark.sql.functions._
import spark.implicits._

df.groupBy($"key")
  .agg(collect_list(struct($"path", $"key", $"region")) as "value")

结果是:

+----+--------------------------------------------------+
|key |value                                             |
+----+--------------------------------------------------+
|key1|[[path1, key1, region1], [path112, key1, region1]]|
|key2|[[path22, key2, region1]]                         |
+----+--------------------------------------------------+

【讨论】:

    猜你喜欢
    • 2017-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-01-15
    • 2016-07-14
    • 1970-01-01
    • 2020-10-04
    • 2017-12-10
    相关资源
    最近更新 更多