PySpark中如何从日期和小时列生成时间戳列?
错误原因分析
你的错误根源是在UDF内部误用了PySpark的Column函数(F.year/F.month等)。UDF接收的是每行的原生Python值(比如date列是datetime.date对象,hour是整数),而非Spark的Column对象。F.year这类函数只能作用于Column,在UDF里调用会因为Python worker进程中没有可用的SparkContext(sc),导致sc._jvm为None,触发AttributeError。
解决方案
方案1:使用Spark内置函数(推荐,性能远优于UDF)
Spark原生提供了高效的时间拼接转换方法,无需自定义UDF:
from pyspark.sql import functions as F # 方法1:字符串拼接后转Timestamp new_df = df.withColumn( "current_ts", F.to_timestamp( F.concat_ws(" ", F.col("date"), F.lpad(F.col("hour"), 2, "0")), "yyyy-MM-dd HH" ) ) # 方法2:Spark 3.0+可用的make_timestamp函数(更直观) new_df = df.withColumn( "current_ts", F.make_timestamp( F.year("date"), F.month("date"), F.dayofmonth("date"), F.col("hour"), F.lit(0), F.lit(0.0) ) )
方案2:修正UDF写法(处理原生Python对象)
如果必须使用UDF,需直接操作Python原生类型,不要调用Spark Column函数:
from datetime import datetime from pyspark.sql import functions as F from pyspark.sql.types import TimestampType def buildTimestamp(date, hour): # 这里的date是datetime.date对象,hour是整数 return datetime(date.year, date.month, date.day, hour, 0, 0) # 显式指定返回类型为TimestampType,避免自动推断出错 buildTimestampUDF = F.udf(buildTimestamp, TimestampType()) new_df = df.withColumn("current_ts", buildTimestampUDF(F.col("date"), F.col("hour")))
内容的提问来源于stack exchange,提问作者user2128702
相关产品推荐
相关产品推荐

