pyspark.pandas中datetime日期加减运算报错该如何处理
pyspark.pandas日期加减运算报错解决方案
报错原因
pyspark.pandas的datetime序列运算逻辑对齐Spark SQL标准,不支持直接和Python原生timedelta、pandas的pd.TimeDelta做算术运算,因此执行加法时会抛出TypeError: Addition can not be applied to datetimes.错误。
正确实现方法
方法1:使用pyspark.pandas自带的Timedelta运算
这是最贴近pandas写法的实现,无需额外调用Spark函数:
import pyspark.pandas as ps import pandas as pd df = pd.DataFrame({'year': [2015, 2016], 'month': [2, 3], 'day': [4, 5]}) df = ps.DataFrame(df) srs = ps.to_datetime(df) # 用ps.Timedelta定义时间差,支持天、小时、分钟等多粒度偏移 result = srs + ps.Timedelta(days=3) print(result)
方法2:调用Spark SQL内置函数(适配复杂偏移场景)
如果涉及更灵活的时间偏移逻辑,可以直接调用Spark SQL的时间函数实现:
from pyspark.sql.functions import expr # 加3天 result = srs.spark.transform(lambda col: expr("date_add(col, 3)")) # 加2小时示例 result = srs.spark.transform(lambda col: expr("timestampadd(hour, 2, col)"))
注意事项
- 调用
to_pandas()转为pandas序列再运算的方式仅适合小数据量测试,该操作会把全量数据拉取到Driver节点,大数据量下会出现Driver内存溢出、任务运行失败的问题,生产环境不推荐使用。 - 所有时间偏移相关的运算尽量在pyspark.pandas的分布式框架内完成,不要做不必要的类型转换。
内容的提问来源于stack exchange,提问作者kaz0621
相关产品推荐
相关产品推荐

