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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 23:24:46