如何在Spark中使用LocalDateTime变量为DataFrame添加额外日期列
错误原因
你的报错本质是字符串插值生成的SQL语句语法错误:你直接把
2020-01-02 00:00:00作为裸值传入expr,SQL解析器会把空格后的00识别为非法的额外Token,所以抛出解析异常。
问题解答
1. 是否可以直接用LocalDateTime新增列
可以,Spark 2.3及以上版本原生支持将java.time.LocalDateTime类型变量直接作为常量列传入DataFrame,不需要手动做字符串格式化转换。
2. 最优实现方式
优先使用lit函数直接传入日期对象,避免手动处理格式带来的解析风险:
import org.apache.spark.sql.functions.lit import java.time.LocalDateTime val loadingDate: LocalDateTime = LocalDateTime.of(2020, 1, 2, 0, 0, 0) val resultDF = DF.withColumn("dttm", lit(loadingDate))
如果使用的是不支持直接传LocalDateTime的旧版本Spark,可以转成java.sql.Timestamp后再传入:
val sqlTimestamp = java.sql.Timestamp.valueOf(loadingDate) val resultDF = DF.withColumn("dttm", lit(sqlTimestamp))
如果一定要用你原本的expr写法,需要给格式化后的日期字符串加上单引号包裹,保证SQL语法合法,但该方式不推荐:
val formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss") DF.withColumn("dttm", expr(s"'${loadingDate.format(formatter)}'").cast("timestamp"))
3. 推荐日期类型选择
- 带时分秒的时间戳场景:使用Spark的
TimestampType,对应Java侧的LocalDateTime类型 - 仅日期不带时间的场景:使用Spark的
DateType,对应Java侧的LocalDate类型
不要用字符串类型存储日期,既会占用更多存储空间,也会导致后续日期计算、过滤等操作的复杂度和出错概率上升。
内容的提问来源于stack exchange,提问作者Alexander Lopatin
相关产品推荐
相关产品推荐

