Spark-SQL高效实现:基于信任列筛选合并单一行(多列场景)
Spark-SQL(Databricks)中多列信任值选取的优化方案
输入数据与需求
输入表结构与数据
表person_details的结构及数据如下(null代表实际空值,trust_on_*列取值为10的整数倍或null,需用COALESCE将此类null替换为0):
| id | name | trust_on_name | age | trust_on_age |
|---|---|---|---|---|
| 1 | John Doe | 90 | null | null |
| 1 | john D. | 50 | 25 | 90 |
| 1 | null | 0 | twenties | 50 |
需求说明
以id为主键,为每个id生成单一行数据:每一列的值,选取对应trust_on_<列名>列最大值所在行的对应值。最终输出结果如下:
| id | name | age |
|---|---|---|
| 1 | John Doe | 25 |
现有方案的局限性
当前采用多CTE关联的实现方式,当表中存在大量<列名>, trust_on_<列名>列对时,需要创建大量CTE并逐一关联,代码冗余度极高:
WITH cte_name AS ( SELECT id, `name` AS name, trust_on_name, null as `age`, trust_on_age FROM person_details QUALIFY trust_on_name = MAX(trust_on_name) OVER (PARTITION BY id) ), cte_age AS ( SELECT id, null as `name`, trust_on_name, age, trust_on_age FROM person_details QUALIFY trust_on_age = MAX(trust_on_age) OVER (PARTITION BY id) ) SELECT A.id, A.name, B.age FROM cte_name A INNER JOIN cte_age B ON A.id = B.id WHERE A.trust_on_name > COALESCE(B.trust_on_name, 0) AND B.trust_on_age > COALESCE(A.trust_on_age, 0)
优化实现方案
针对大量列对的场景,推荐以下几种更简洁的实现方式:
方案一:窗口函数+条件聚合
先通过窗口函数计算每个id下各trust_on_*列的最大值,再通过条件聚合筛选出对应列的目标值,无需创建多个CTE:
SELECT id, -- 选取trust_on_name最大值对应的name MAX(CASE WHEN COALESCE(trust_on_name, 0) = max_trust_name THEN name END) AS name, -- 选取trust_on_age最大值对应的age MAX(CASE WHEN COALESCE(trust_on_age, 0) = max_trust_age THEN age END) AS age FROM ( SELECT *, MAX(COALESCE(trust_on_name, 0)) OVER (PARTITION BY id) AS max_trust_name, MAX(COALESCE(trust_on_age, 0)) OVER (PARTITION BY id) AS max_trust_age FROM person_details ) t GROUP BY id, max_trust_name, max_trust_age
方案二:逆透视+透视(适配大量列场景)
通过将列对转换为行数据(逆透视),按id和列名筛选最高信任值对应的内容,再转换回原列结构(透视),适合列对数量较多的场景:
-- 1. 逆透视:将列对转为行数据 WITH unpivoted AS ( SELECT id, 'name' AS col_name, name AS col_value, COALESCE(trust_on_name, 0) AS trust_score FROM person_details UNION ALL SELECT id, 'age' AS col_name, age AS col_value, COALESCE(trust_on_age, 0) AS trust_score FROM person_details ), -- 2. 按id和列名排序,取最高信任值的行 ranked AS ( SELECT id, col_name, col_value, ROW_NUMBER() OVER (PARTITION BY id, col_name ORDER BY trust_score DESC) AS rn FROM unpivoted ) -- 3. 透视回原列结构 SELECT id, MAX(CASE WHEN col_name = 'name' THEN col_value END) AS name, MAX(CASE WHEN col_name = 'age' THEN col_value END) AS age FROM ranked WHERE rn = 1 GROUP BY id
方案三:动态生成SQL(适配超大量列场景)
如果列对数量极多,可通过Spark的元数据动态生成SQL,避免手动编写重复代码:
// 读取目标表 val df = spark.table("person_details") // 提取所有<列名, trust_on_列名>的配对 val columnPairs = df.columns .filter(_.startsWith("trust_on_")) .map(trustCol => (trustCol.replace("trust_on_", ""), trustCol)) // 生成条件聚合语句 val caseStatements = columnPairs.map { case (col, trustCol) => s"MAX(CASE WHEN COALESCE(`$trustCol`, 0) = max_$trustCol THEN `$col` END) AS `$col`" }.mkString(",\n ") // 生成最大信任值计算语句 val maxTrustStatements = columnPairs.map { case (col, trustCol) => s"MAX(COALESCE(`$trustCol`, 0)) AS max_$trustCol" }.mkString(",\n ") // 拼接完整SQL val dynamicSql = s""" SELECT id, $caseStatements FROM ( SELECT *, $maxTrustStatements FROM person_details ) t GROUP BY id, ${columnPairs.map(_._2).map(c => s"max_$c").mkString(", ")} """ // 执行SQL并展示结果 spark.sql(dynamicSql).show()
内容的提问来源于stack exchange,提问作者Lingesh.K
相关产品推荐
相关产品推荐

