【问题标题】:Column wise comparison between Spark Dataframe using Spark core使用 Spark 核心的 Spark Dataframe 之间的逐列比较
【发布时间】:2021-02-11 17:40:50
【问题描述】:

给定示例,但要查找 N 列数比较两个数据框之间的列数。

给定 5 行 3 列的示例,以 EmpID 作为主键。

如何在 Spark 核心中进行这种比较?

输入Df1:

|EMPID |Dept     |  Salary
--------------------------
|1     |HR       |   100
|2     |IT       |   200
|3     |Finance  |   250
|4     |Accounts |   200
|5     |IT       |   150

输入DF2:

|EMPID |Dept       |Salary
------------------------------
|1     |HR         | 100
|2     |IT         | 200
|3     |FIN        | 250
|4     |Accounts   | 150
|5     |IT         | 150

预期结果 DF:

|EMPID   |Dept      |Dept      |status      |Salary     |Salary   |status
--------------------------------------------------------------------
|1       |HR        |HR        | TRUE       | 100       | 100     | TRUE
|2       |IT        |IT        | TRUE       | 200       | 200     | TRUE
|3       |Finance   |FIN       | False      | 250       | 250     | TRUE
|4       |Accounts  |Accounts  | TRUE       | 200       | 150     | FALSE
|5       |IT        |IT        | TRUE       | 150       | 150     | TRUE

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    您可以使用 EMPID 进行连接并比较结果列:

    val result = df1.alias("df1").join(
        df2.alias("df2"), "EMPID"
    ).select(
        $"EMPID",
        $"df1.Dept", $"df2.Dept",
        ($"df1.Dept" === $"df2.Dept").as("status"),
        $"df1.Salary", $"df2.Salary",
        ($"df1.Salary" === $"df2.Salary").as("status")
    )
    
    result.show
    +-----+--------+--------+------+------+------+------+
    |EMPID|    Dept|    Dept|status|Salary|Salary|status|
    +-----+--------+--------+------+------+------+------+
    |    1|      HR|      HR|  true|   100|   100|  true|
    |    2|      IT|      IT|  true|   200|   200|  true|
    |    3| Finance|     FIN| false|   250|   250|  true|
    |    4|Accounts|Accounts|  true|   200|   150| false|
    |    5|      IT|      IT|  true|   150|   150|  true|
    +-----+--------+--------+------+------+------+------+
    

    请注意,您可能希望重命名列,因为将来无法查询重复的列名。

    【讨论】:

      【解决方案2】:

      您可以使用 join 然后遍历 df.columns 来选择所需的输出列:

      val df_final = df1.alias("df1")
        .join(df2.alias("df2"), "EMPID")
        .select(
            Seq(col("EMPID")) ++
            df1.columns.filter(_ != "EMPID")
              .flatMap(c =>
                Seq(
                  col(s"df1.$c").as(s"df1_$c"),
                  col(s"df2.$c").as(s"df2_$c"),
                  (col(s"df1.$c") === col(s"df2.$c")).as(s"status_$c")
                )
            ): _*
      )
      
      df_final.show
      
      //+-----+--------+--------+-----------+----------+----------+-------------+
      //|EMPID|df1_Dept|df2_Dept|status_Dept|df1_Salary|df2_Salary|status_Salary|
      //+-----+--------+--------+-----------+----------+----------+-------------+
      //|    1|      HR|      HR|       true|       100|       100|         true|
      //|    2|      IT|      IT|       true|       200|       200|         true|
      //|    3| Finance|     FIN|      false|       250|       250|         true|
      //|    4|Accounts|Accounts|       true|       200|       150|        false|
      //|    5|      IT|      IT|       true|       150|       150|         true|
      //+-----+--------+--------+-----------+----------+----------+-------------+
      

      【讨论】:

        【解决方案3】:

        您也可以通过以下方式执行此操作:

        //Source data
        val df = Seq((1,"HR",100),(2,"IT",200),(3,"Finance",250),(4,"Accounts",200),(5,"IT",150)).toDF("EMPID","Dept","Salary")
        val df1 = Seq((1,"HR",100),(2,"IT",200),(3,"Fin",250),(4,"Accounts",150),(5,"IT",150)).toDF("EMPID","Dept","Salary")
        
        //joins and other operations
        val finalDF = df.as("d").join(df1.as("d1"),Seq("EMPID"),"inner")
        .withColumn("DeptStatus",$"d.Dept" === $"d1.Dept")
        .withColumn("Salarystatus",$"d.Salary" === $"d1.Salary")
        .selectExpr("EMPID","d.Dept","d1.Dept","DeptStatus as 
        Status","d.Salary","d1.Salary","SalaryStatus as Status")
        display(finalDF)
        

        你可以看到如下输出:

        【讨论】:

          猜你喜欢
          • 2018-10-03
          • 1970-01-01
          • 2019-08-09
          • 1970-01-01
          • 1970-01-01
          • 2021-07-05
          • 1970-01-01
          • 1970-01-01
          • 2022-06-11
          相关资源
          最近更新 更多