【问题标题】:Spark Dataframes join with 2 columns using or operatorSpark Dataframes 使用 or 运算符加入 2 列
【发布时间】:2018-08-18 04:58:44
【问题描述】:

数据框 1:

+---------+---------+
|login_Id1|login_Id2|
+---------+---------+
|  1234567|  1234568|
|  1234567|     null|
|     null|  1234568|
|  1234567|  1000000|
|  1000000|  1234568|
|  1000000|  1000000|
+---------+---------+

数据帧 2:

+--------+---------+-----------+
|login_Id|user_name| user_Email|
+--------+---------+-----------+
| 1234567|TestUser1|user1_Email|
| 1234568|TestUser2|user2_Email|
| 1234569|TestUser3|user3_Email|
| 1234570|TestUser4|user4_Email|
+--------+---------+-----------+

预期输出

+---------+---------+--------+---------+-----------+
|login_Id1|login_Id2|login_Id|user_name| user_Email|
+---------+---------+--------+---------+-----------+
|  1234567|  1234568| 1234567|TestUser1|user1_Email|
|  1234567|     null| 1234567|TestUser1|user1_Email|
|     null|  1234568| 1234568|TestUser2|user2_Email|
|  1234567|  1000000| 1234567|TestUser1|user1_Email|
|  1000000|  1234568| 1234568|TestUser2|user2_Email|
|  1000000|  1000000|    null|     null|       null|
+---------+---------+--------+---------+-----------+

我的要求是我必须加入两个数据框,以便从 DataFrame 2 获取每个登录 ID 的附加信息。login_Id1 或 login_Id2 都将有数据(在大多数情况下)。有时这两列也可能有数据。在这种情况下,我想使用 login_Id1 执行连接。当两列不匹配时,我希望 null 作为结果

我点击了这个链接

Join in spark dataframe (scala) based on not null values

我试过了

DataFrame1.join(broadcast(DataFrame2), DataFrame1("login_Id1") === DataFrame2("login_Id") || DataFrame1("login_Id2") === DataFrame2("login_Id") )

我得到的输出是

+---------+---------+--------+---------+-----------+
|login_Id1|login_Id2|login_Id|user_name| user_Email|
+---------+---------+--------+---------+-----------+
|  1234567|  1234568| 1234567|TestUser1|user1_Email|
|  1234567|  1234568| 1234568|TestUser2|user2_Email|
|  1234567|     null| 1234567|TestUser1|user1_Email|
|     null|  1234568| 1234568|TestUser2|user2_Email|
|  1234567|  1000000| 1234567|TestUser1|user1_Email|
|  1000000|  1234568| 1234568|TestUser2|user2_Email|
|  1000000|  1000000|    null|     null|       null|
+---------+---------+--------+---------+-----------+

当任一列都有值时,我得到了预期的行为。当它们都有值时,对两个列(Row1,Row3)执行连接。在这种情况下 ||不会短路吗?

有没有办法获得预期的数据帧?

到目前为止,我有一个 udf 函数,它检查 login_Id1 是否有值(返回 login_Id1)或 login_Id2 是否有值(返回 login_Id2),如果它们都有值,我将返回 loginId1,并添加 udf 函数的结果作为 DataFrame1 的另一列(Filtered_Login_id)。

使用 udf 添加 FilteredId 列后的 Dataframe1

+--------+---------+-----------+
|loginId1|loginId2 | FilteredId|
+--------+---------+-----------+
| 1234567|1234568  |1234567    |
| 1234567|null     |1234567    |
| null   |1234568  |1234568    |
| 1234567|1000000  |1234567    |
| 1000000|1234568  |1000000    |
| 1000000|1000000  |1000000    |
+--------+---------+-----------+

然后我根据 FilteredId ===loginId 执行 join 并得到结果

DataFrame1.join(broadcast(DataFrame2), DataFrame1("FilteredId") === DataFrame2("login_Id"),"left_outer" )

有没有更好的方法可以在不使用 udf 的情况下实现此结果?仅使用 join(其行为类似于短路或运算符)?

包括 Leo 指出的用例。我的 udf 方法错过了 Leo 指出的用例。我的确切要求是 2 个输入列值中的任何一个(login_Id1,login_Id2)是否与 Dataframe2 的 login_Id 匹配,即应获取 loginId 数据。如果任一列不匹配,则应添加 null(类似于左外连接)

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    我不清楚您的示例数据是否已经涵盖了login_Id-pairs 的所有场景。如果是这样,专注于null 检查的解决方案就足够了;否则它需要稍微复杂一些的东西(比如你使用UDF)。

    不依赖UDF 的一种方法是应用left_outer 加入df1left_semi 加入df2,每个附加一个flag 列用于优先排序,通过union 组合它们,加入df2 用于包含非键列,最后根据flag 消除重复行。

    这里是示例代码,示例数据更通用:

    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.expressions.Window
    
    val df1 = Seq(
      ("1234567", "1234568"),
      ("1234567", null),
      (null, "1234568"),
      ("1234569", "1000000"),
      ("1000000", "1234570"),
      ("1000000", "1000000")
    ).toDF("login_Id1", "login_Id2")
    
    val df2 = Seq(
      ("1234567", "TestUser1", "user1_Email"),
      ("1234568", "TestUser2", "user2_Email"),
      ("1234569", "TestUser3", "user3_Email"),
      ("1234570", "TestUser4", "user4_Email")
    ).toDF("login_Id", "user_name", "user_Email")
    
    val dfOuter = df1.join(df2, $"login_Id1" === df2("login_Id"), "left_outer").
      withColumn("flag", when($"login_Id".isNull, lit(9)).otherwise(lit(1))).
      select("login_Id1", "login_Id2", "flag")
    // +---------+---------+----+
    // |login_Id1|login_Id2|flag|
    // +---------+---------+----+
    // |  1234567|  1234568|   1|
    // |  1234567|     null|   1|
    // |     null|  1234568|   9|
    // |  1234569|  1000000|   1|
    // |  1000000|  1234570|   9|
    // |  1000000|  1000000|   9|
    // +---------+---------+----+
    
    val dfSemi = df1.join(df2, $"login_Id2" === df2("login_Id"), "left_semi").
      withColumn("flag", lit(2))
    // +---------+---------+----+
    // |login_Id1|login_Id2|flag|
    // +---------+---------+----+
    // |  1234567|  1234568|   2|
    // |     null|  1234568|   2|
    // |  1000000|  1234570|   2|
    // +---------+---------+----+
    
    val window = Window.partitionBy("login_Id1", "login_Id2").orderBy("flag")
    
    (dfOuter union dfSemi).
      withColumn("row_num", row_number.over(window)).
      where($"row_num" === 1).
      withColumn("login_Id", when($"flag" === 1, $"login_Id1").
        otherwise(when($"flag" === 2, $"login_Id2"))
      ).
      join(df2, Seq("login_Id"), "left_outer").
      select("login_Id1", "login_Id2", "login_Id", "user_name", "user_Email")
    // +---------+---------+--------+---------+-----------+
    // |login_Id1|login_Id2|login_Id|user_name| user_Email|
    // +---------+---------+--------+---------+-----------+
    // |  1000000|  1000000|    null|     null|       null|
    // |  1000000|  1234570| 1234570|TestUser4|user4_Email|
    // |  1234567|  1234568| 1234567|TestUser1|user1_Email|
    // |  1234569|  1000000| 1234569|TestUser3|user3_Email|
    // |  1234567|     null| 1234567|TestUser1|user1_Email|
    // |     null|  1234568| 1234568|TestUser2|user2_Email|
    // +---------+---------+--------+---------+-----------+
    

    请注意,如果broadcastdf1 相比明显更小,您可以像在现有示例代码中一样将df2 应用到df2。如果df2 小到足以成为collect-ed,则可以将其简化为以下内容:

    val loginIdList = df2.collect.map(r => r.getAs[String](0))
    
    val df1Unmatched = df1.where(
      !$"login_Id1".isin(loginIdList: _*) && !$"login_Id2".isin(loginIdList: _*)
    )
    
    (df1 except df1Unmatched).
      join( broadcast(df2), $"login_Id1" === $"login_Id" ||
        ($"login_Id2" === $"login_Id" &&
          ($"login_Id1".isNull || !$"login_Id1".isin(loginIdList: _*))
        )
      ).
      union(
        df1Unmatched.join(df2, $"login_Id2" === $"login_Id", "left_outer")
      )
    

    【讨论】:

    • 感谢 Leo 捕捉到这个用例,其中任一列的值在 Dataframe2 中不存在。这也是要求的一部分。错过了。将更新我的问题以包括这个案例.我会尝试你的解决方案,如果它对我有用,我会告诉你
    • 两种解决方案都有效,只是想知道使用 collect 的 dataframe2 应该有多小?在我的情况下,dataframe2 有 25000k 记录..
    • 很大程度上取决于硬件配置(尤其是内存)。我通常不会将 collect 应用于任何具有数百万行的 RDD/DataFrame。
    • 哦,谢谢,如何将您的方法转换为左外连接?我试过但无法修改它。我已经相应地编辑了我的问题。感谢您的帮助
    • @Aradhana,包含login_Id-pair(df2 中不包含id)的要求使事情变得更加复杂。请看我修改后的答案。
    【解决方案2】:

    如果第一列为空,您只需要第二列,将该条件添加到您的连接子句中:

    @ df1.join(df2, df1("login_Id1") <=> df2("login_Id") || (df1("login_Id1").isNull && df1("login_Id2") <=> df2("login_Id"))).show()
    +---------+---------+--------+---------+-----------+
    |login_Id1|login_Id2|login_Id|user_name| user_Email|
    +---------+---------+--------+---------+-----------+
    |  1234567|  1234568| 1234567|TestUser1|user1_Email|
    |  1234567|     null| 1234567|TestUser1|user1_Email|
    |     null|  1234568| 1234568|TestUser2|user2_Email|
    +---------+---------+--------+---------+-----------+
    

    注意:右侧仅找到此行:

    @ df1.join(df2, df1("login_Id1").isNull && df1("login_Id2") <=> df2("login_Id")).show()
    +---------+---------+--------+---------+-----------+
    |login_Id1|login_Id2|login_Id|user_name| user_Email|
    +---------+---------+--------+---------+-----------+
    |     null|  1234568| 1234568|TestUser2|user2_Email|
    +---------+---------+--------+---------+-----------+
    

    【讨论】:

    • 感谢您的回答,解决了 null 的情况,根据我最初的问题,我错过了包含 Leo 提到的用例 :)
    【解决方案3】:

    您可以使用coalesce 函数创建一个新值,即login_Id1(如果它不为空)或login_Id2(如果1 为空) - 并将该结果与login_Id 进行比较:

    import org.apache.spark.sql.functions._
    import spark.implicits._
    
    val res = DataFrame1.join(DataFrame2, coalesce($"login_Id1", $"login_Id2") === $"login_Id")
    
    res.show()
    +---------+---------+--------+---------+-----------+
    |login_Id1|login_Id2|login_Id|user_name| user_Email|
    +---------+---------+--------+---------+-----------+
    |  1234567|     null| 1234567|TestUser1|user1_Email|
    |  1234567|  1234568| 1234567|TestUser1|user1_Email|
    |     null|  1234568| 1234568|TestUser2|user2_Email|
    +---------+---------+--------+---------+-----------+
    

    【讨论】:

    • 感谢您的回答,解决了 null 的情况,根据我最初的问题,我错过了包含 Leo 提到的用例 :)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-07-08
    • 2012-04-23
    • 2019-07-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-09-08
    相关资源
    最近更新 更多