如何使用PySpark求取两个DataFrame同ID数据的交集与对称差集
PySpark按ID匹配实现DataFrame交集与对称差集
以下示例默认两个待处理DataFrame名为df1、df2,关联匹配的主键列名为id,先给出测试用样例数据方便复现:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, coalesce spark = SparkSession.builder.appName("df_set_operation").getOrCreate() # 构造测试数据 df1 = spark.createDataFrame( [(1, "a"), (2, "b"), (3, "c")], schema=["id", "val"] ) df2 = spark.createDataFrame( [(2, "b"), (3, "d"), (4, "e")], schema=["id", "val"] )
按ID取交集
交集指两个DataFrame中均存在对应ID的匹配记录,两种常用实现:
- 内连接实现(最常用,适合需要同时取两边字段的场景)
# 直接按id做内连接,自动保留两边id匹配的记录 intersect_df = df1.join(df2, on="id", how="inner") intersect_df.show()
执行后返回id为2、3的两条匹配记录。
- ID交集过滤实现(适合仅需要保留单表字段的场景)
# 先取出两表共有的ID集合,再过滤单表即可 common_id_set = df1.select("id").intersect(df2.select("id")) intersect_df_simple = df1.join(common_id_set, on="id")
按ID取对称差集
对称差集指仅在其中一个DataFrame存在、另一表无对应ID的所有记录,即「df1独有+df2独有」的合集,两种常用实现:
- 全外连接实现(适合需要同时对比两边同ID字段差异的场景)
full_join_res = df1.alias("t1").join( df2.alias("t2"), on=col("t1.id") == col("t2.id"), how="full_outer" ) # 过滤出任意一边id为空的记录,即为单边独有数据 symmetric_diff_df = full_join_res.filter( col("t1.id").isNull() | col("t2.id").isNull() ).select( # 合并两边非空的id为统一列 coalesce(col("t1.id"), col("t2.id")).alias("id"), col("t1.val").alias("val_from_df1"), col("t2.val").alias("val_from_df2") ) symmetric_diff_df.show()
执行后返回id为1(仅df1存在)、id为4(仅df2存在)的两条记录。
- 反连接合并实现(执行效率更高,适合快速取单边独有数据的场景)
# left_anti是PySpark内置的差集join类型,性能优于left join后过滤空值 only_in_df1 = df1.join(df2.select("id"), on="id", how="left_anti") only_in_df2 = df2.join(df1.select("id"), on="id", how="left_anti") # 对齐字段后合并两边独有数据即为对称差集 symmetric_diff_df_fast = only_in_df1.unionByName(only_in_df2, allowMissingColumns=True)
补充说明:如果不需要按ID匹配,而是要取整行内容完全一致的交集/差集,可直接调用DataFrame内置API:
df1.intersect(df2)取整行交集、df1.exceptAll(df2).unionByName(df2.exceptAll(df1))取整行对称差集,无需手动写join逻辑。
内容的提问来源于stack exchange,提问作者sun tony
相关产品推荐
相关产品推荐

