Scala DataFrame如何基于timestamp列筛选获取最近6个月的数据
Scala Spark DataFrame 按时间戳列过滤最近6个月数据的实现方案
首先确认你的timestampCol列类型为TimestampType,若为字符串格式可先做类型转换:
import org.apache.spark.sql.functions._ val dfWithTs = df.withColumn("timestampCol", to_timestamp(col("timestampCol"), "yyyy-MM-dd HH:mm:ss"))
方案1:内置日期函数直接过滤(通用场景首选)
直接利用Spark SQL内置的时间函数计算阈值,无需额外处理月份天数、时区等兼容问题:
import org.apache.spark.sql.functions._ val filteredDf = df.filter( col("timestampCol") >= add_months(current_timestamp(), -6) )
- 优点:写法简洁,自动适配不同月份的天数差异,支持时区配置
- 若需指定时区,可在执行过滤前设置会话时区:
spark.conf.set("spark.sql.session.timeZone", "Asia/Shanghai")
方案2:预计算阈值广播过滤(超大数据量场景优化)
如果你的数据量极大,可提前在Driver端计算好阈值,广播到所有Executor避免每行重复执行时间计算逻辑,提升性能:
import java.time.LocalDateTime import java.sql.Timestamp // Driver端预计算6个月前的时间阈值 val sixMonthsAgoThreshold = Timestamp.valueOf(LocalDateTime.now().minusMonths(6)) // 广播阈值减少重复传输 val broadcastThreshold = spark.sparkContext.broadcast(sixMonthsAgoThreshold) val filteredDf = df.filter( col("timestampCol") >= lit(broadcastThreshold.value) )
扩展:按自然月过滤
如果你需要保留的是完整自然月的数据(比如当前是2024年7月,保留2024年1月及以后的所有数据,不考虑7月的具体时间),可使用trunc函数对日期做月份截断:
val filteredDf = df.filter( trunc(col("timestampCol"), "month") >= trunc(add_months(current_timestamp(), -6), "month") )
内容的提问来源于stack exchange,提问作者NAVNEET MATHPAL
相关产品推荐
相关产品推荐

