Spark中将指定列数据转为LocalDate解决日期过滤错误问题
Spark 日期过滤错误修复方案
错误原因
你当前的逻辑是直接对dd/MM/yyyy格式的字符串做大小比较,字符串比较规则是逐位按ASCII值判断,02/12/2022的首字符0小于22/09/2021的首字符2,因此会被误判为满足小于条件,出现错误结果。
推荐方案:使用Spark内置日期函数(性能最优)
优先使用Spark原生的to_date函数完成字符串到日期类型的转换,无需自定义UDF,执行效率更高:
// 导入所需内置函数 import org.apache.spark.sql.functions.{to_date, current_date} val filteredDF = sampleData.where( // 将字符串列按dd/MM/yyyy格式转换为Date类型,和当前日期直接比较 to_date(sampleData("data_riferimento_condizioni"), "dd/MM/yyyy") < current_date() )
如果需要固定对比某一指定日期而非运行时的当前日期,也可以传入LocalDate实例:
import java.time.LocalDate import org.apache.spark.sql.functions.{to_date, lit} import org.apache.spark.sql.types.DateType val targetDate = LocalDate.of(2021,9,22) // 也可以直接用LocalDate.now获取当前日期 val filteredDF = sampleData.where( to_date(sampleData("data_riferimento_condizioni"), "dd/MM/yyyy") < lit(targetDate).cast(DateType) )
可选方案:自定义UDF使用isBefore方法
如果你明确需要调用LocalDate.isBefore方法实现判断,可以通过自定义UDF实现,注意该方案性能低于内置函数,适合小数据量或逻辑复杂的场景:
import java.time.LocalDate import java.time.format.DateTimeFormatter import org.apache.spark.sql.functions.udf // 全局复用Formatter,避免重复初始化 private val dateFormatter = DateTimeFormatter.ofPattern("dd/MM/yyyy") // 定义UDF,可按需添加空值、非法格式的异常处理逻辑 val isBeforeToday = udf((dateStr: String) => { if (dateStr == null || dateStr.trim.isEmpty) { false // 空值处理规则可按需调整 } else { try { LocalDate.parse(dateStr, dateFormatter).isBefore(LocalDate.now()) } catch { case e: Exception => false // 非法格式日期处理规则可按需调整 } } }) // 执行过滤 val filteredDF = sampleData.where(isBeforeToday(sampleData("data_riferimento_condizioni")))
注意事项
- 若
data_riferimento_condizioni列存在大量非法格式的日期,建议先做数据清洗再执行过滤逻辑 - 生产环境优先选择内置函数方案,避免UDF带来的性能损耗
内容的提问来源于stack exchange,提问作者Miko
相关产品推荐
相关产品推荐

