判断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", ???)
因日期数据可能超千条,递归方案性能不足,需高效实现判断逻辑。
解决方案
思路
核心通过日期转换与集合运算实现高效判断,避免递归:
- 将日期转换为年月格式(
yyyy-MM),去重后得到每个用户的行为月份集合; - 以用户自身的最新行为日期为基准,计算近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
相关产品推荐
相关产品推荐

