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

如何去除PySpark DataFrame自连接及三列结果中的重复行?

Fixing Duplicate Rows in PySpark Self-Joined Movie Cast Data

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:

  1. Exact duplicate rows (all three columns identical)
  2. 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_name filter 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 08:55:09