PySpark在for循环中执行过滤返回空DataFrame如何解决?
问题根因
你的代码中过滤国家的逻辑写死了匹配值,F.col('country') == '+country+'里的'+country+'是固定字符串,不会读取for循环的country变量值,所以没有匹配到任何数据,最终输出空DataFrame。
修复方法
把国家过滤行的'+country+'替换为变量country即可,修正后的对应代码段:
for country in COUNTRY_LIST: print('Begin '+country) # 修正国家过滤条件,使用循环变量country poi_co = poi_reference.filter(F.col('country') == country) poi_ret = poi_co.filter(F.col('l1_name') == 'Retail') poi_final = poi_ret.dropDuplicates(['unique_id']) # 其余逻辑保持不变
附加排查步骤
如果修改后依然输出空结果,可按以下步骤验证:
- 在过滤国家后添加计数打印,确认是否有匹配数据:
poi_co = poi_reference.filter(F.col('country') == country) print(f"{country} 国家匹配数据量:{poi_co.count()}")
- 确认
poi_reference表中的country字段取值和COUNTRY_LIST的元素大小写、前后空格完全一致,存在差异的话可调整匹配逻辑:
# 去除前后空格+统一大写后匹配 poi_co = poi_reference.filter(F.upper(F.trim(F.col('country'))) == country)
- 同样验证
l1_name字段是否存在Retail取值,可加print(f"Retail分类匹配数据量:{poi_ret.count()}")确认。
内容的提问来源于stack exchange,提问作者Zaid
相关产品推荐
相关产品推荐

