Skip to main content
PYSPARK • LESSON 212

Duplicate Record Detection & Flagging

How can we flag duplicate records across business keys without immediately dropping them?

Advanced3 Minutes810 XP
🤔 THE QUESTION

How can we flag duplicate records across business keys without immediately dropping them?

💡 WHAT IS IT?

Using window count() partitioned by business keys assigns an occurrence count to identify duplicates.

🎯 WHAT IS IT USED FOR?

Auditing data quality issues, identifying upstream replication bugs, and routing duplicates to quarantine.

💻 EXAMPLE
w = Window.partitionBy("order_id", "line_item_id")
df_flagged = df.withColumn("duplicate_count", count("order_id").over(w)).withColumn("is_duplicate", col("duplicate_count") > 1)

🎯 Mission Objectives

Practice typing production-grade PySpark code for Duplicate Record Detection & Flagging.

  • Window partitioning by business key
  • Duplicate occurrence count
  • Audit quality flagging