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

Spark中"partial"窗口函数:限定范围计算营收的实现问题

解决方案:Spark中实现带时间窗口的聚合并过滤目标数据

刚好之前处理过类似的广告归因聚合场景,这个需求的核心是先用充足的历史数据覆盖窗口计算逻辑,再精准过滤出最终需要保留的结果集,下面是具体的实现思路和代码:

核心思路拆解

  • 加载7天数据:确保最早的目标数据(最近3天的第一天)能获取到完整的4天窗口数据(比如最近3天是day5-day7,day5的窗口需要day2-day5的数据,所以必须加载day2-day7,为了保险通常加载连续7天)
  • 定义时间范围窗口:按campaign分区,基于impressionTime的时间范围(而非行号)计算近4天的revenue总和,避免因某天无数据导致窗口计算错误
  • 过滤最近3天数据:计算7天数据中的最大日期,只保留该日期往前推2天(含当天)的结果

具体实现代码

Scala版本

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.DateType
import org.apache.spark.sql.expressions.Window
import java.sql.Date
import java.time.LocalDate

// 1. 加载7天数据并标准化时间字段
val rawData = spark.read
  .format("parquet") // 根据你的数据源替换,比如csv、delta等
  .load("/path/to/7days/data")
  // 将时间字段转换为日期类型(如果是timestamp可以转成date方便后续计算)
  .withColumn("impressionDate", to_date(col("impressionTime")).cast(DateType))

// 2. 计算过滤阈值:取7天数据中的最新日期,往前推2天即为保留数据的起始日期
val maxDate = rawData.select(max("impressionDate")).first().getDate(0)
val cutoffDate = Date.valueOf(maxDate.toLocalDate.minusDays(2))

// 3. 定义窗口规范:按campaign分区,按日期的时间戳排序,范围是当前日期往前3天(共4天)
// 注意:将日期转成long型(秒数)才能使用rangeBetween的时间范围计算
val campaignWindow = Window
  .partitionBy("campaign")
  .orderBy(col("impressionDate").cast("long"))
  .rangeBetween(-3 * 86400, 0) // 3天的总秒数:3*24*3600=86400*3

// 4. 计算近4天revenue总和,然后过滤出最近3天的数据
val finalResult = rawData
  .withColumn("last_4d_revenue_sum", sum("revenue").over(campaignWindow))
  .filter(col("impressionDate") >= cutoffDate)
  .drop("impressionDate") // 移除临时日期字段
  .select("eventId", "impressionTime", "campaign", "revenue", "last_4d_revenue_sum")

// 查看结果或写入存储
finalResult.show()
// finalResult.write.format("parquet").save("/path/to/result")

Python版本

from pyspark.sql import functions as F
from pyspark.sql.types import DateType
from pyspark.sql.window import Window
from datetime import timedelta

# 1. 加载7天数据并处理时间字段
raw_data = spark.read \
    .format("parquet") \
    .load("/path/to/7days/data") \
    .withColumn("impressionDate", F.to_date(F.col("impressionTime")).cast(DateType))

# 2. 计算过滤起始日期
max_date = raw_data.select(F.max("impressionDate")).first()[0]
cutoff_date = max_date - timedelta(days=2)

# 3. 定义时间范围窗口
campaign_window = Window \
    .partitionBy("campaign") \
    .orderBy(F.col("impressionDate").cast("long")) \
    .rangeBetween(-3 * 86400, 0)

# 4. 计算聚合并过滤结果
final_result = raw_data \
    .withColumn("last_4d_revenue_sum", F.sum("revenue").over(campaign_window)) \
    .filter(F.col("impressionDate") >= cutoff_date) \
    .drop("impressionDate") \
    .select("eventId", "impressionTime", "campaign", "revenue", "last_4d_revenue_sum")

# 输出或保存结果
final_result.show()
# final_result.write.format("parquet").save("/path/to/result")

关键注意事项

  • 为什么用rangeBetween而非rowsBetween?:rowsBetween是按行的顺序计算窗口,若某天无数据会导致窗口范围错误;rangeBetween基于时间的数值(秒数)计算,能精准覆盖近4天的时间范围,不受数据缺失影响
  • 性能优化:如果数据是按日期分区的,加载时直接过滤出7天的分区数据(比如filter(impressionDate >= '2024-01-01' and impressionDate <= '2024-01-07')),避免加载全量数据
  • 时间精度处理:如果impressionTime是毫秒级timestamp,需将rangeBetween的数值调整为毫秒(-3*86400*1000)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:43:14