PySpark函数传参时无法创建目标DataFrame问题排查
问题描述
我希望通过函数创建DataFrame:先从df_correct_countries中提取去重的NewCountry列并收集为Row对象列表;编写create_df函数接收原DataFrame和国家参数,过滤出对应国家的行并返回DataFrame;通过循环传参调用函数时无报错,但无法访问以国家命名的DataFrame(如EquatorialGuinea),提示NameError;但直接在函数外执行相同过滤逻辑却能正常生成DataFrame,请问问题出在哪里?
相关代码
获取国家列表代码
country_list=df_correct_countries.select('NewCountry').dropDuplicates().collect() for i in country_list: print(i)
输出
Row(NewCountry='Senegal') Row(NewCountry='Algeria') Row(NewCountry='Nigeria') Row(NewCountry='Morocco') Row(NewCountry='Ethiopia')
create_df函数代码
def create_df(df,cnt): cnt=str(cnt) cnt=df.where(col("NewCountry")==str(cnt)) return cnt
函数调用代码
for j in country_list: create_df(df_correct_countries,j['NewCountry'])
报错信息
display(EquatorialGuinea) ###Error NameError: name 'EquatorialGuinea' is not defined
外部正常执行代码
df_correct_countries.where(col("NewCountry")=='EquatorialGuinea')
问题原因
- 变量未被赋值:调用
create_df函数后,你没有将返回的DataFrame赋值给对应国家名称的变量。函数返回的结果如果不被接收,会直接被丢弃,不会自动创建EquatorialGuinea这类命名的变量。 - 局部变量不作用于外部:函数内部的
cnt是局部变量,仅在函数执行期间存在,返回后若不赋值给外部变量,不会进入当前命名空间。
解决方法
方法1:动态创建对应命名的变量
使用locals()动态创建变量,将函数返回的DataFrame赋值给以国家名命名的变量:
for j in country_list: country_name = j['NewCountry'] # 动态创建变量并赋值 locals()[country_name] = create_df(df_correct_countries, country_name)
之后就能直接通过display(EquatorialGuinea)访问对应的DataFrame。
方法2:用字典统一存储(更推荐)
避免大量零散变量导致命名混乱,用字典集中管理所有国家的DataFrame:
country_dfs = {} for j in country_list: country_name = j['NewCountry'] country_dfs[country_name] = create_df(df_correct_countries, country_name)
访问时通过字典键获取:display(country_dfs['EquatorialGuinea'])
优化create_df函数
原函数中cnt=str(cnt)属于冗余操作(传入的j['NewCountry']已经是字符串),可以简化:
def create_df(df, cnt): return df.where(col("NewCountry") == cnt)
内容的提问来源于stack exchange,提问作者Ajay S Pal
相关产品推荐
相关产品推荐

