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

判断Spark DataFrame中用户近3个月行为日期是否连续

问题:判断用户近3个月是否持续有行为

原始数据

现有包含用户ID与行为日期的DataFrame:

+----+----------+
|  id|      date|
+----+----------+
|   1|2022-09-01|
|   1|2022-10-01|
|   1|2022-11-01|
|   2|2022-07-01|
|   2|2022-10-01|
|   2|2022-11-01|
|   3|2022-09-01|
|   3|2022-10-01|
|   3|2022-11-01|
+----+----------+

需求

判断每个用户是否在近3个月内持续有行为(示例中用户2缺失9月数据,不符合条件),期望输出:

+----+------------------------------------+------+
|  id|                               dates|3month|
+----+------------------------------------+------+
|   1|[2022-09-01, 2022-10-01, 2022-11-01]|  true|
|   2|[2022-07-01, 2022-10-01, 2022-11-01]| false|
|   3|[2022-09-01, 2022-10-01, 2022-11-01]|  true|
+----+------------------------------------+------+

已完成分组收集日期数组的代码:

data.groupBy(col("id")).agg(collect_list("date") as "dates").withColumn("3month", ???)

因日期数据可能超千条,递归方案性能不足,需高效实现判断逻辑。


解决方案

思路

核心通过日期转换与集合运算实现高效判断,避免递归:

  1. 将日期转换为年月格式(yyyy-MM),去重后得到每个用户的行为月份集合;
  2. 以用户自身的最新行为日期为基准,计算近3个月的目标月份(当月及前两个月);
  3. 检查用户的行为月份集合是否包含所有目标月份。

完整代码实现(Scala)

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.DateType

// 1. 预处理:转换日期类型、提取年月,收集用户的日期列表、最新日期、去重行为月份
val processedData = data
  .withColumn("date", col("date").cast(DateType))
  .withColumn("year_month", date_format(col("date"), "yyyy-MM"))
  .groupBy("id")
  .agg(
    collect_list("date").alias("dates"),
    max("date").alias("latest_date"),
    collect_set("year_month").alias("active_months")
  )

// 2. 计算需要检查的近3个目标月份
val targetMonthsData = processedData
  .withColumn("target_month1", date_format(col("latest_date"), "yyyy-MM"))
  .withColumn("target_month2", date_format(add_months(col("latest_date"), -1), "yyyy-MM"))
  .withColumn("target_month3", date_format(add_months(col("latest_date"), -2), "yyyy-MM"))
  .withColumn("target_months", array(col("target_month1"), col("target_month2"), col("target_month3")))

// 3. 验证行为月份是否覆盖所有目标月份,生成结果字段
val result = targetMonthsData
  .withColumn("3month", array_except(col("target_months"), col("active_months")).isEmpty)
  .select("id", "dates", "3month")

result.show(false)

代码说明

  • collect_set("year_month"):去重收集用户的行为月份,避免同一月多条记录干扰判断;
  • add_months:Spark内置日期函数,高效计算偏移月份,无需手动处理年月转换逻辑;
  • array_except:计算目标月份与行为月份的差集,若差集为空则说明所有目标月份都有行为,返回true,否则false。

该方案完全基于Spark内置函数实现,分布式计算性能优异,即使面对超千条日期数据也能高效处理。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 22:45:37