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()
验证说明
对你给出的示例数据执行后,输出结果和你期望的完全一致:
| accountid | startdateacceleration | enddateacceleration |
|---|---|---|
| 0011t00000MYFRKAA5 | 2021-05-25 00:00:00.000000 | NULL |
| 0011t00000MYFRKAA5 | 2021-09-01 00:00:00.000000 | 2022-05-26 00:00:00.000000 |
内容的提问来源于stack exchange,提问作者aarigoni
相关产品推荐
相关产品推荐

