如何在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
相关产品推荐
相关产品推荐

