PySpark中关联两个时间序列DataFrame匹配右表最近有效条目
Spark DataFrame 关联实现方案
核心思路
- 先对df2按
id和date字段去重,同一个id同一个日期的重复日志仅保留1条,降低后续计算开销 - 左关联df1与去重后的df2,关联条件为
id等值匹配,且df2的日期小于等于df1当前行的日期 - 按df1的
id和date分组,取df2日期的最大值,即为同id下小于等于df1当前日期的最近日期
代码实现
通用实现(兼容所有Spark版本)
from pyspark.sql import functions as F # 第一步:统一将date字段转为日期类型,避免字符串比较出错 df1 = df1.withColumn("date", F.to_date(F.col("date"))) df2 = df2.withColumn("date", F.to_date(F.col("date"))) # 第二步:对df2按id、date去重 df2_distinct = df2.dropDuplicates(["id", "date"]) # 第三步:关联分组计算 result = df1.alias("a") \ .join( df2_distinct.alias("b"), (F.col("a.id") == F.col("b.id")) & (F.col("b.date") <= F.col("a.date")), how="left" ) \ .groupBy("a.id", "a.date") \ .agg(F.max("b.date").alias("date(df2)")) \ .select("a.date", "a.id", "date(df2)") \ .orderBy(F.col("a.id"), F.col("a.date").desc())
Spark SQL 写法
如果习惯用SQL语法,可按如下方式实现:
# 注册临时视图 df1.createOrReplaceTempView("df1") df2_distinct.createOrReplaceTempView("df2_distinct") # 执行SQL查询 result = spark.sql(""" SELECT a.date, a.id, MAX(b.date) AS `date(df2)` FROM df1 a LEFT JOIN df2_distinct b ON a.id = b.id AND b.date <= a.date GROUP BY a.id, a.date ORDER BY a.id, a.date DESC """)
Spark 3.0+ 高性能Asof Join实现
针对大数据量场景,可使用Spark 3.0推出的Asof Join优化,性能比普通关联提升明显:
result = df1.join( df2_distinct.hint("ASOF"), (df1.id == df2_distinct.id) & (df2_distinct.date <= df1.date), how="left" ).select(df1.date, df1.id, df2_distinct.date.alias("date(df2)")) \ .orderBy(df1.id, df1.date.desc())
内容的提问来源于stack exchange,提问作者Into Numbers
相关产品推荐
相关产品推荐

