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

Spark中基于2天日期差范围聚合数据的实现方案

Spark 基于2天日期范围聚合数据(无循环/UDF实现)

问题场景

给定按日期升序排列的输入DataFrame,需要将后续日期与当前聚合组最早日期差在2天内的记录合并聚合,聚合后取组内最早日期作为输出日期,且不能使用循环或自定义UDF实现。

输入DataFrame

InputDateQuantity
2019-04-09 10:10:101
2019-04-11 10:10:105
2019-04-12 10:10:104
2019-04-15 10:10:105
2019-04-18 10:10:105
2019-04-20 10:10:105

期望输出DataFrame

OutputDateQuantity
2019-04-09 10:10:106
2019-04-12 10:10:104
2019-04-15 10:10:105
2019-04-18 10:10:1010

解决方案

利用Spark内置窗口函数和日期函数,通过动态标记分组起始日期的方式实现聚合,无需循环或UDF。

实现思路

  1. 将字符串日期转换为Timestamp类型,便于日期计算
  2. 使用窗口函数last()获取当前记录之前的分组起始日期
  3. 通过datediff()判断当前日期与分组起始日期的差值,超过2天则更新分组起始日期
  4. 按分组起始日期聚合,求和数量并重命名列

完整PySpark代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 初始化SparkSession
spark = SparkSession.builder.appName("date-range-agg").getOrCreate()

# 构建输入数据
data = [
    ("2019-04-09 10:10:10", 1),
    ("2019-04-11 10:10:10", 5),
    ("2019-04-12 10:10:10", 4),
    ("2019-04-15 10:10:10", 5),
    ("2019-04-18 10:10:10", 5),
    ("2019-04-20 10:10:10", 5)
]

df = spark.createDataFrame(data, ["InputDate", "Quantity"])

# 转换日期格式为Timestamp
df = df.withColumn("InputDate", F.to_timestamp("InputDate"))

# 定义按日期排序的窗口
date_window = Window.orderBy("InputDate")

# 动态生成分组起始日期
df = df.withColumn(
    "group_start",
    F.coalesce(
        # 如果当前日期与前分组起始日期差超过2天,更新为当前日期
        F.when(
            F.datediff(F.col("InputDate"), F.last("group_start", ignorenulls=True).over(date_window)) > 2,
            F.col("InputDate")
        ),
        # 第一条记录的分组起始日期为自身
        F.col("InputDate")
    )
)

# 按分组起始日期聚合求和
result_df = df.groupBy("group_start").agg(
    F.sum("Quantity").alias("Quantity")
).withColumnRenamed("group_start", "OutputDate")

# 展示结果(按日期排序)
result_df.orderBy("OutputDate").show(truncate=False)

代码说明

  • F.last("group_start", ignorenulls=True).over(date_window):获取当前记录之前所有记录的最新分组起始日期,忽略null值确保第一条记录能正确取值
  • F.datediff(end, start):计算两个日期的天数差,这里用来判断当前日期是否超出分组时间范围
  • 分组起始日期group_start会自动继承或更新,确保所有符合"2天内"规则的记录归为同一组
  • 最后按group_start聚合,直接得到每组的最早日期和数量总和

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 08:35:19