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

SQL左连接查询转Scala版Spark DataFrame实现问题咨询

Spark左连接统计右表字段的Scala DataFrame API实现

需求说明

需将现有左连接SQL查询完全转换为Spark DataFrame API实现,不在代码中直接执行SQL语句,最终在spark-shell中运行得到统计结果。核心逻辑为左连接后同时统计左表、右表的非空记录数,原开发过程中踩坑点为聚合右表字段时未正确设置别名,导致关联后引用右表字段时报列不存在错误。


原SQL逻辑

select count(x.user_id) user_id_count, count(w.id2) current_id2_count
from
    (select
        user_id
    from
        tb1
    where
        year='2021'
        and month=1
        
    ) x
left join
    (select id1, max(id2) id2 from tb2 group by id1) w
on
    x.user_id=w.id1;

初始错误写法及问题

错误代码

var x = spark.sqlContext.table("tb1").where("year='2021' and month=1")
var w= spark.sqlContext.table("tb2").groupBy("id1").agg(max("id2")).alias("id2")
var joined = x.join(w, x("user_id")===w("id1"), "left")

错误原因

对tb2聚合计算max(id2)时,未给聚合结果列显式设置别名,聚合后的DataFrame中该列默认名称为max(id2)而非预期的id2,后续关联后引用id2字段时就会报列不存在错误。


正确可运行实现

本地测试样例代码

// 构造左表测试数据,对应tb1过滤后的结果
var x = Seq("1","2","3","4").toDF("user_id")

// 构造右表测试数据,对应tb2原始数据
var w = Seq (("1", 1), ("1",2), ("3",10),("1",5),("5",4)).toDF("id1", "id2")

// 右表聚合计算max(id2),显式设置别名id2
var z= w.groupBy("id1").agg(max("id2").alias("id2"))

// 左连接后统计两个指标
val xJoinsZ= x.join(z, x("user_id") === z("id1"), "left").select(count(col("user_id").alias("user_id_count")), count(col("id2").alias("current_id2_count")))

运行验证结果

左表x数据

scala> x.show(false)
+-------+
|user_id|
+-------+
|1      |
|2      |
|3      |
|4      |
+-------+

右表聚合后z数据

scala> z.show(false)
+---+---+                                                                       
|id1|id2|
+---+---+
|3  |10 |
|5  |4  |
|1  |5  |
+---+---+

最终统计结果

scala> xJoinsZ.show(false)
+---------------------------------+---------------------------------+
|count(user_id AS `user_id_count`)|count(id2 AS `current_id2_count`)|
+---------------------------------+---------------------------------+
|4                                |2                                |
+---------------------------------+---------------------------------+

如果需要更简洁的输出列名,可将alias方法移到count函数外侧,调整为count(col("user_id")).alias("user_id_count")即可。


内容的提问来源于stack exchange,提问作者RRy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 20:06:03