PySpark中两个非平衡DataFrame的子集列差异查询问题
解决PySpark中非平衡DataFrame子集列对比并返回完整行的问题
针对你提出的需求——对比两个非平衡DataFrame的子集列(name和age),找出df_a中与df_b对应name的age不匹配的行并返回完整数据,subtract确实无法满足,因为它仅支持整行内容的对比,且要求两个DataFrame的列结构完全一致。下面提供两种可行的解决方案:
方法一:关联(Join)+ 过滤
这是最直接的实现方式,通过name列关联两个DataFrame,再过滤出age不匹配的行,最后保留df_a的完整列:
from pyspark.sql import functions as F # 基于name关联两个DataFrame,筛选age不匹配的行 result_df = df_a.join(df_b, on="name", how="inner") \ .filter(df_a.age != df_b.age) \ .select(df_a["id"], df_a["name"], df_a["age"]) # 查看结果 result_df.show()
执行后输出结果:
+---+----+---+ | id|name|age| +---+----+---+ | 2| b| 5| | 3| c| 7| +---+----+---+
如果只需要返回id列,只需修改select语句为:
.select(df_a["id"])
方法二:Spark SQL子查询
如果你更熟悉SQL语法,可以通过注册临时视图,用EXISTS子查询实现逻辑:
# 注册临时视图 df_a.createOrReplaceTempView("df_a") df_b.createOrReplaceTempView("df_b") # 执行SQL查询 result_df = sql_context.sql(""" SELECT id, name, age FROM df_a WHERE EXISTS ( SELECT 1 FROM df_b WHERE df_b.name = df_a.name AND df_b.age != df_a.age ) """) result_df.show()
该方法会得到与上述方案完全一致的结果。
总结
这类基于关联键的列值对比场景,必须通过关联(Join)或子查询实现——本质上子查询的底层逻辑也是关联操作。subtract仅适用于整行去重/对比的简单场景,无法满足你需要的"基于子集列对比并返回完整行"的需求。
内容的提问来源于stack exchange,提问作者Question-er XDD
相关产品推荐
相关产品推荐

