Spark DataFrame基于现有列新增日期列UDF报错如何解决
错误原因
- 你定义UDF时绑定的lambda仅接收1个入参,但调用UDF时传入了
date_range_start、date_range_end两个参数,直接触发参数数量不匹配的报错。 - 返回类型声明错误:
get_dates函数返回的是日期列表,你指定的返回类型是单个日期类型DateType(),即便参数问题修复后续也会出现类型不匹配报错。 - 附加问题:pandas
date_range返回的是DatetimeIndex类型,没有转成Spark可识别的Python date对象列表,会触发序列化错误。
修复方案
方案1:修正自定义UDF
修改后的完整代码如下:
from datetime import datetime, timedelta import pandas as pd from pyspark.sql.types import ArrayType, DateType def get_dates(s, e): start = datetime.strptime(s, '%Y-%m-%d').date() end = datetime.strptime(e, '%Y-%m-%d').date() # 将pandas日期对象转为Python原生date列表,适配Spark序列化 return [d.date() for d in pd.date_range(start, end - timedelta(days=1), freq='d')] # 直接传入函数名,无需额外lambda,声明返回类型为日期数组 udf_get_dates = udf(get_dates, ArrayType(DateType())) df = df.withColumn('date_bet_dates', udf_get_dates(df['date_range_start'], df['date_range_end'])) df.show(3, truncate=False)
方案2:使用Spark原生函数(优先推荐)
Spark 2.4及以上版本支持原生sequence函数实现相同需求,无需自定义UDF,性能高数十倍:
from pyspark.sql import functions as F df = df.withColumn('date_bet_dates', F.sequence( F.to_date('date_range_start'), F.date_sub(F.to_date('date_range_end'), 1), F.expr('INTERVAL 1 DAY') ) ) df.show(3, truncate=False)
内容的提问来源于stack exchange,提问作者9ganzi
相关产品推荐
相关产品推荐

