PySpark:不修改原DataFrame批量复制列值及优化实现
问题解答
1. 不修改初始DataFrame的实现方法
完全可以做到。你原来的代码有两个问题:
- 每次循环都基于原始的
df生成新的df_return,导致前面循环添加的列都被覆盖,最后只保留了最后一列的修改; F.lit(df.share)是错误用法:lit()用于将常量值转为列对象,而这里要引用已有的share列,应该用F.col("share")或者直接df.share。
修正后的代码全程不会修改原始的df:
import pyspark.sql.functions as F # 先把初始df赋值给df_return,作为操作的起始点 df_return = df for column in columns: # 每次在df_return的基础上添加/更新目标列 df_return = df_return.withColumn(column, F.col("share"))
这样df始终保持原始状态,所有修改都在df_return上累积,最终就能得到所有目标列都被赋值为share列值的结果。
2. 更简单高效的实现方式
Spark DataFrame是不可变结构,多次循环调用withColumn虽然可行,但可以用以下两种更简洁高效的方式,减少中间临时对象的生成:
方法一:用reduce链式处理
借助functools.reduce,从初始df开始依次为每个目标列赋值,代码更紧凑:
from functools import reduce import pyspark.sql.functions as F df_return = reduce(lambda temp_df, col: temp_df.withColumn(col, F.col("share")), columns, df)
方法二:用select一次性生成所有列
如果需要保留原DataFrame的所有列,同时新增/替换目标列,可以直接用select一次性完成:
import pyspark.sql.functions as F # "*" 保留原所有列,后面的列表生成所有目标列并赋值为share的值 df_return = df.select("*", *[F.col("share").alias(col) for col in columns])
如果目标列已经存在于原DataFrame中,这个写法会直接替换该列的值;如果不存在,则会新增列,效果和withColumn完全一致。
内容的提问来源于stack exchange,提问作者paulo
相关产品推荐
相关产品推荐

