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

请求解析Scala编写的Spark DataFrame代码的SQL执行逻辑

代码逻辑解析与Spark SQL实现

一、原代码的SQL逻辑拆解

这段Scala版Spark API代码的核心执行逻辑如下:

  • 从fact数据集中筛选出x、y、d、f、s这五个字段
  • 将筛选后的结果与unionTempDF数据集,以x、y、d、f为关联键执行内连接(INNER JOIN)
  • 对连接后的全量结果集去重,保留唯一行

二、对应的Spark SQL实现

你可以先将两个DataFrame注册为临时视图,再用SQL语句实现完全一致的逻辑:

// 把DataFrame注册为临时视图
unionTempDF.createOrReplaceTempView("union_temp")
fact.createOrReplaceTempView("fact_table")

// 执行Spark SQL得到结果
val tempFactDF = spark.sql("""
    SELECT DISTINCT ut.*, ft.s
    FROM union_temp ut
    INNER JOIN fact_table ft
    ON ut.x = ft.x 
    AND ut.y = ft.y 
    AND ut.d = ft.d 
    AND ut.f = ft.f
""")

说明:SQL里的DISTINCT和原代码中的dropDuplicates作用完全一致,都是对整行数据去重;如果只需要针对特定字段去重,可替换为SELECT * FROM ... DROP DUPLICATES(x, y, d, f),但原代码是全字段去重,所以用DISTINCT更匹配。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 11:01:13