PySpark内连接不等条件失效,求df1与df2 studyid不匹配记录统计方法
PySpark 统计 df1 中 studyid 未在 df2 出现的记录总数解决方案
原有代码错误原因
你写的内连接不等值匹配逻辑完全不符合需求:df1.join(df2, df1.studyid != df2.studyid, how='inner')会生成两个表所有 studyid 不相等的行的笛卡尔积,最终结果会包含大量重复的 df1 行,无法筛选出 df1 独有的 studyid 对应的记录。
正确实现方案
方法1:left anti 连接(最优方案,适合大数据量场景)
PySpark 提供的left_anti连接就是专门为这种「取左表未在右表匹配到关联键的行」的场景设计的,分布式执行性能高,不会把数据拉取到 Driver 端,也不会产生多余字段。
# 按studyid关联,仅保留df1中没有在df2匹配到对应studyid的行 unmatched_df = df1.join(df2, on="studyid", how="left_anti") # 统计符合条件的记录总数 unmatched_count = unmatched_df.count()
方法2:isin 取反(仅适合 df2 的 studyid 量级较小的场景)
如果 df2 的 studyid 去重后总量很小,可以先把 df2 的 studyid 收集为本地集合,再过滤 df1。注意该方法会把 df2 的 studyid 全量拉到 Driver 端,数据量大时容易触发 OOM。
from pyspark.sql.functions import col # 收集df2所有去重后的studyid df2_studyid_list = df2.select("studyid").distinct().rdd.flatMap(lambda x: x).collect() # 过滤df1中studyid不在上述集合的行 unmatched_df = df1.filter(~col("studyid").isin(df2_studyid_list)) # 统计符合条件的记录总数 unmatched_count = unmatched_df.count()
内容的提问来源于stack exchange,提问作者subhajit saha
相关产品推荐
相关产品推荐

