如何去除PySpark DataFrame自连接及三列结果中的重复行?
Got it, let's tackle this problem step by step—dealing with duplicate rows from a PySpark self-join (like bidirectional person pairs and exact row duplicates) is super common when working with co-occurrence data like movie casts. Here's how to clean up your DataFrame properly:
First, Let's Contextualize the Scenario
Suppose your original DataFrame (movie_cast) links movie IDs to cast members, like this:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, least, greatest spark = SparkSession.builder.appName("CleanCoActors").getOrCreate() # Sample input data (includes an exact duplicate row for demonstration) data = [ ("movie1", "Alice"), ("movie1", "Bob"), ("movie1", "Charlie"), ("movie2", "Alice"), ("movie2", "Bob"), ("movie2", "Dave"), ("movie1", "Alice") ] movie_cast = spark.createDataFrame(data, ["movie_id", "person_name"])
When you self-join to get co-actor pairs, you'll end up with two types of duplicates:
- Exact duplicate rows (all three columns identical)
- Reversed person pairs (e.g.,
(Alice, Bob, movie1)and(Bob, Alice, movie1)for the same movie)
Let's fix both issues systematically.
Step 1: Remove Exact Duplicate Rows First
First, run your self-join (with a critical filter to avoid pairing a person with themselves), then use dropDuplicates() to eliminate exact row duplicates:
# Perform self-join, filtering out same-person pairs joined_df = movie_cast.alias("a")\ .join(movie_cast.alias("b"), (col("a.movie_id") == col("b.movie_id")) & (col("a.person_name") != col("b.person_name")), "inner")\ .select( col("a.person_name").alias("person1"), col("b.person_name").alias("person2"), col("a.movie_id").alias("movie_id") ) # Remove exact duplicate rows deduped_exact = joined_df.dropDuplicates(["person1", "person2", "movie_id"])
Step 2: Eliminate Bidirectional Person Pair Duplicates
The trick here is to standardize the order of person names so reversed pairs become identical. We'll use PySpark's least() and greatest() functions—these will always put the lexicographically smaller name first, then the larger one. We can then deduplicate based on this standardized pair:
# Standardize the order of person names in each pair standardized_df = deduped_exact\ .withColumn("person1_standard", least(col("person1"), col("person2")))\ .withColumn("person2_standard", greatest(col("person1"), col("person2"))) # Deduplicate based on the standardized pair + movie ID final_clean_df = standardized_df\ .dropDuplicates(["person1_standard", "person2_standard", "movie_id"])\ .select( col("person1_standard").alias("person1"), col("person2_standard").alias("person2"), col("movie_id") ) # View the cleaned result final_clean_df.show()
The output will have unique, non-reversed co-actor pairs per movie, with no exact duplicates:
+-------+-------+--------+ |person1|person2|movie_id| +-------+-------+--------+ | Alice| Bob| movie1| | Alice|Charlie| movie1| | Bob|Charlie| movie1| | Alice| Bob| movie2| | Alice| Dave| movie2| | Bob| Dave| movie2| +-------+-------+--------+
Key Notes
- Don't skip the
a.person_name != b.person_namefilter in your self-join—it avoids useless rows where a person is paired with themselves. If you forget it, add a final filter:final_clean_df = final_clean_df.filter(col("person1") != col("person2")) - The
least()/greatest()logic works for string names, but you can adapt it to numeric person IDs too—just use the same standardization approach.
内容的提问来源于stack exchange,提问作者silviacamplani

