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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 16:36:05