如何基于日期匹配关联主键不同的表并避免数据膨胀?
解决方案
Spark-SQL 实现
利用 LATERAL VIEW OUTER(Spark 2.4及以上版本支持)可以让子查询直接引用表Y的字段,同时通过排序+限制行数确保每个Y的行仅关联一行符合条件的X数据:
SELECT y.*, x.FieldINeed, x.WerktijdID FROM Y y LEFT JOIN LATERAL VIEW OUTER ( SELECT x.FieldINeed, x.WerktijdID FROM X x WHERE x.OmgevingId = y.OmgevingId AND x.AdministratieKantoorID = y.AdministratieKantoorID AND x.WerkgeverID = y.WerkgeverID AND x.DienstverbandID = y.DienstverbandID AND x.StartDate <= y.Period AND (x.EndDate >= y.Period OR x.EndDate IS NULL) -- 若同一组有多个符合条件的行,可调整排序规则(比如按StartDate降序),再取第一行 ORDER BY x.WerktijdID DESC LIMIT 1 ) x ON true
PySpark DataFrame API 实现
通过连接+筛选+窗口函数去重的步骤,实现关联并避免数据膨胀:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, row_number spark = SparkSession.builder.appName("Y_X_Join").getOrCreate() # 假设已加载表Y和X为DataFrame:df_y、df_x # 1. 按共同键左连接,筛选符合日期条件的行 temp_joined = df_y.join( df_x, on=[ "OmgevingId", "AdministratieKantoorID", "WerkgeverID", "DienstverbandID" ], how="left" ).filter( (col("StartDate") <= col("Period")) & (col("EndDate") >= col("Period") | col("EndDate").isNull()) ) # 2. 定义窗口:按Y的复合主键分组,对符合条件的X行排序 window_spec = Window.partitionBy( "OmgevingId", "AdministratieKantoorID", "WerkgeverID", "DienstverbandID", "VerloningsNummer" ).orderBy(col("WerktijdID").desc()) # 3. 取每组第一行,去掉冗余列 filtered_x = temp_joined.withColumn("row_num", row_number().over(window_spec)) \ .filter(col("row_num") == 1) \ .drop("row_num", "StartDate", "EndDate") # 4. 与原Y表左连接,保留所有Y行 final_result = df_y.join( filtered_x, on=["OmgevingId", "AdministratieKantoorID", "WerkgeverID", "DienstverbandID", "VerloningsNummer"], how="left" ) # 查看结果 final_result.show()
说明
- 若同一组存在多个符合日期条件的X行,可根据业务需求调整排序规则(比如按
StartDate降序、EndDate降序等),确保取到最符合需求的那一行。 - 两种方案均避免了直接m:n连接导致的数据膨胀,实现了Y到X的m:1关联。
内容的提问来源于stack exchange,提问作者Psychotechnopath
相关产品推荐
相关产品推荐

