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

Spark中group by分组后如何按指定日期区间过滤分组内数据

Spark分组后按日期规则过滤方案

实现思路

通过按accountid分组后对日期排序,校验每行的前置记录日期是否命中过滤区间,最终保留符合要求的行。提供SQL和DataFrame API两种实现方式,可按需选择。

具体实现

方式1:Spark SQL实现

直接对原始表执行如下SQL即可:

WITH temp AS (
    SELECT 
        *,
        -- 取同分组内排序后上一行的起始日期
        LAG(startdateacceleration, 1) OVER (PARTITION BY accountid ORDER BY startdateacceleration) AS prev_start_date,
        -- 计算当前行起始日期对应的上月首日
        DATE_TRUNC('month', ADD_MONTHS(startdateacceleration, -1)) AS prev_month_first_day
    FROM 你的表名
)
SELECT accountid, startdateacceleration, enddateacceleration
FROM temp
-- 保留没有前置行的分组首行,以及前置行日期早于上月首日的行
WHERE prev_start_date IS NULL OR prev_start_date < prev_month_first_day

方式2:Scala DataFrame API实现

如果是代码开发场景,可以用如下实现:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

// 定义分组排序窗口
val win = Window.partitionBy("accountid").orderBy("startdateacceleration")

val resultDF = 原始数据集DF
  .withColumn("prev_start_date", lag("startdateacceleration", 1).over(win))
  .withColumn("prev_month_first_day", date_trunc("month", add_months(col("startdateacceleration"), -1)))
  .filter(col("prev_start_date").isNull || col("prev_start_date") < col("prev_month_first_day"))
  .select("accountid", "startdateacceleration", "enddateacceleration")

特殊场景适配

如果需要基于已经保留的上一行做判断(而非原始排序的上一行),可以用mapGroups算子实现状态遍历:

case class AccRecord(accountid: String, startdateacceleration: java.sql.Timestamp, enddateacceleration: java.sql.Timestamp)

val resultDF = 原始数据集DF
  .as[AccRecord]
  .groupByKey(_.accountid)
  .flatMapGroups { (_, records) =>
    val sortedRecords = records.toList.sortBy(_.startdateacceleration)
    sortedRecords.foldLeft(List.empty[AccRecord]) { (retained, cur) =>
      if (retained.isEmpty) {
        cur :: retained
      } else {
        val lastRetainedStart = retained.head.startdateacceleration
        val curPrevMonthFirst = java.sql.Timestamp.valueOf(
          java.time.LocalDateTime.ofInstant(cur.startdateacceleration.toInstant, java.time.ZoneId.systemDefault())
            .minusMonths(1).withDayOfMonth(1).toLocalDate.atStartOfDay()
        )
        if (lastRetainedStart.before(curPrevMonthFirst)) cur :: retained else retained
      }
    }.reverse
  }.toDF()

验证说明

对你给出的示例数据执行后,输出结果和你期望的完全一致:

accountidstartdateaccelerationenddateacceleration
0011t00000MYFRKAA52021-05-25 00:00:00.000000NULL
0011t00000MYFRKAA52021-09-01 00:00:00.0000002022-05-26 00:00:00.000000

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 17:27:03