如何使用Spark SQL按固定参数(id,date,name)实现表扁平连接
Spark 数据集全外连接并填充缺失值
Spark SQL 实现
假设两个数据集已注册为临时视图 d1 和 d2,可通过以下SQL语句完成需求:
SELECT COALESCE(d1.id, d2.id) AS id, COALESCE(d1.date, d2.date) AS date, COALESCE(d1.name, d2.name) AS name, COALESCE(d1.val1, 0) AS val1, COALESCE(d2.val2, 0) AS val2 FROM d1 FULL OUTER JOIN d2 ON d1.id = d2.id AND d1.date = d2.date AND d1.name = d2.name ORDER BY id, date, name;
逻辑说明
FULL OUTER JOIN保留两个数据集中所有的(id, date, name)组合,无论对方数据集是否存在匹配项COALESCE函数处理空值:id/date/name优先取d1的值,d1无数据时取d2的值val1和val2无匹配项时,用0填充空值
ORDER BY用于对齐示例输出的排序顺序,属于非必需的格式化操作
DataFrame API 实现
Scala 版本
import org.apache.spark.sql.functions._ val result = d1.join(d2, Seq("id", "date", "name"), "full_outer") .withColumn("val1", coalesce(col("val1"), lit(0))) .withColumn("val2", coalesce(col("val2"), lit(0))) .orderBy("id", "date", "name") result.show()
Python 版本
from pyspark.sql.functions import coalesce, lit result = d1.join(d2, on=["id", "date", "name"], how="full_outer") \ .withColumn("val1", coalesce("val1", lit(0))) \ .withColumn("val2", coalesce("val2", lit(0))) \ .orderBy("id", "date", "name") result.show()
逻辑说明
- 通过
join方法指定连接键和全外连接类型(full_outer) - 利用
coalesce函数将空值字段替换为0 orderBy保证结果排序与示例输出一致
内容的提问来源于stack exchange,提问作者Vikneswaran S J
相关产品推荐
相关产品推荐

