【问题标题】:generic join between two dataframes spark/scala两个数据框 spark/scala 之间的通用连接
【发布时间】:2019-01-27 03:05:06
【问题描述】:

我有两个数据框,即左和右。我有我的问题的工作解决方案。我需要一种使它通用的方法。我的问题在这里结束。

左DF:

+------+---------+-------+-------+
|leftId|leftAltId|leftCur|leftAmt|
+------+---------+-------+-------+
|1     |100      |USD    |20     |
|2     |200      |INR    |100    |
|4     |500      |MXN    |100    |
+------+---------+-------+-------+

右DF:

+-------+----------+--------+--------+
|rightId|rightAltId|rightCur|rightAmt|
+-------+----------+--------+--------+
|1      |300       |USD     |20      |
|3      |400       |MXN     |100     |
|4      |600       |MXN     |200     |
+-------+----------+--------+--------+

我想在这两个数据帧之间执行连接,我希望有四个数据帧作为输出

  1. 存在于 leftDF 而不在 rightDF 中的事务

  2. 存在于 rightDF 中而不存在于 leftDF 中的事务

  3. 在两个数据帧之间具有至少一个公共 ID 的事务

    3.a 严格匹配:相同货币,两个数据框之间的金额。示例:id 为 1 的交易。

    3.b 宽松匹配:具有相同 ID 但货币/金额组合不同的交易。 id 为 4 的示例交易。

这是所需的输出:

  1. leftDF 中存在的事务而不是 rightDF 中的事务

    +------+---------+-------+-------+-------+----------+--------+--------+
    |leftId|leftAltId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt|
    +------+---------+-------+-------+-------+----------+--------+--------+
    |2     |200      |INR    |100    |null   |null      |null    |null    |
    +------+---------+-------+-------+-------+----------+--------+--------+
    
  2. 存在于 rightDF 中而不存在于 leftDF 中的事务

    +------+---------+-------+-------+-------+----------+--------+--------+
    |leftId|leftAtId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt|
    +------+---------+-------+-------+-------+----------+--------+--------+
    |null  |null     |null   |null   |3      |400       |MXN     |100     |
    +------+---------+-------+-------+-------+----------+--------+--------+
    
  3. 在两个数据帧之间具有至少一个公共 ID 的事务

    +------+---------+-------+-------+-------+----------+--------+--------+
    |leftId|leftAltId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt|
    +------+---------+-------+-------+-------+----------+--------+--------+
    |1     |100      |USD    |20     |1      |300       |USD     |20      |
    |4     |500      |MXN    |100    |4      |600       |MXN     |200     |
    +------+---------+-------+-------+-------+----------+--------+--------+
    

    3.a 严格匹配:相同货币,两个数据框之间的金额。示例:id 为 1 的交易。

    +------+---------+-------+-------+-------+----------+--------+--------+        
    |leftId|leftAltId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt|
    +------+---------+-------+-------+-------+----------+--------+--------+
    |1     |100      |USD    |20     |1      |300       |USD     |20      |
    +------+---------+-------+-------+-------+----------+--------+--------+
    

    3.b 宽松匹配:具有相同 ID 但货币/金额组合不同的交易。 id 为 4 的示例交易。

     +------+---------+-------+-------+-------+----------+--------+--------+
    |leftId|leftAltId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt|
    +------+---------+-------+-------+-------+----------+--------+--------+
    |4     |500      |MXN    |100    |4      |600       |MXN     |200     |
    +------+---------+-------+-------+-------+----------+--------+--------+
    

这是我的工作代码:

import sparkSession.implicits._

val leftDF: DataFrame = Seq((1, 100, "USD", 20), (2, 200, "INR", 100), (4, 500, "MXN", 100)).toDF("leftId", "leftAltId", "leftCur", "leftAmt")
val rightDF: DataFrame = Seq((1, 300, "USD", 20), (3, 400, "MXN", 100), (4, 600, "MXN", 200)).toDF("rightId", "rightAltId", "rightCur", "rightAmt")

leftDF.show(false)
rightDF.show(false)
val idMatchQuery = leftDF("leftId") === rightDF("rightId") || leftDF("leftAltId") === rightDF("rightAltId")
val currencyMatchQuery = leftDF("leftCur") === rightDF("rightCur") && leftDF("leftAmt") === rightDF("rightAmt")
val leftOnlyQuery = (col("leftId").isNotNull && col("rightId").isNull) || (col("leftAltId").isNotNull && col("rightAltId").isNull)
val rightOnlyQuery = (col("rightId").isNotNull && col("leftId").isNull) || (col("rightAltId").isNotNull && col("leftAltId").isNull)
val matchQuery = (col("rightId").isNotNull && col("leftId").isNotNull) || (col("rightAltId").isNotNull && col("leftAltId").isNotNull)

val result = leftDF.join(rightDF, idMatchQuery, "fullouter")

val leftOnlyDF = result.filter(leftOnlyQuery)
val rightOnlyDF = result.filter(rightOnlyQuery)

val matchDF = result.filter(matchQuery)
val strictMatchDF = matchDF.filter(currencyMatchQuery.equalTo(true))
val relaxedMatchDF = matchDF.filter(currencyMatchQuery.equalTo(false))

leftOnlyDF.show(false)
rightOnlyDF.show(false)
matchDF.show(false)
strictMatchDF.show(false)
relaxedMatchDF.show(false)

我在寻找什么:

我希望能够将列名作为列表加入,并使代码通用。

例如

    val relaxedJoinList = Array(("leftId", "rightId"), ("leftAltId", "rightAltId"))
    val strictJoinList = Array(("leftCur", "rightCur"), ("leftAmt", "rightAmt"))

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    我希望能够将列名作为列表加入,并使代码通用。

    这不是一个完美的建议,但肯定会帮助您概括。建议使用foldLeft

    val relaxedJoinList = Array(("leftId", "rightId"), ("leftAltId", "rightAltId"))
    val rHead = relaxedJoinList.head
    
    val strictJoinList = Array(("leftCur", "rightCur"), ("leftAmt", "rightAmt"))
    val sHead = strictJoinList.head
    
    val idMatchQuery = relaxedJoinList.tail.foldLeft(leftDF(rHead._1) === rightDF(rHead._2)){(x, y) => x || leftDF(y._1) === rightDF(y._2)}
    val currencyMatchQuery = strictJoinList.tail.foldLeft(leftDF(sHead._1) === rightDF(sHead._2)){(x, y) => x && leftDF(y._1) === rightDF(y._2)}
    val leftOnlyQuery = relaxedJoinList.tail.foldLeft(col(rHead._1).isNotNull && col(rHead._2).isNull){(x, y) => x || col(y._1).isNotNull && col(y._2).isNull}
    val rightOnlyQuery = relaxedJoinList.tail.foldLeft(col(rHead._1).isNull && col(rHead._2).isNotNull){(x, y) => x || col(y._1).isNull && col(y._2).isNotNull}
    val matchQuery = relaxedJoinList.tail.foldLeft(col(rHead._1).isNotNull && col(rHead._2).isNotNull){(x, y) => x || col(y._1).isNotNull && col(y._2).isNotNull}
    

    其余的代码是你的

    希望回答对你有帮助

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-01-04
      • 2017-08-16
      • 1970-01-01
      • 2020-05-22
      • 2019-06-06
      • 1970-01-01
      相关资源
      最近更新 更多