【发布时间】: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 |
+-------+----------+--------+--------+
我想在这两个数据帧之间执行连接,我希望有四个数据帧作为输出
存在于 leftDF 而不在 rightDF 中的事务
存在于 rightDF 中而不存在于 leftDF 中的事务
-
在两个数据帧之间具有至少一个公共 ID 的事务
3.a 严格匹配:相同货币,两个数据框之间的金额。示例:id 为 1 的交易。
3.b 宽松匹配:具有相同 ID 但货币/金额组合不同的交易。 id 为 4 的示例交易。
这是所需的输出:
-
leftDF 中存在的事务而不是 rightDF 中的事务
+------+---------+-------+-------+-------+----------+--------+--------+ |leftId|leftAltId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt| +------+---------+-------+-------+-------+----------+--------+--------+ |2 |200 |INR |100 |null |null |null |null | +------+---------+-------+-------+-------+----------+--------+--------+ -
存在于 rightDF 中而不存在于 leftDF 中的事务
+------+---------+-------+-------+-------+----------+--------+--------+ |leftId|leftAtId|leftCur|leftAmt|rightId|rightAltId|rightCur|rightAmt| +------+---------+-------+-------+-------+----------+--------+--------+ |null |null |null |null |3 |400 |MXN |100 | +------+---------+-------+-------+-------+----------+--------+--------+ -
在两个数据帧之间具有至少一个公共 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