Spark DataFrame时间戳列操作无报错但执行失败
Spark DataFrame时间戳处理失败的问题排查与修复
首要问题:语法不完整
你写的两行withColumn代码都没给when函数加闭合括号。Spark是惰性执行模式,定义DataFrame转换逻辑时不会立刻校验语法,只有当你执行show()、count()这类触发计算的操作时才会暴露错误,这就是为什么没提前报错但执行失败的原因。
第二个问题:数值溢出风险
sys.maxsize是Python的整数最大值(64位系统里是9223372036854775807),但Spark默认的整数类型是32位Int,最大值只有2147483647,直接把- (sys.maxsize - 1)赋值给列会导致数值溢出,最终操作失败。
修正后的代码
import sys from pyspark.sql.functions import current_date, col, when, months_between, days_between # 补全括号,替换为Spark兼容的最小Int值 aggregate = aggregate.withColumn( 'DaysSinceFirstUsage', when(months_between(current_date(), col('FirstUsage')) > 120, -2147483648) .otherwise(days_between(current_date(), col('FirstUsage'))) ) aggregate = aggregate.withColumn( 'DaysSinceLastUsage', when(months_between(current_date(), col('LastUsage')) > 120, -2147483648) .otherwise(days_between(current_date(), col('LastUsage'))) )
可选优化
如果需要更大的数值范围,避免溢出,可以把列类型指定为Long:
from pyspark.sql.functions import lit aggregate = aggregate.withColumn( 'DaysSinceFirstUsage', when(months_between(current_date(), col('FirstUsage')) > 120, lit(-sys.maxsize).cast("long")) .otherwise(days_between(current_date(), col('FirstUsage')).cast("long")) )
另外要确认FirstUsage和LastUsage列是Date或Timestamp类型,如果是字符串格式,需要先转成日期类型:col('FirstUsage').cast("date")
内容的提问来源于stack exchange,提问作者Vaibhav Raj
相关产品推荐
相关产品推荐

