Pyspark提取Timestamp列年份出现TypeError/AttributeError报错求助
PySpark从timestamp列提取年份的修正方案
你报错的核心原因是混用了Python原生datetime库的API和Spark分布式列处理API,Python原生datetime方法只能处理单个时间值,不能直接作用于Spark的Column类型对象。
最优实现方案(性能最好,推荐使用)
直接使用Spark SQL内置的year函数处理已生成的REPORT_TIMESTAMP列即可:
- 导入依赖函数
from pyspark.sql.functions import year
- 执行列新增操作
jsonDf = jsonDf.withColumn("YEAR", year("REPORT_TIMESTAMP"))
原有错误写法的修正说明
第一种写法修正
原写法报错是因为datetime.fromtimestamp要求传入Python整型数值,而你传入的是Spark Column类型。如果一定要用Python datetime逻辑实现,需要封装为UDF使用:
from pyspark.sql.functions import udf from pyspark.sql.types import IntegerType import datetime # 定义处理单值的UDF year_extract_udf = udf(lambda ts: datetime.fromtimestamp(ts).year if ts is not None else None, IntegerType()) # 调用UDF处理列 jsonDf = jsonDf.withColumn("YEAR", year_extract_udf(jsonDF.reportData.timestamp.cast("integer")))
注意:该方案性能远低于内置
year函数,非必要不推荐使用。
第二种写法说明
原写法不存在可修正的方案,datetime.date.to_timestamp是不存在的API,属于将Spark的to_timestamp函数和Python datetime库API混淆误用,建议直接使用上述内置函数方案。
- 注意代码中变量大小写保持一致,你示例代码中同时出现
jsonDf和jsonDF,大小写不一致会触发未定义变量报错。
内容的提问来源于stack exchange,提问作者Moritz
相关产品推荐
相关产品推荐

