如何通过Spark按指定天数区间生成窗口内日期年龄映射列表?
Spark 滑动时间窗口收集历史映射实现方案
你原有代码的问题是未限制窗口的时间范围,默认会收集同性别分区内从第一条到当前行的所有记录,只要给窗口函数加上基于天数的滑动范围限制,即可实现自定义天数区间的需求,具体实现如下:
前置校验
提前将表中date字段转换为DateType类型,或确保日期字符串格式统一为yyyy-MM-dd。
方式1:Spark SQL 实现
将SQL中的${自定义天数}替换为你需要的数值(如15、30、45)即可直接运行:
SELECT Age, Gender, collect_list(map(date, Age)) OVER ( PARTITION BY Gender ORDER BY CAST(unix_timestamp(date, 'yyyy-MM-dd') / 86400 AS INT) RANGE BETWEEN ${自定义天数} PRECEDING AND CURRENT ROW ) AS date_age_map FROM master
逻辑说明:将日期转换为1970年以来的天数整数,通过RANGE BETWEEN限制只收集当前行日期往前N天内的同性别记录。
方式2:Spark DataFrame API 实现(Scala版)
可通过变量灵活配置天数参数:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 自定义天数,可按需修改 val intervalDays: Int = 15 // 定义窗口规范 val timeWindow = Window .partitionBy("Gender") .orderBy(cast(unix_timestamp(col("date"), "yyyy-MM-dd") / 86400 as "int")) .rangeBetween(-intervalDays, Window.currentRow) // 计算结果 val result = spark.table("master") .withColumn("date_age_map", collect_list(map(col("date"), col("Age"))).over(timeWindow)) .select("Age", "Gender", "date_age_map")
可选:格式化为指定字符串样式
如果需要输出你示例中日期->[age数组]的字符串格式,可加简单UDF处理:
val formatMap = udf((maps: Seq[Map[String, Seq[Int]]]) => { maps.map(m => s"${m.head._1}->${m.head._2}") }) val finalResult = result.withColumn("date_age_map", formatMap(col("date_age_map")))
该实现完全匹配你给出的15天区间的预期输出结果,调整intervalDays参数即可适配不同时间区间需求。
内容的提问来源于stack exchange,提问作者punam
相关产品推荐
相关产品推荐

