Pyspark编写Case When语句报错,如何实现按字段值取对应列数据?
PySpark多条件分支取值报错修复方案
你代码的核心错误是列判断的语法使用错误:F.col()方法仅接收列名字符串作为入参,你将等值判断逻辑写在了F.col()的括号内部,会导致PySpark错误尝试读取名为MINS_TRA == "1"的不存在列,直接触发运行报错。
修复后代码
ABC_SCHED = ABC_SCHED.withColumn("ABC123", F.when(F.col("MINS_TRA") == "1", ABC_SCHED.T1_0to30) .when(F.col("MINS_TRA") == "2", ABC_SCHED.T2_31to60) .when(F.col("MINS_TRA") == "3", ABC_SCHED.T3_61to90) .when(F.col("MINS_TRA") == "4", ABC_SCHED.T4_91to120) .when(F.col("MINS_TRA") == "5", ABC_SCHED.T4_120ORMORE) .otherwise(F.lit(None)) )
可选简化写法(适合后续分支规则扩展场景)
如果后续判断规则会频繁新增调整,可以用字典映射的方式简化代码,避免重复写when分支:
# 定义匹配规则 match_rule = { "1": "T1_0to30", "2": "T2_31to60", "3": "T3_61to90", "4": "T4_91to120", "5": "T4_120ORMORE" } # 构造映射列自动匹配,不存在的key自动返回空值 map_col = F.create_map(*[F.lit(item) for kv in match_rule.items() for item in kv]) ABC_SCHED = ABC_SCHED.withColumn("ABC123", map_col[F.col("MINS_TRA")])
内容的提问来源于stack exchange,提问作者confused101
相关产品推荐
相关产品推荐

