如何将PySpark指定列去重逻辑转换为兼容dbt的SQL语句?
如何将PySpark的drop_duplicates逻辑转换为SQL(适用于dbt)
你原来的SQL写法逻辑完全错误——GROUP BY column_3是按column_3分组聚合,和PySpark里按column_1、column_2组合去重的逻辑根本不匹配,结果自然不一致。
正确的思路是:针对column_1和column_2的每一组重复值,保留其中一行对应的column_3。下面是两种通用的实现方式,适配大多数SQL方言(包括dbt常用的Spark SQL、PostgreSQL等):
方法1:窗口函数(通用所有SQL引擎)
用ROW_NUMBER()窗口函数给每组column_1+column_2的行编号,然后取编号为1的行:
WITH ranked_rows AS ( SELECT column_1, column_2, column_3, -- 按column_1和column_2分组,给每组内的行编号 ROW_NUMBER() OVER (PARTITION BY column_1, column_2 ORDER BY (SELECT NULL)) AS row_num FROM my_table ) SELECT column_1, column_2, column_3 FROM ranked_rows WHERE row_num = 1
ORDER BY (SELECT NULL)表示不指定排序规则,和PySpark默认保留第一个出现的行逻辑一致;如果需要稳定保留某一行(比如最大/最小的column_3),可以把排序条件改成ORDER BY column_3 DESC或ASC。
方法2:DISTINCT ON(适用于Spark SQL、PostgreSQL)
如果你的dbt对接的是支持DISTINCT ON的SQL引擎,可以用更简洁的写法:
SELECT DISTINCT ON (column_1, column_2) column_1, column_2, column_3 FROM my_table -- 可选:指定保留哪一行,比如保留column_3最大的那行 -- ORDER BY column_1, column_2, column_3 DESC
这个语法会直接返回每组column_1+column_2的第一行,和PySpark的drop_duplicates行为一致。
内容的提问来源于stack exchange,提问作者The Singularity
相关产品推荐
相关产品推荐

