【发布时间】:2017-08-11 02:52:53
【问题描述】:
我有以下DataSet,结构如下。
case class Person(age: Int, gender: String, salary: Double)
我想通过gender 和age 确定平均工资,因此我将DS 用两个键分组。我遇到了两个主要问题,一个是两个键混合在一个中,但我想将它们保留在两个不同的列中,另一个是 aggregated 列的名称很长而且我不能弄清楚如何重命名它(显然as 和alias 不起作用)所有这些都使用DS API。
val df = sc.parallelize(List(Person(100000.00, "male", 27),
Person(120000.00, "male", 27),
Person(95000, "male", 26),
Person(89000, "female", 31),
Person(250000, "female", 51),
Person(120000, "female", 51)
)).toDF.as[Person]
df.groupByKey(p => (p.gender, p.age)).agg(typed.avg(_.salary)).show()
+-----------+------------------------------------------------------------------------------------------------+
| key| TypedAverage(line2503618a50834b67a4b132d1b8d2310b12.$read$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$Person)|
+-----------+------------------------------------------------------------------------------------------------+
|[female,31]| 89000.0...
|[female,51]| 185000.0...
| [male,27]| 110000.0...
| [male,26]| 95000.0...
+-----------+------------------------------------------------------------------------------------------------+
【问题讨论】:
标签: scala apache-spark apache-spark-dataset