请求解析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
相关产品推荐
相关产品推荐

