You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 12:24:04