PySpark如何根据ERROR列值对DataFrame多列做条件更新
PySpark 按列条件实现多列更新的实现方案
不需要找专门支持多列同时更新的特殊语法,两种成熟写法都能满足需求,且可以完全匹配给出的预期输出效果:
方案1:逐列套用when().otherwise()(性能优先推荐)
when-otherwise本身是按行做条件判断返回列值,只需要给每个需要更新的列单独套一层相同的判断逻辑即可,Spark底层会自动优化判断逻辑,不会重复扫描数据,是性能最好的实现方式。
示例代码:
from pyspark.sql.functions import current_timestamp, lit, col df = df.withColumn( "LAST_UPDATE_DATE", # ERROR为空时更新为当前时间,否则保留列原有值 when(col("ERROR").isNull(), current_timestamp()).otherwise(col("LAST_UPDATE_DATE")) ).withColumn( "ADDR_1", # ERROR为空时设为ADDR_1,非空时按需求设为"0",如果要保留原值就换成col("ADDR_1") when(col("ERROR").isNull(), lit("ADDR_1")).otherwise(lit("0")) ).withColumn( "ADDR_2", # ERROR为空时设为ADDR_2,否则保留列原有值 when(col("ERROR").isNull(), lit("ADDR_2")).otherwise(col("ADDR_2")) )
如果要完全匹配贴出的样例输出(ERROR非空时ADDR_1显示null),只需要把ADDR_1列otherwise分支的lit("0")改成col("ADDR_1")即可,样例中null值是原始数据本身的值,不是更新逻辑生成的。
方案2:拆分数据集分别处理后合并(适合多列复杂更新场景)
如果单个条件分支下要修改的列数量很多,逐列写判断逻辑太繁琐,可以按判断条件把DataFrame拆成两个子集,分别执行不同的更新逻辑后再用unionByName合并,代码可读性更高。
示例代码:
from pyspark.sql.functions import current_timestamp, lit, col # 处理ERROR为空的子集 df_null_err = df.filter(col("ERROR").isNull()) \ .withColumn("LAST_UPDATE_DATE", current_timestamp()) \ .withColumn("ADDR_1", lit("ADDR_1")) \ .withColumn("ADDR_2", lit("ADDR_2")) # 处理ERROR非空的子集 df_nonnull_err = df.filter(col("ERROR").isNotNull()) \ .withColumn("ADDR_1", lit("0")) # 如果不需要修改ADDR_1,直接删掉上面的withColumn即可保留原始值 # 合并两个子集,自动按列名匹配,避免列顺序不一致导致的数据错误 result_df = df_null_err.unionByName(df_nonnull_err)
注意:使用拆分合并方案时,不要用普通
union,union是按列位置匹配值,列顺序变了就会出现数据错位,unionByName按列名匹配更稳妥。
内容的提问来源于stack exchange,提问作者Sonali Bisht
相关产品推荐
相关产品推荐

