How can we flag duplicate records across business keys without immediately dropping them?
Using window count() partitioned by business keys assigns an occurrence count to identify duplicates.
Auditing data quality issues, identifying upstream replication bugs, and routing duplicates to quarantine.
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)Practice typing production-grade PySpark code for Duplicate Record Detection & Flagging.