【问题标题】:Scala Transformation and actionScala 转换和操作
【发布时间】:2019-03-19 21:07:16
【问题描述】:

我有一个 RDD List[(String, List[Int])] 像 List(("A",List(1,2,3,4)),("B",List(5,6,7 )))

如何将它们转换为 List(("A",1),("A",2),("A",3),("A",4),("B",5), ("B",6),("B",7))

然后操作将通过 key 减少并生成 List(("A",2.5)("B",6)) 之类的结果

我尝试过使用 map(e=>List(e._1,e._2)) 但它没有给出想要的结果。

“A”平均为 2.5,“B”平均为 6

帮助我完成这些转变和行动。 提前致谢

【问题讨论】:

  • 后面的部分我已经弄清楚它如下所示:- val rdd_toreduce = spark.sparkContext.parallelize(List(("A",1.0),("A",2.0),("A ",3.0),("A",4.0),("B",5.0),("B",6.0),("B",7.0))) .mapValues(value => (value, 1)) .reduceByKey { case ((sumL, countL), (sumR, countR)) => (sumL + sumR, countL + countR) } .mapValues { case (sum , count) => sum / count } .collect rdd_toreduce.foreach( println)

标签: scala apache-spark rdd


【解决方案1】:

有几种方法可以得到你想要的。你也可以使用for comprehension,但我首先想到的是这个实现:

val l = List(("A", List(1, 2, 3)), ("B", List(1, 2, 3)))

val flattenList = l.flatMap {
  case (elem, _elemList) =>
    _elemList.map((elem, _))
}

输出:

List((A,1), (A,2), (A,3), (B,1), (B,2), (B,3))

【讨论】:

  • 感谢您的解决方案。这正是我想要的。
【解决方案2】:

如果您最终想要的是每个列表的平均值,那么没有必要使用flatMap 将它们分解为单独的元素。对大型列表这样做会不必要地用大型数据集打乱大量数据。

由于它们已经通过键聚合,只需将它们转换为如下所示:

val l = spark.sparkContext.parallelize(Seq(
  ("A", List(1, 2, 3, 4)),
  ("B", List(5, 6, 7))
))

val avg = l.map(r => {
    (r._1, (r._2.sum.toDouble / r._2.length.toDouble))
})

avg.collect.foreach(println)

请记住,如果您的任何列表长度为 0,这将失败。如果您有一些0 长度列表,则必须在映射中放置一个检查条件。

上面的代码给你:

(A,2.5)
(B,6.0)

【讨论】:

  • 谢谢特拉维斯,它也很有效。实际上,今天我面临一些采访,他们在两个单独的部分中问过我这个问题。一个是平面图,另一个是聚合,但我无法实现它的第一部分和我能够做的下一部分,这就是我这样问的原因。
【解决方案3】:

你可以试试explode()

scala> val df = List(("A",List(1,2,3,4)),("B",List(5,6,7))).toDF("x","y")
df: org.apache.spark.sql.DataFrame = [x: string, y: array<int>]

scala> df.withColumn("z",explode('y)).show(false)
+---+------------+---+
|x  |y           |z  |
+---+------------+---+
|A  |[1, 2, 3, 4]|1  |
|A  |[1, 2, 3, 4]|2  |
|A  |[1, 2, 3, 4]|3  |
|A  |[1, 2, 3, 4]|4  |
|B  |[5, 6, 7]   |5  |
|B  |[5, 6, 7]   |6  |
|B  |[5, 6, 7]   |7  |
+---+------------+---+


scala> val df2 = df.withColumn("z",explode('y))
df2: org.apache.spark.sql.DataFrame = [x: string, y: array<int> ... 1 more field]

scala> df2.groupBy("x").agg(sum('z)/count('z) ).show(false)
+---+-------------------+
|x  |(sum(z) / count(z))|
+---+-------------------+
|B  |6.0                |
|A  |2.5                |
+---+-------------------+


scala>

【讨论】:

  • 谢谢哥们。这也是更好的方法..!!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-08-21
  • 2021-12-01
  • 2012-08-09
  • 2015-07-11
  • 2013-07-01
相关资源
最近更新 更多