PySpark内连接异常求助:同键同行数却返回空DataFrame
问题原因分析
你的问题是Spark Catalyst优化器的谓词下推或列消除优化导致的意外行为:
当你先对table_2执行filter(筛选日期)再执行select(仅保留部分列)时,Spark的优化器可能会重新排列操作顺序,将filter推到select之后执行。由于select已经移除了application_date列,此时filter(F.col('application_date') == '2023-10-09')会因列不存在被视为永假条件,导致table_2实际返回0行,最终内连接结果为0。
而你测试的其他有效方法,本质上都是阻止了这种错误的优化:
table_2.select('*')保留了application_date列,filter条件能正常生效,table_2保持473行;table_1.select('*')或显式指定连接条件table_1['application_id']==table_2['application_id'],改变了优化器的执行计划逻辑,避免了filter被错误下推;- 左连接会保留
table_1的所有行,即使table_2为空也会返回473行(但这不是解决内连接问题的根本办法); - 转Pandas后合并是在本地执行,不受Spark优化器影响。
验证方法
你可以通过查看执行计划确认这个问题:
table_2.explain(extended=True)
如果输出中显示Filter操作位于Project(即select)之后,且引用了application_date列,即可验证是优化器的问题。
解决方案
以下几种方法可以解决这个问题:
1. 调整操作顺序(推荐)
先选择包含application_date的列,执行filter后再选择目标列,确保filter在select移除application_date之前执行:
table_2 = spark.table(table_2_path) # 先选包含过滤所需列的集合,过滤后再选目标列 table_2 = table_2.select('application_id', 'application_date', 'col1', 'col2', 'col3') \ .filter(F.col('application_date') == '2023-10-09') \ .select(*cols_of_interest)
2. 使用显式连接条件
直接指定连接的列等式,绕过字符串连接键的优化逻辑:
print(table_1.join(table_2, on=table_1.application_id == table_2.application_id, how='inner').count())
3. 缓存中间结果阻止优化
对过滤后的table_2执行缓存,强制Spark保留过滤后的数据集,避免优化器重新排列操作:
table_2 = spark.table(table_2_path) table_2 = table_2.filter(F.col('application_date') == '2023-10-09').cache() table_2 = table_2.select(*cols_of_interest)
4. 临时禁用谓词下推(不推荐用于生产)
使用hint阻止优化器将filter下推到select之后:
table_2 = spark.table(table_2_path) table_2 = table_2.filter(F.col('application_date') == '2023-10-09').hint("NO_PUSH_DOWN") table_2 = table_2.select(*cols_of_interest)
内容的提问来源于stack exchange,提问作者Funy
相关产品推荐
相关产品推荐

