PySpark多输出循环中DataFrame结果存取及函数改造技术问询
问题分析与解决方案
咱们直接拆解你的两个疑问,顺便修正代码里的小细节问题:
疑问1:为什么TRANSFORMS里的结果无法正常访问?
你当前的代码里,TRANSFORMS.append(multi_output)是把未执行的函数对象添加到了列表中,而不是函数运行后得到的DataFrame结果。这就是为什么你打印TRANSFORMS[0]时看到的是晦涩的函数描述——转换操作确实没执行,你只是把函数存起来了,根本没调用它。
另外还有个容易踩的小坑:在table_name=="TWO"的分支里,你写的是output_table== Input_table.drop("gender"),这是相等比较运算符,不是赋值操作,会导致这个分支返回布尔值,后续调用函数时必然会报错(其他分支返回的是DataFrame,类型完全不匹配)。
修正后的TRANSFORMS访问方式
修改循环中的代码,在添加到列表前主动执行函数,传入对应的输入表:
TRANSFORMS = [] DATASETS = { "ONE" : df_1, "TWO" : df_2, "THREE" : df_3, } for table_name, table_location in list(DATASETS.items()): def multi_output(Input_table, table_name=table_name): if table_name=="ONE": output_table = Input_table.drop("name") elif table_name=="TWO": output_table = Input_table.drop("gender") # 把==改成=,修复赋值错误 elif table_name=="THREE": output_table = Input_table.drop("salary") return output_table # 执行函数并将结果存入列表,而不是存函数本身 TRANSFORMS.append(multi_output(table_location))
现在你执行print(TRANSFORMS[0].show())就能看到第一个转换后的DataFrame内容了。
疑问2:如何生成命名规范的结果数据集?
不推荐动态生成df_1_result这类全局变量(容易让代码变得混乱,后续维护起来很麻烦),更优雅的方式是用字典存储结果,键名可以完全匹配你想要的命名规则:
方案1:用字典存储结果(强烈推荐)
RESULT_DATASETS = {} for table_name, table_location in DATASETS.items(): if table_name=="ONE": result_df = table_location.drop("name") RESULT_DATASETS["df_1_result"] = result_df elif table_name=="TWO": result_df = table_location.drop("gender") RESULT_DATASETS["df_2_result"] = result_df elif table_name=="THREE": result_df = table_location.drop("salary") RESULT_DATASETS["df_3_result"] = result_df # 后续分析时直接通过字典键访问 RESULT_DATASETS["df_1_result"].show() RESULT_DATASETS["df_2_result"].describe().show()
方案2:动态生成全局变量(不推荐)
如果你一定要生成df_1_result这类独立变量,可以用globals()函数动态创建,但这种方式会让代码可读性变差,不建议在生产代码中使用:
for idx, (table_name, table_location) in enumerate(DATASETS.items(), 1): if table_name=="ONE": result_df = table_location.drop("name") elif table_name=="TWO": result_df = table_location.drop("gender") elif table_name=="THREE": result_df = table_location.drop("salary") # 动态生成符合规则的变量名 globals()[f"df_{idx}_result"] = result_df # 现在可以直接访问这些变量 df_1_result.show() df_3_result.printSchema()
优化建议
其实可以把转换逻辑封装成一个独立函数,避免在循环内定义函数(容易触发闭包陷阱,比如之前的table_name参数可能会被循环最后一次的值意外覆盖),优化后的完整代码:
import pyspark from pyspark.sql import SparkSession from pyspark.sql.functions import * spark = SparkSession.builder.appName('Sparky').getOrCreate() # 创建初始DataFrame data = [("James","M",60000),("Michael","M",70000), ("Robert",None,400000),("Maria","F",500000), ("Jen","",None)] columns = ["name","gender","salary"] df_when = spark.createDataFrame(data = data, schema = columns) # 创建三个相同的数据集 df_1 = df_when df_2 = df_when df_3 = df_when # 独立的转换函数,逻辑更清晰 def transform_dataset(input_df, table_name): if table_name == "ONE": return input_df.drop("name") elif table_name == "TWO": return input_df.drop("gender") elif table_name == "THREE": return input_df.drop("salary") else: return input_df # 默认返回原表,避免无匹配时报错 DATASETS = { "ONE" : df_1, "TWO" : df_2, "THREE" : df_3, } # 存储转换结果 TRANSFORMS = [] RESULT_DATASETS = {} for table_name, table_df in DATASETS.items(): transformed_df = transform_dataset(table_df, table_name) TRANSFORMS.append(transformed_df) # 自动匹配命名规则存入字典 idx = list(DATASETS.keys()).index(table_name) + 1 RESULT_DATASETS[f"df_{idx}_result"] = transformed_df # 验证结果 TRANSFORMS[0].show() RESULT_DATASETS["df_2_result"].show()
内容的提问来源于stack exchange,提问作者Largo Terranova
相关产品推荐
相关产品推荐

