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

如何基于时间范围条件在PySpark中关联两张数据表?

在PySpark中实现基于时间区间的表关联

要实现每个id匹配满足startTime <= Time <= endTime的记录,直接使用PySpark的**不等值连接(Non-Equi Join)**即可,这种方式天然支持范围条件匹配,无需额外复杂逻辑。

步骤1:创建示例数据

先构造测试用的DataFrame,模拟你的两张表结构:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, TimestampType

# 初始化Spark会话
spark = SparkSession.builder.appName("TimeRangeJoin").getOrCreate()

# 表1:包含id、startTime、endTime
schema_table1 = StructType([
    StructField("id", StringType(), nullable=True),
    StructField("startTime", TimestampType(), nullable=True),
    StructField("endTime", TimestampType(), nullable=True)
])
data_table1 = [
    ("A", "2023-01-01 00:00:00", "2023-01-01 02:00:00"),
    ("B", "2023-01-01 01:00:00", "2023-01-01 03:00:00"),
    ("C", "2023-01-01 04:00:00", "2023-01-01 05:00:00")
]
df_table1 = spark.createDataFrame(data_table1, schema_table1)

# 表2:包含Time、value
schema_table2 = StructType([
    StructField("Time", TimestampType(), nullable=True),
    StructField("value", StringType(), nullable=True)
])
data_table2 = [
    ("2023-01-01 00:30:00", "v1"),
    ("2023-01-01 01:30:00", "v2"),
    ("2023-01-01 02:30:00", "v3"),
    ("2023-01-01 04:30:00", "v4")
]
df_table2 = spark.createDataFrame(data_table2, schema_table2)

步骤2:执行时间区间关联

使用join方法,通过逻辑与(&)连接两个时间范围条件,即可完成匹配:

# 执行内连接,保留双方满足条件的记录;若需保留表1所有id,可改用how="left"
result_df = df_table1.join(
    df_table2,
    (df_table1.startTime <= df_table2.Time) & (df_table2.Time <= df_table1.endTime),
    how="inner"
)

# 查看结果
result_df.show(truncate=False)

输出结果

上述代码的输出如下,可见重叠时间区间的记录会被正确匹配:

+---+-------------------+-------------------+-------------------+-----+
|id |startTime          |endTime            |Time               |value|
+---+-------------------+-------------------+-------------------+-----+
|A  |2023-01-01 00:00:00|2023-01-01 02:00:00|2023-01-01 00:30:00|v1   |
|A  |2023-01-01 00:00:00|2023-01-01 02:00:00|2023-01-01 01:30:00|v2   |
|B  |2023-01-01 01:00:00|2023-01-01 03:00:00|2023-01-01 01:30:00|v2   |
|B  |2023-01-01 01:00:00|2023-01-01 03:00:00|2023-01-01 02:30:00|v3   |
|C  |2023-01-01 04:00:00|2023-01-01 05:00:00|2023-01-01 04:30:00|v4   |
+---+-------------------+-------------------+-------------------+-----+

性能优化建议

如果数据量较大,可通过以下方式提升性能:

  • 广播小表:若表2数据量远小于表1,使用broadcast函数广播表2,减少shuffle开销:
    from pyspark.sql.functions import broadcast
    
    result_df = df_table1.join(
        broadcast(df_table2),
        (df_table1.startTime <= df_table2.Time) & (df_table2.Time <= df_table1.endTime),
        how="inner"
    )
    
  • 分区/分桶:针对时间字段(startTime、Time)进行分区或分桶,缩小join时的数据扫描范围。

内容的提问来源于stack exchange,提问作者user6386155

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 00:26:37