PySpark如何基于唯一字段与日期范围匹配列并标记设备场景
解决方案
首先需要将两个DataFrame关联以获取设备转移日期,再结合窗口函数完成三类场景的判定:
步骤1:关联DataFrame拿到transferdate
通过userId和deviceID将df与df2关联,确保每条登录记录都能匹配到对应的设备转移日期:
from pyspark.sql import functions as f, Window # 关联两个DataFrame,左连接避免丢失df中的登录记录 df_joined = df.join(df2, on=["userId", "deviceID"], how="left")
注:如果df中存在df2无匹配的设备记录,transferdate会为null,这类设备直接判定为不属于当前用户(无合法转移记录)。
步骤2:标记设备归属状态
新增列标记当前登录的设备是否属于用户:
df_joined = df_joined.withColumn( "is_owned", # 转移日期晚于登录日/无转移记录,均判定为非自有设备 f.when( (f.col("transferdate").isNull()) | (f.col("transferdate") > f.col("Clean_date")), False ).otherwise(True) )
步骤3:计算分组维度的关键指标
用窗口函数计算两个维度的核心指标:
- 同一用户+同一登录日的设备数量、是否存在非自有设备
- 同一用户的总登录设备数量、是否存在非自有设备
# 窗口1:按用户+登录日期分组 w_daily = Window.partitionBy("userId", "Clean_date") # 窗口2:按用户分组 w_user = Window.partitionBy("userId") df_with_metrics = df_joined.withColumn( "daily_device_count", f.size(f.collect_set("deviceID").over(w_daily)) ).withColumn( "has_unowned_daily", f.max(f.col("is_owned") == False).over(w_daily) ).withColumn( "total_device_count", f.size(f.collect_set("deviceID").over(w_user)) ).withColumn( "has_unowned_total", f.max(f.col("is_owned") == False).over(w_user) )
步骤4:生成identifier列
按照预设规则判定场景:
df_final = df_with_metrics.withColumn( "identifier", f.when( (f.col("daily_device_count") > 1) & f.col("has_unowned_daily"), "P1" ).when( (f.col("total_device_count") > 1) & f.col("has_unowned_total"), "P2" ).when( (f.col("total_device_count") > 1) & ~f.col("has_unowned_total"), "NA" ).otherwise("NA") # 单设备登录场景也归为NA )
逻辑说明
- P1:当日登录多设备,且至少一台设备不属于用户(转移日期晚于登录日/无转移记录)
- P2:跨日登录多设备,且至少一台设备不属于用户(但当日登录设备数不满足P1条件)
- NA:所有登录设备均属于用户(无论同日或跨日多设备),或仅登录单设备
内容的提问来源于stack exchange,提问作者Akash Pujara
相关产品推荐
相关产品推荐

