如何使用PySpark按日期及时间区间关联合并两张数据表
PySpark 按日期区间关联两个表的实现方案
你要的需求可以直接通过带多关联条件的左连接实现,不需要引入额外依赖,适配你给出的样例场景代码如下:
首先确认已经导入必要依赖:
import datetime from pyspark.sql import SparkSession import pyspark.sql.functions as F spark = SparkSession.builder.appName("date_range_join").getOrCreate()
你给出的样例数据初始化代码无需修改,直接执行关联逻辑即可:
# 多条件左连接:id相等 + 事件日期落在属性表的起止区间内 df_result = df1.join( df2, (df1.id == df2.id) & (df1.event_date >= df2.start_period) & (df1.event_date <= df2.end_period), how="left" ).select(df1["*"], df2.info1, df2.info2) # 查看结果 df_result.show()
执行后输出的结果和你给出的预期完全一致:
+---+-----+-------------------+------+------+ | id|event| event_date| info1| info2| +---+-----+-------------------+------+------+ | 1| a|2021-01-01 00:00:00| Xxz45| XX013| | 1| b|2021-01-05 00:00:00| Xxz45| XX013| | 1| c|2021-01-24 00:00:00| Xbbd| XX015| | 2| d|2021-01-10 00:00:00| null| null| | 2| e|2021-01-15 00:00:00| null| null| +---+-----+-------------------+------+------+
扩展场景适配
- 如果存在单个事件日期匹配到多个区间的情况,可通过窗口函数取指定优先级的匹配结果,示例如下(取最新开始时间的区间):
from pyspark.sql.window import Window window_spec = Window.partitionBy("id", "event_date").orderBy(F.desc("start_period")) df_result = df_result.withColumn("rn", F.row_number().over(window_spec))\ .filter(F.col("rn") == 1)\ .drop("rn")
- 如果df2数据量较小,可使用广播优化关联性能,将join逻辑中的
df2替换为F.broadcast(df2)即可,避免shuffle开销。
内容的提问来源于stack exchange,提问作者Alex Germain
相关产品推荐
相关产品推荐

