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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 12:06:06