【问题标题】:Spark Aggregating multiple columns (possible to array) from join outputSpark从连接输出聚合多个列(可能是数组)
【发布时间】:2019-06-02 02:28:21
【问题描述】:

我有以下数据集

表1

表2

现在我想得到下面的数据集。我试过左外连接 Table1.id == Table2.departmentid 但是,我没有得到想要的输出。

稍后,我需要使用此表来获取多个计数并将数据转换为 xml 。我将使用地图进行此转换。

任何帮助将不胜感激。

【问题讨论】:

标签: apache-spark apache-spark-sql


【解决方案1】:

仅加入并不足以获得所需的输出。可能你遗漏了一些东西,每个嵌套数组的最后一个元素可能是departmentid。假设嵌套数组的最后一个元素是departmentid,我通过以下方式生成了输出:

import org.apache.spark.sql.{Row, SparkSession}
import org.apache.spark.sql.functions.collect_list

case class department(id: Integer, deptname: String)
case class employee(employeid:Integer, empname:String, departmentid:Integer)

val spark = SparkSession.builder().getOrCreate()
import spark.implicits._
val department_df = Seq(department(1, "physics")
                            ,department(2, "computer") ).toDF()
val emplyoee_df = Seq(employee(1, "A", 1)
                      ,employee(2, "B", 1)
                      ,employee(3, "C", 2)
                      ,employee(4, "D", 2)).toDF()

val result = department_df.join(emplyoee_df, department_df("id") === emplyoee_df("departmentid"), "left").
      selectExpr("id", "deptname", "employeid", "empname").
      rdd.map {
        case Row(id:Integer, deptname:String, employeid:Integer, empname:String) => (id, deptname, Array(employeid.toString, empname, id.toString))
      }.toDF("id", "deptname", "arrayemp").
          groupBy("id", "deptname").
          agg(collect_list("arrayemp").as("emplist")).
        orderBy("id", "deptname")

输出如下:

result.show(false)
+---+--------+----------------------+
|id |deptname|emplist               |
+---+--------+----------------------+
|1  |physics |[[2, B, 1], [1, A, 1]]|
|2  |computer|[[4, D, 2], [3, C, 2]]|
+---+--------+----------------------+

说明:如果我将最后一个数据帧转换分解为多个步骤,它可能会清楚输出是如何生成的。

department_df 和employee_df 之间的左外连接

val df1 = department_df.join(emplyoee_df, department_df("id") === emplyoee_df("departmentid"), "left").
      selectExpr("id", "deptname", "employeid", "empname")
df1.show()
    +---+--------+---------+-------+
| id|deptname|employeid|empname|
+---+--------+---------+-------+
|  1| physics|        2|      B|
|  1| physics|        1|      A|
|  2|computer|        4|      D|
|  2|computer|        3|      C|
+---+--------+---------+-------+

使用 df1 数据框中的某些列的值创建数组

val df2 = df1.rdd.map {
                case Row(id:Integer, deptname:String, employeid:Integer, empname:String) => (id, deptname, Array(employeid.toString, empname, id.toString))
              }.toDF("id", "deptname", "arrayemp")
df2.show()
            +---+--------+---------+
        | id|deptname| arrayemp|
        +---+--------+---------+
        |  1| physics|[2, B, 1]|
        |  1| physics|[1, A, 1]|
        |  2|computer|[4, D, 2]|
        |  2|computer|[3, C, 2]|
        +---+--------+---------+

使用 df2 数据框创建聚合多个数组的新列表

val result = df2.groupBy("id", "deptname").
              agg(collect_list("arrayemp").as("emplist")).
              orderBy("id", "deptname")
result.show(false)
            +---+--------+----------------------+
        |id |deptname|emplist               |
        +---+--------+----------------------+
        |1  |physics |[[2, B, 1], [1, A, 1]]|
        |2  |computer|[[4, D, 2], [3, C, 2]]|
        +---+--------+----------------------+

【讨论】:

    【解决方案2】:
    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.types._
    import org.apache.spark.sql.Row
    
    val df = spark.sparkContext.parallelize(Seq(
       (1,"Physics"),
       (2,"Computer"),
       (3,"Maths")
     )).toDF("ID","Dept")
    
     val schema = List(
        StructField("EMPID", IntegerType, true),
        StructField("EMPNAME", StringType, true),
        StructField("DeptID", IntegerType, true)
      )
    
      val data = Seq(
        Row(1,"A",1),
        Row(2,"B",1),
        Row(3,"C",2),
        Row(4,"D",2) ,
        Row(5,"E",null)
      )
    
      val df_emp = spark.createDataFrame(
        spark.sparkContext.parallelize(data),
        StructType(schema)
      )
    
      val newdf =  df_emp.withColumn("CONC",array($"EMPID",$"EMPNAME",$"DeptID")).groupBy($"DeptID").agg(expr("collect_list(CONC) as emplist"))
    
      df.join(newdf,df.col("ID") === df_emp.col("DeptID")).select($"ID",$"Dept",$"emplist").show()
    
    ---+--------+--------------------+
    | ID|    Dept|             listcol|
    +---+--------+--------------------+
    |  1| Physics|[[1, A, 1], [2, B...|
    |  2|Computer|[[3, C, 2], [4, D...|
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2022-01-24
      • 2022-11-30
      • 1970-01-01
      • 1970-01-01
      • 2013-10-31
      • 1970-01-01
      • 2014-10-15
      相关资源
      最近更新 更多