Spark全外连接DataFrame合并name列优先取df2值无匹配取df1值方法
实现方式
核心用Spark内置的coalesce函数实现取值优先级逻辑,该函数会按传入参数的顺序返回第一个非null值,刚好匹配「优先取df2的name值,df2无匹配时取df1对应name值」的规则,提供两种常用实现:
写法1:最小改动适配原有逻辑
在原有全外连接代码的基础上,先给两个DataFrame设置别名避免重名列引用冲突,join完成后用coalesce合并两个name列即可:
from delta.tables import * from pyspark.sql import functions as F vals1 = [(1, "a"), (2, "b"), (3, "c"), (4, "d") ] columns1 = ["id","name"] df1 = spark.createDataFrame(data=vals1, schema=columns1) vals2 = [(1, "k"), (2, "l"), (3, "m") ] columns2 = ["id","name"] df2 = spark.createDataFrame(data=vals2, schema=columns2) # 给两个df设置别名,避免重名的name列引用冲突 df_joined = df1.alias("df1").join(df2.alias("df2"), on="id", how="full") # 按优先级合并name列:优先取df2的name,为null时取df1的name df_result = df_joined.select( "id", F.coalesce(F.col("df2.name"), F.col("df1.name")).alias("name") # 如果需要固定按id升序得到k、l、m、d的输出,加下面这行排序 # ).orderBy("id") ) df_result.show()
运行输出:
+---+----+ | id|name| +---+----+ | 1| k| | 3| m| | 2| l| | 4| d| +---+----+
写法2:拆分join逻辑(适合大数据量场景)
全外连接的shuffle成本较高,如果你的场景只需要保留df1中存在的id(和示例预期输出完全匹配,不需要保留仅在df2存在的id),可以拆分join步骤减少性能开销:
from pyspark.sql import functions as F # 拿到df1和df2能匹配上的记录,直接取df2的name值 df_matched = df1.alias("df1").join( df2.alias("df2"), on=F.col("df1.id") == F.col("df2.id"), how="inner" ).select(F.col("df1.id"), F.col("df2.name").alias("name")) # 拿到df1中没有匹配到df2的记录,取df1自身的name值 df_unmatched = df1.join(df2, on="id", how="left_anti").select("id", "name") # 合并两部分结果,需要固定顺序就加orderBy("id") df_result = df_matched.unionByName(df_unmatched) df_result.show()
运行后输出和写法1完全一致。
说明:Spark的shuffle类操作(比如join、union)默认不保证行输出顺序,如果需要严格得到
k、l、m、d的顺序,一定要加orderBy("id")做显式排序。
内容的提问来源于stack exchange,提问作者user626528
相关产品推荐
相关产品推荐

