Databricks(PySpark)中两DataFrame按ID和日期条件关联映射求助
在Databricks(PySpark)中实现基于条件的DataFrame映射(最近日期匹配)
需求概述
当ID匹配时,将DF2中Date1小于等于DF1.Date且最接近该Date的行的相关列映射到DF1;若存在多个符合条件的行,取最新的Date1对应的行;若无符合条件的行,则填充null。要求避免Python循环,适配大数据量处理。
实现步骤与代码
1. 准备示例数据(可替换为你的实际数据)
from pyspark.sql import functions as F, Window # 创建DF1示例数据 data1 = [ ("21/06/2022", "ABC", "XZ18610", ""), ("22/05/2022", "ABC", "XZ18610", ""), ("22/04/2022", "ABC", "XZ18610", ""), ("05/05/2022", "DEF", "XZ25277", ""), ("04/02/2022", "DEF", "XZ25277", ""), ("28/06/2022", "GHI", "XZU6S19", ""), ("18/07/2022", "JKL", "XZ54866", ""), ("27/07/2022", "MNO", "XZ82434", ""), ("20/06/2022", "PQR", "XZ78433", "") ] df1 = spark.createDataFrame(data1, ["Date", "name", "ID", "其他列"]) # 创建DF2示例数据 data2 = [ ("30/05/2022", "XZ18610", "B", ""), ("21/06/2021", "XZ18610", "A", ""), ("05/01/2021", "XZ25277", "B", ""), ("28/07/2022", "XZU6S19", "E", ""), ("18/05/2022", "XZ54866", "D", ""), ("27/07/2022", "XZ82434", "F", ""), ("20/06/2022", "XZ78433", "I", "") ] df2 = spark.createDataFrame(data2, ["Date1", "ID1", "Value", "其他列1"])
2. 预处理:转换日期格式与列名对齐
将字符串类型的日期转为Spark日期类型(确保日期比较逻辑正确),并重命名DF2的ID列与DF1对齐:
# 转换DF1的Date列为日期类型 df1 = df1.withColumn("Date", F.to_date(F.col("Date"), "dd/MM/yyyy")) # 重命名DF2的ID列,转换Date1列为日期类型 df2 = df2.withColumnRenamed("ID1", "ID") \ .withColumn("Date1", F.to_date(F.col("Date1"), "dd/MM/yyyy"))
3. 左连接+窗口函数实现最近日期匹配
通过左连接保留DF1所有行,再用窗口函数筛选每个DF1行对应的最近符合条件的DF2记录:
# 左连接两个DataFrame,过滤出DF2中Date1 <= DF1.Date的记录(含无匹配的null行) joined_df = df1.join(df2, on="ID", how="left") \ .filter((F.col("Date1") <= F.col("Date")) | F.col("Date1").isNull()) # 定义窗口:按DF1的唯一标识(Date/name/ID/其他列)分组,按Date1降序排序 window_spec = Window.partitionBy(df1["Date"], df1["name"], df1["ID"], df1["其他列"]) \ .orderBy(F.desc("Date1")) # 添加行号,取每个分组的第一行(即最近的符合条件的DF2记录) result_df = joined_df.withColumn("row_num", F.row_number().over(window_spec)) \ .filter(F.col("row_num") == 1) \ .drop("row_num")
4. 可选:将日期转回原始字符串格式
如果需要和输入格式一致,可将日期列转回字符串:
result_df = result_df.withColumn("Date", F.date_format(F.col("Date"), "dd/MM/yyyy")) \ .withColumn("Date1", F.when(F.col("Date1").isNotNull(), F.date_format(F.col("Date1"), "dd/MM/yyyy")))
5. 查看结果
result_df.show()
关键说明
- 日期类型转换:必须将字符串日期转为Spark的
date类型,否则字符串比较会出现逻辑错误(如"21/06/2022"和"30/05/2022"的字符串比较结果不符合日期逻辑)。 - 窗口函数优势:基于Spark分布式计算,避免Python循环,可高效处理海量数据。
- 左连接保留全量行:确保DF1中无匹配条件的行不会被丢弃,最终填充null。
内容的提问来源于stack exchange,提问作者ASD
相关产品推荐
相关产品推荐

