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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 09:27:06