使用PySpark关联两表提取事件日期前对应用户最新体重
PySpark实现方案
前置处理
首先需要将两个表的日期字段转换为日期类型,避免字符串比较出现异常:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_date, row_number from pyspark.sql.window import Window # 转换日期字段类型 events_df = Events.withColumn("event_date", to_date(col("Date"), "yyyy-MM-dd")) user_df = User.withColumn("weight_date", to_date(col("Date"), "yyyy-MM-dd"))
方法1:窗口函数法(全版本兼容)
该方法兼容性最强,支持所有PySpark版本,逻辑清晰易读:
# 左关联两个表,关联条件为同用户且体重记录日期早于事件发生日期 joined_df = events_df.join( user_df, (events_df.User_id == user_df.User_id) & (user_df.weight_date < events_df.event_date), how="left" ) # 定义窗口:按事件ID分组,按体重记录日期倒序排序,最新的记录排在第一位 window_spec = Window.partitionBy("Event_id").orderBy(col("weight_date").desc()) # 取每个事件对应的最新体重记录,无匹配的体重字段自动为null result_df = joined_df.withColumn("rn", row_number().over(window_spec)) \ .filter(col("rn") == 1) \ .select("Event_id", "Date", "User_id", "Weight-kg")
方法2:Asof Join法(PySpark 3.0+ 高性能方案)
如果你的PySpark版本在3.0以上,推荐用专门用于最近值匹配的Asof Join,性能远高于先全量关联再开窗的方案,适合大数据量场景:
# 两个表先按关联键和排序字段排序 events_df = events_df.orderBy("User_id", "event_date") user_df = user_df.orderBy("User_id", "weight_date") # 执行asof join直接匹配最近记录 result_df = events_df.join( user_df.hint("asof"), on=["User_id"], condition=(user_df.weight_date < events_df.event_date), how="left" ).select("Event_id", "Date", "User_id", "Weight-kg")
两种方法输出结果均和你给出的预期完全一致。
内容的提问来源于stack exchange,提问作者danimille
相关产品推荐
相关产品推荐

