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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 08:33:24