Pyspark DataFrame如何替换列中指定值并生成新列
问题根因
你的代码核心问题是多次调用withColumn重写aa列时,otherwise分支始终取原始A列的值,导致前序的替换结果被后续操作完全覆盖:
- 第一次执行完
withColumn后,OTH/CON对应的aa确实被设为了Collect - 第二次执行时,只有
A等于Freight Collect的行才会被设为Collect,其余行(包括已经改对的OTH/CON行)的aa会被重新赋值为原始A列的值,第一次的修改直接失效 - 第三次执行时同理,只有
DBG行会被修改,其余行又被重置为原始A列的值,最终只剩最后一次匹配的规则生效
正确实现方式
推荐写法:链式调用when(仅一次列计算,性能最优)
直接把所有判断条件串在一次withColumn操作中,避免多次覆盖的问题:
from pyspark.sql import functions as F df = df.withColumn("aa", F.when(F.col("A").isin(["OTH/CON", "Freight Collect"]), F.lit("Collect")) .when(F.col("A") == "DBG", F.lit("Dispose")) .otherwise(F.col("A")) )
兼容分步修改的写法
如果确实需要拆分规则分步实现,把otherwise的取值从原始A列改为已经修改过的aa列即可:
from pyspark.sql import functions as F df = df.withColumn("aa", F.when(F.col("A").isin(["OTH/CON"]), F.lit("Collect")).otherwise(F.col("A"))) df = df.withColumn("aa", F.when(F.col("A").isin(["Freight Collect"]), F.lit("Collect")).otherwise(F.col("aa"))) df = df.withColumn("aa", F.when(F.col("A").isin(["DBG"]), F.lit("Dispose")).otherwise(F.col("aa")))
批量映射写法(适合替换规则多的场景)
如果后续需要新增大量替换规则,用字典映射的方式维护更清晰:
from pyspark.sql import functions as F replace_map = { "OTH/CON": "Collect", "Freight Collect": "Collect", "DBG": "Dispose" } map_col = F.create_map([F.lit(x) for item in replace_map.items() for x in item]) df = df.withColumn("aa", F.coalesce(map_col[F.col("A")], F.col("A")))
以上任意一种写法执行后,都可以得到你期望的输出结果。
内容的提问来源于stack exchange,提问作者EnigmAI
相关产品推荐
相关产品推荐

