Spark DataFrame的dropDuplicates方法在Spark SQL中的等价实现及代码转换问询
Great question! Let's break down exactly what your Spark code does, then map it step-by-step to equivalent Spark SQL.
First, let's clarify the behavior of your original code:
- It runs an inner join between
fooandbarusingcol1,col2, andcol3as join keys. UsingSeq("col1","col2","col3")in Spark's join means these columns only appear once in the resulting DataFrame (no duplicatedfoo.col1/bar.col1columns). - Then
dropDuplicates("col1","col2")removes rows wherecol1andcol2match, keeping only one arbitrary row per(col1, col2)pair. Spark doesn't guarantee which row is retained unless you specify ordering, which your code doesn't.
You're right that SELECT DISTINCT col1, col2, * isn't equivalent (and it's invalid SQL syntax anyway). DISTINCT applies to all selected columns, so it only removes rows where every single column is identical—not just rows sharing the same col1 and col2. That's a critical difference from dropDuplicates.
Solution 1: Window Functions (Works for Spark 3.0+)
The most accurate way to replicate your code's behavior is to use a ROW_NUMBER() window function. This assigns a unique number to each row in the (col1, col2) group, then we filter for rows where the number is 1. We'll use the USING clause to match Spark's join behavior (avoiding duplicate join columns):
WITH joined_data AS ( -- Replicate the inner join with merged join columns SELECT * FROM foo INNER JOIN bar USING (col1, col2, col3) ), deduped_data AS ( -- Assign row numbers to each (col1, col2) group SELECT * FROM ( SELECT *, -- ORDER BY (SELECT NULL) mimics Spark's arbitrary row selection in dropDuplicates ROW_NUMBER() OVER (PARTITION BY col1, col2 ORDER BY (SELECT NULL)) AS rn FROM joined_data ) t -- Keep only one row per (col1, col2) group WHERE rn = 1 ) -- Remove the temporary row number column SELECT * EXCEPT (rn) FROM deduped_data;
Solution 2: Simplified with QUALIFY (Spark 3.1+)
If you're on Spark 3.1 or later, you can use the QUALIFY clause to skip the intermediate CTE and make the query more concise:
SELECT * FROM foo INNER JOIN bar USING (col1, col2, col3) -- Filter directly after calculating the window function QUALIFY ROW_NUMBER() OVER (PARTITION BY col1, col2 ORDER BY (SELECT NULL)) = 1;
Key Notes
ORDER BY (SELECT NULL)tells Spark not to enforce any specific order when picking which row to keep, exactly matchingdropDuplicates's behavior (since row order isn't guaranteed in distributed Spark environments).- If you want to explicitly retain a specific row (e.g., the one with the largest
col3value), replace(SELECT NULL)with a column and order direction likeORDER BY col3 DESC.
内容的提问来源于stack exchange,提问作者Randomize

