PySpark UDF比较日期列写入时报类型错误如何解决
解决方案
报错的核心原因有两个:
- 你的数据中
x、CRE_DT、cont_DT三个字段存在空值(None),show方法默认只加载前20条数据,没有命中空值行,全量写入时触发空值比较错误 - 你在Python UDF中混用了位运算符
&做逻辑判断,且直接将日期类型的cont_DT和字符串'9999-12-31'比较,存在类型匹配隐患
方案1:修改UDF兼容空值和类型问题
def getStatus(x, CRE_DT, cont_DT): # 先做空值判断,空值返回状态可根据业务需求调整,示例返回Unknown if x is None or CRE_DT is None or cont_DT is None: return 'Unknown' if x < CRE_DT: return 'Not Active' # 日期类型转成字符串后再比对,也可直接和datetime.date(9999,12,31)比对 elif str(cont_DT) == '9999-12-31': return 'Active' # Python逻辑判断用and,不要用位运算符& elif x > CRE_DT and cont_DT > x and str(cont_DT) != '9999-12-31': return 'Active' return 'Stop' func_udf = udf(getStatus, StringType())
方案2:替换为Spark原生函数(更推荐)
UDF会带来额外的Python-JVM序列化/反序列化开销,用原生when函数性能更好,也不会出现跨环境类型转换问题,代码示例:
from pyspark.sql.functions import when, lit, to_date, col for i in range(2015, 2022): end_date = to_date(lit(f'{i}-12-31')) # 先写正常逻辑 status_col = when(end_date < col("CRE_DT"), "Not Active") \ .when(col("cont_DT") == to_date(lit("9999-12-31")), "Active") \ .when((end_date > col("CRE_DT")) & (col("cont_DT") > end_date), "Active") \ .otherwise("Stop") # 单独处理空值场景 status_col = when(col("CRE_DT").isNull() | col("cont_DT").isNull(), "Unknown").otherwise(status_col) df = df.withColumn(f'{i}-12-31', status_col)
内容的提问来源于stack exchange,提问作者Sreel
相关产品推荐
相关产品推荐

