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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 08:06:03