如何在PySpark中使用withColumns基于单个条件创建多列?
解决PySpark一次性基于条件创建多列的问题
你的写法问题所在
你原来的代码无法运行,核心原因是withColumns要求传入列名到列表达式的字典,而不是直接把when的返回值作为参数——when返回的是单个Column对象,并非多列的字典结构,因此不符合withColumns的参数要求。
正确实现方式
完全不需要使用UDF,用PySpark内置函数就能一次性完成多列创建,只需给每个新列单独定义when/otherwise逻辑,再把所有列的表达式打包成字典传给withColumns即可:
新增列(原DataFrame无目标列)
如果是新增列,不符合条件时可以设为null或其他默认值:
import pyspark.sql.functions as F df = df.withColumns({ "new_c1": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c2"))).otherwise(None), "new_c2": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c3"))).otherwise(None), "new_c3": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c4"))).otherwise(None) })
更新已有列
如果目标列已经存在,想保留不符合条件时的原列值,可以这样写:
df = df.withColumns({ "new_c1": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c2"))).otherwise(F.col("new_c1")), "new_c2": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c3"))).otherwise(F.col("new_c2")), "new_c3": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c4"))).otherwise(F.col("new_c3")) })
关于withColumns和when/otherwise的疑问
withColumns本身不支持直接把when作为整体参数,但每个列的表达式可以单独使用when/otherwise——本质是给每个新列单独定义条件逻辑,再批量传入withColumns实现一次性处理,效果和多次调用withColumn一致,但代码更简洁。
后续添加多条件的扩展
如果之后要增加更多条件,直接在when后面链式调用即可,示例如下:
"new_c1": F.when(F.col("age") < 6, F.least(F.col("c1"), F.col("c2"))) .when(F.col("age").between(6, 12), F.least(F.col("c2"), F.col("c3"))) .otherwise(None)
内容的提问来源于stack exchange,提问作者Chuck
相关产品推荐
相关产品推荐

