【问题标题】:How to unpack multiple keys in a Spark DataSet如何解压 Spark 数据集中的多个键
【发布时间】: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


    【解决方案1】:

    别名是一种无类型的操作,因此您必须在之后重新键入它。解压密钥的唯一方法是在之后通过选择或其他方式进行:

    df.groupByKey(p => (p.gender, p.age))
      .agg(typed.avg[Person](_.salary).as("average_salary").as[Double])
      .select($"key._1",$"key._2",$"average_salary").show
    

    【讨论】:

      【解决方案2】:

      实现这两个目标的最简单方法是将map()从聚合结果再次传递到Person实例:

      .map{case ((gender, age), salary) => Person(gender, age, salary)}
      

      如果在案例类的构造函数中稍微重新排列参数的顺序,结果会看起来最好:

      case class Person(gender: String, age: Int, salary: Double)
      
      +------+---+--------+
      |gender|age|  salary|
      +------+---+--------+
      |female| 31| 89000.0|
      |female| 51|185000.0|
      |  male| 27|110000.0|
      |  male| 26| 95000.0|
      +------+---+--------+
      

      完整代码:

      import session.implicits._
      val df = session.sparkContext.parallelize(List(
        Person("male", 27, 100000),
        Person("male", 27, 120000),
        Person("male", 26, 95000),
        Person("female", 31, 89000),
        Person("female", 51, 250000),
        Person("female", 51, 120000)
      )).toDS
      
      import org.apache.spark.sql.expressions.scalalang.typed
      df.groupByKey(p => (p.gender, p.age))
        .agg(typed.avg(_.salary))
        .map{case ((gender, age), salary) => Person(gender, age, salary)}
        .show()
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2023-04-02
        • 2017-03-24
        • 2016-05-06
        • 1970-01-01
        • 1970-01-01
        • 2023-03-25
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多