You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在PySpark中连接两个DataFrame且不产生重复行?

PySpark Outer Join 合并同Key行解决方案

问题背景

你基于column1、column2、column3对两个DataFrame执行outer join后,同key的记录被拆分为两行,但期望将column4和column5的非空值合并到同一行。

实现方法

方法1:分组聚合(通用方案)

先执行全外连接,再以关联key列为分组依据,聚合提取各列的非空值:

from pyspark.sql import functions as F

# 执行全外连接
joined_df = df1.join(df2, on=["column1", "column2", "column3"], how="outer")

# 分组聚合,合并同key的非空字段
result_df = joined_df.groupBy("column1", "column2", "column3") \
    .agg(
        F.first("column4", ignorenulls=True).alias("column4"),
        F.first("column5", ignorenulls=True).alias("column5")
    )

result_df.show()

方法2:直接关联+去重(适用于key无重复场景)

若确认column1/column2/column3在两个DataFrame中唯一,可简化为:

from pyspark.sql import functions as F

result_df = df1.join(df2, on=["column1", "column2", "column3"], how="outer") \
    .select("column1", "column2", "column3", "column4", "column5") \
    .groupBy("column1", "column2", "column3") \
    .agg(
        F.max("column4").alias("column4"),
        F.max("column5").alias("column5")
    )

result_df.show()

排查提示

如果outer join后出现重复行,可能是column1/column2/column3的值存在隐形差异(如空格、大小写),可通过以下代码验证:

# 检查两个DataFrame的key列是否存在差异
df1.select("column1", "column2", "column3").exceptAll(df2.select("column1", "column2", "column3")).show()

若输出为空,说明key值完全一致,使用上述方法即可得到期望结果。


内容的提问来源于stack exchange,提问作者PRuss

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 00:12:36