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

如何在Spark 2(Scala)的DataFrame中筛选近10小时数据?

在Spark 2(Scala)中筛选最近10小时的DataFrame行

嘿,这个需求其实挺常见的,我来给你一步步拆解实现方法:

首先,核心思路是计算出当前时间往前推10小时对应的Unix时间戳阈值,然后用这个阈值去过滤DataFrame里的时间戳列,保留大于等于该阈值的行就行。

步骤1:导入必要的Spark函数

首先得把需要用到的Spark SQL内置函数导进来:

import org.apache.spark.sql.functions.{unix_timestamp, current_timestamp, col}

步骤2:计算10小时前的Unix时间戳阈值

我们用current_timestamp()获取当前集群的系统时间,用minusHours(10)减去10小时,最后通过unix_timestamp()把这个时间转换成秒级的Unix时间戳:

val tenHoursAgoUnix = unix_timestamp(current_timestamp().minusHours(10))

步骤3:过滤DataFrame

假设你的DataFrame里存储Unix时间戳的列名叫event_timestamp,直接用filter方法筛选出时间戳大于等于阈值的行:

val filteredDF = yourOriginalDF.filter(col("event_timestamp") >= tenHoursAgoUnix)

特殊情况处理:如果时间戳是毫秒级的

如果你的Unix时间戳列是毫秒级的(比如13位数字),那需要把计算出来的阈值乘以1000,转换成毫秒级再比较:

// 毫秒级时间戳的筛选逻辑
val filteredDF = yourOriginalDF.filter(col("event_timestamp") >= tenHoursAgoUnix * 1000)

完整示例代码

这里给你一个可运行的完整示例,假设我们从CSV文件读取数据,并且时间戳列需要转成Long类型:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{unix_timestamp, current_timestamp, col}

object FilterRecentData {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("FilterRecent10HoursData")
      .master("local[*]") // 本地测试用,生产环境请移除
      .getOrCreate()

    // 读取原始数据,假设CSV里的时间戳列叫event_time
    val originalDF = spark.read
      .option("header", "true")
      .csv("path/to/your/data.csv")
      .withColumn("event_time", col("event_time").cast("long"))

    // 计算10小时前的Unix时间戳(秒级)
    val threshold = unix_timestamp(current_timestamp().minusHours(10))

    // 筛选最近10小时的数据
    val recentDF = originalDF.filter(col("event_time") >= threshold)

    // 查看结果
    recentDF.show()

    spark.stop()
  }
}

小提示:current_timestamp()获取的是Spark集群节点的系统时间,要确保集群时区和你的业务需求时区一致哦。

内容的提问来源于stack exchange,提问作者Markus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:09:25