Spark中用复杂when逻辑替代withColumn的Python2.7兼容方案问询
兼容Python 2.7的Spark Select优化方案
问题原因
Python 2.7不支持在函数参数中对多个可迭代对象使用*解包语法,这就是你调用df.select(*cols_to_keep, *cols_transformed1, *cols_transformed2)时触发语法错误的根本原因。
解决方案
方案1:合并所有列列表为单个列表
将需要保留和转换的列列表合并成一个大列表,再一次性传递给select方法:
from pyspark.sql import functions as F cols_to_keep = [c for c in df.columns if c not in column_list] cols_transformed1 = [ F.when(...).otherwise( F.when(...).otherwise(F.when(...).otherwise(...)) ).alias(c + "_new_name") for c in column_list ] cols_transformed2 = [ (df["result"] + F.abs(df[col_name])).alias("result") for col_name in column_list ] # 合并所有列列表 all_cols = cols_to_keep + cols_transformed1 + cols_transformed2 # 传递合并后的列表给select df = df.select(*all_cols)
方案2:使用itertools.chain拼接迭代器
如果不想手动合并列表,可以用itertools.chain来拼接多个列迭代器,同样适配Python 2.7:
from pyspark.sql import functions as F import itertools cols_to_keep = [c for c in df.columns if c not in column_list] cols_transformed1 = [ F.when(...).otherwise( F.when(...).otherwise(F.when(...).otherwise(...)) ).alias(c + "_new_name") for c in column_list ] cols_transformed2 = [ (df["result"] + F.abs(df[col_name])).alias("result") for col_name in column_list ] # 用chain拼接所有列迭代器 df = df.select(*itertools.chain(cols_to_keep, cols_transformed1, cols_transformed2))
注意事项
cols_transformed2中必须给df["result"] + F.abs(df[col_name])加上括号,否则alias会仅作用在abs(df[col_name])上,导致逻辑错误。- 虽然
withColumn循环写法能运行,但每次调用withColumn都会生成新的DataFrame执行计划,多次循环会导致计划膨胀、性能下降,因此优先推荐上述select优化方案。
内容的提问来源于stack exchange,提问作者user2441441
相关产品推荐
相关产品推荐

