Spark SQL左连接保留左表全量且基于右表聚合取对应列的问题咨询
通用SQL左连接按右表聚合条件筛选实现方案(Spark SQL为例)
核心需求
- 保留左表t1的全量行数据
- 基于右表t2的聚合结果选取对应target字段
测试表结构与数据
左表t1
+---+--------+ |pk1|constant| +---+--------+ | a|constant| | b|constant| | c|constant| | d|constant| +---+--------+
右表t2
+---+---------+------+ |fk1|condition|target| +---+---------+------+ | a| 1|check1| | a| 2|check2| | b| 1|check1| | b| 2|check2| +---+---------+------+
已尝试的错误方案
方案1
spark.sql(""" select pk1, constant, target from t1 left join t2 on pk1 = fk1 group by pk1, constant, target having min(condition) """).show
输出结果:
+---+--------+------+ |pk1|constant|target| +---+--------+------+ | b|constant|check1| | a|constant|check2| | a|constant|check1| | b|constant|check2| +---+--------+------+
问题说明:过滤逻辑导致t1中pk1为'c'、'd'的行丢失,左连接退化为内连接。
方案2
spark.sql(""" select pk1, constant, min(condition), target from t1 left join t2 on pk1 = fk1 group by pk1, constant, target """).show
输出结果:
+---+--------+--------------+------+ |pk1|constant|min(condition)|target| +---+--------+--------------+------+ | a|constant| 1|check1| | a|constant| 2|check2| | b|constant| 2|check2| | b|constant| 1|check1| | c|constant| null| null| | d|constant| null| null| +---+--------+--------------+------+
问题说明:min聚合未起到筛选作用,同一个pk1对应多行t2数据全部被返回。
期望输出结果
+---+--------+--------------+------+ |pk1|constant|min(condition)|target| +---+--------+--------------+------+ | a|constant| 1|check1| | b|constant| 1|check1| | c|constant| null| null| | d|constant| null| null| +---+--------+--------------+------+
注:min(condition)列可根据需要自行删除。
单查询最优实现方案
方案1:窗口函数法(兼容性最强,适配绝大多数SQL方言)
对t2按外键分组取condition最小的行,再左连接t1即可:
spark.sql(""" select t1.pk1, t1.constant, t2_filter.min_condition, t2_filter.target from t1 left join ( select fk1, condition as min_condition, target, row_number() over(partition by fk1 order by condition asc) as rn from t2 ) t2_filter on t1.pk1 = t2_filter.fk1 and t2_filter.rn = 1 """).show
方案2:右表预聚合再关联(性能更优,适合大数据量场景)
先聚合t2拿到每个fk1对应的最小condition,再回表关联t2拿到对应target,最后左连t1:
spark.sql(""" select t1.pk1, t1.constant, t2_agg.min_condition, t2.target from t1 left join ( select fk1, min(condition) as min_condition from t2 group by fk1 ) t2_agg on t1.pk1 = t2_agg.fk1 left join t2 on t2_agg.fk1 = t2.fk1 and t2_agg.min_condition = t2.condition """).show
附录:测试表构建代码
val columns1 = Seq("pk1", "constant") val columns2 = Seq("fk1","condition","target") val data1 = Seq( ("a","constant"), ("b","constant"), ("c","constant"), ("d","constant") ) val data2 = Seq( ("a",1,"check1"), ("a",2,"check2"), ("b",1,"check1"), ("b",2,"check2") ) val t1 = spark.createDataFrame(data1).toDF(columns1:_*) val t2 = spark.createDataFrame(data2).toDF(columns2:_*)
内容的提问来源于stack exchange,提问作者Blaisem
相关产品推荐
相关产品推荐

