【发布时间】:2018-10-13 15:27:12
【问题描述】:
您好,我有两张这样的桌子。
源表
orig1 orig2 orig3 xref1 xref2 xref3
1 1 1 2 2 2
1 1 1 3 3 3
23 23 23 12 12 12
目标表:
orig1 orig2 orig3 xref1 xref2 xref3 version
1 1 1 1 1 1 0
我需要如下输出
1) 我需要匹配(source(orig1 orig2 orig3) == target(orig1 orig2 orig3)),
如果它的 macthing 我们需要通过将版本增加 1 来从源表追加到目标表
如果不匹配,只需将版本附加为 '0'
预期输出为:
orig1 orig2 orig3 xref1 xref2 xref3 version
1 1 1 1 1 1 0
1 1 1 2 2 2 1
1 1 1 3 3 3 2
23 23 23 12 12 12 0
我尝试了数据框级别。但它没有按预期工作。任何帮助将不胜感激。
我尝试了以下方法。
val source = spark.sql("select xref1,xref2,xref3,orig1,orig2,orig3 from default.source")
val target = spark.sql("select xref1,xref2,xref3,orig1,orig2,orig3 from default.target")
val target10 = spark.sql("select xref1,xref2,xref3,orig1,orig2,orig3,version from default.target")
val diff=( source.select("xref1","xref2","xref3","orig1","orig2","orig3") == target.select("xref1","xref2","xref3","orig1","orig2","orig3"))
if ( diff == false ){
val diff1 = source.select("orig1","orig2","orig3").except(target.select("orig1","orig2","orig3"))
if ( diff1.count > 0 ) {
val ver = target10.groupBy("orig1","orig2","orig3").max("version")
val common = source.select("orig1","orig2","orig3").intersect(target.select("orig1","orig2","orig3"))
val result = common.join(ver, common("orig1") === ver("orig1") && common("orig2") === ver("orig2") && common("orig3") === ver("orig3"), "inner").select(ver("orig1"),ver("orig2"),ver("orig3"),(ver("max(version)") + 1
) as "version")
val result1 = result.join(source, result("orig1") === source("orig1") && result("orig2") === source("orig2") && result("orig3") === source("orig3"), "inner").select(source("orig1"),source("orig2"),source("orig3"),result("version"),source("xref1"),source("xref2"),source("xref3"))
val result2=source.select("orig1","orig2","orig3").except(target.select("orig1","orig2","orig3")).withColumn("version",lit(0))
val execpettarget=result2.select($"orig1".alias("DIV"),$"orig2".alias("SEC"),$"orig3".alias("UN"),$"version".alias("VER"))
val result23 = execpettarget.join(source, execpettarget("DIV") === source("orig1") && execpettarget("SEC") === source("orig2") && execpettarget("UN") === source("orig3"), "inner").select(source("orig1"),source("orig2"),source("orig3"),execpettarget("VER"),source("orig1"),source("orig2"), source("orig3"))
val final_result = result1.union(result23)
final_result.show()
}else{
println("else")
val ver1 = target10.groupBy("orig1","orig2","orig3").max("version")
val common1 = source.select("orig1","orig2","orig3").intersect(target.select("orig1","orig2","orig3"))
val result11 = common1.join(ver1, common1("orig1") === ver1("orig1") && common1("orig2") === ver1("orig2") && common1("orig3") === ver1("orig3"), "inner").select(ver1("orig1"),ver1("orig2"),ver1("orig3"),(ver1("max(version)") + 1) as "version")
val result3 = result11.join(source, result11("orig1") === source("orig1") && result11("orig2") === source("orig2") && result11("orig3") === source("orig3"), "inner").select(source("orig1"),source("orig2"),source("orig3"),result11("version"),source("xref1"),source("xref2"),source("xref3"))
result3.show()
}}
但在最终连接中,源有 2 个重复行。因此,当将源与目标连接时,我会得到多行。
【问题讨论】:
-
看起来像是两个数据框的联合(相应地添加或删除一列)然后
row_number() over partition by orig1, orig2, orig3 order by xref1, xref2, xref3 -1 as version -
我尝试过类似的方法。但问题是当我加入源表时。我有前 3 列的重复行。这是个问题。
-
如果你展示一个例子会很有用。
-
请添加您的尝试和您得到的输出。 “但它没有按预期工作”相当模糊。
-
我现在就发帖
标签: python python-3.x apache-spark dataframe pyspark