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

PySpark实现DataFrame匹配最近前向Timestamp并新增对应值列

PySpark 前向最近时间匹配实现方案

你需要实现的是时间序列的前向最近匹配需求,适配PySpark的两种实现方式如下:


方案1:原生ASOF JOIN(推荐,PySpark 3.1+支持)

该方案为Spark官方针对有序时间匹配场景的原生实现,性能远高于自定义逻辑,适配绝大多数生产环境需求:

from pyspark.sql import SparkSession
from pyspark.sql.functions import to_timestamp, col

# 1. 拆分两个独立时间序列数据集(从你已有的合并DF中提取、去重)
df1 = merged_df.select("TimeStamp1", "Value1").dropDuplicates()
df2 = merged_df.select("TimeStamp2", "Value2").dropDuplicates()

# 2. 确保时间字段为Timestamp类型,如原有为字符串可执行转换
df1 = df1.withColumn("TimeStamp1", to_timestamp(col("TimeStamp1")))
df2 = df2.withColumn("TimeStamp2", to_timestamp(col("TimeStamp2")))

# 3. 执行前向ASOF JOIN,direction="forward"即匹配>=当前TimeStamp1的最小TimeStamp2
result_df = df1.join(
    df2,
    on=df1.TimeStamp1 >= df2.TimeStamp2,
    how="asof",
    direction="forward"
).select("TimeStamp1", "Value1", col("Value2").alias("Value 2"))

方案2:窗口函数实现(兼容PySpark 3.1以下版本)

如果你的Spark版本不支持ASOF JOIN,可以通过窗口函数实现相同逻辑:

from pyspark.sql import Window
from pyspark.sql.functions import row_number, col

# 1. 先过滤出所有TimeStamp2大于等于当前TimeStamp1的候选行
filtered_df = merged_df.filter(col("TimeStamp2") >= col("TimeStamp1"))

# 2. 按TimeStamp1分组,取每组TimeStamp2最小的行即为匹配结果
window_spec = Window.partitionBy("TimeStamp1", "Value1").orderBy("TimeStamp2")
result_df = filtered_df.withColumn("rn", row_number().over(window_spec)) \
    .filter(col("rn") == 1) \
    .select("TimeStamp1", "Value1", col("Value2").alias("Value 2"))

两种方案输出结果均与你给出的预期输出完全一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 02:24:03