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
相关产品推荐
相关产品推荐

