Where duplicates come from, why DISTINCT is not a strategy, and the three questions to answer first
A pipeline can finish successfully. The dashboard loads. Every job shows green. And your numbers can still be wrong.
One of the most common reasons is also one of the quietest: duplicate records.
Start with the smallest possible example. Your orders table has Order 101 for $100. A retry sends the same order a second time, so now there are two rows. Nothing crashes. The schema is valid. The pipeline completes in the usual four minutes. Then finance calculates revenue, and that $100 order becomes $200.
Duplicates rarely cause technical failures. They cause business failures. That is why they survive so long in production.
https://youtu.be/iZZCF5YDK9E
A broken pipeline is loud. Someone gets paged, the run turns red, and the fix happens the same day.
A duplicated row is silent. It flows into the gold layer, into a dashboard, into a board deck, into a forecast. By the time somebody notices that revenue looks slightly high, the number has already been used to make a decision.
This is the same failure shape we wrote about in the pipeline was green, the numbers were wrong. Correctness is not the same thing as completion. A green run only tells you the code finished. It says nothing about whether each row in the table deserves to be there.
Before you write any dedup logic, it helps to know which of these four doors the duplicate walked through. Each one needs a different response.
A batch is written, then the acknowledgement fails on the way back. The sending system does the safe thing and resends. Now the same file, or the same batch of rows, lands twice.
This is why we keep coming back to safe reruns and idempotency. A pipeline that can be run twice without changing the result is a pipeline that survives retries. A pipeline that appends blindly will grow a duplicate every time the network hiccups.
Most streaming and messaging systems promise at-least-once delivery. Read that promise carefully. It means the platform would rather send an event twice than lose it once.
So in an event pipeline, duplicates are not an anomaly. They are the documented behaviour. If your design assumes exactly-once arrival, the design is wrong, not the platform. We cover the mechanics of this in Streaming and in Incremental Processing.
The website sends an order. The mobile app sends the same order. A partner system sends it again in a nightly file, with a different column order and a slightly different timestamp.
Three rows, one business event. Nothing is technically duplicated because no two rows are byte-identical. But the business truth was created once.
This is the duplicate that most people miss, because the source data was clean.
One customer has three addresses. You join customers to addresses. That one customer is now three rows. Then somebody sums revenue by customer, and the total is tripled.
Nothing was duplicated at ingestion. The duplicate was created by the data model, in your own code. Joins that fan out are the most common source of wrong totals in a warehouse, which is why we spend real time on it in Joins and Aggregations.
-- One order, three address rows, revenue counted three times
SELECT
c.customer_id,
SUM(o.order_amount) AS revenue
FROM orders o
JOIN customers c ON o.customer_id = c.customer_id
JOIN addresses a ON c.customer_id = a.customer_id
GROUP BY c.customer_idThe join to addresses adds nothing to the calculation. It only changes the row count. And changing the row count changes the answer.
Before removing a single duplicate, you have to answer one question: what is one row supposed to represent?
That is the grain. It is the contract of the table.
If the grain is one row per order, then order_id should be unique, and a second row with the same order_id is a duplicate.
If the grain is one row per order line item, then a single order can appear five times legitimately, and the identity of a row is order_id plus line_number.
Without the grain, deduplication is guesswork. You will either remove valid rows or keep invalid ones, and you will not be able to tell which you did. Writing the grain down as a one-line comment at the top of every table definition is one of the cheapest habits in data engineering.
| Grain | Unique key | A second row with the same key means |
|---|---|---|
| One row per order | order_id | Duplicate |
| One row per order line | order_id, line_number | Duplicate |
| One row per order status change | order_id, status, changed_at | Valid history |
| One row per customer, current state | customer_id | Duplicate, needs a merge |
Notice the third row. In a history table, repeated order_id values are the point. This is the world of SCD patterns, where keeping the old version is a feature, not a defect. The same physical shape, repeated keys, is correct in one table and a bug in another. Only the grain tells you which.
SELECT DISTINCT is the first thing most people reach for, and it is almost always the wrong tool. It compares whole rows, and whole rows are not what identifies a business event.
It fails in both directions.
It keeps duplicates you wanted gone. The same order arrives twice with two different ingested_at timestamps. The rows are not identical, so DISTINCT keeps both, and revenue is still doubled.
It removes rows you needed. A customer genuinely buys the same product twice in the same second at the same price. Two real events, identical values, and DISTINCT silently deletes one. You have just lost revenue instead of inflating it.
There is a deeper problem. DISTINCT hides the symptom without explaining the cause. Whatever created those extra rows is still running, and now nobody can see it.
Good deduplication is three decisions, made in order. The SQL is the easy part.
One: define the business key. What identifies one unique business event? An order_id. An order_id plus item_id. An event_id from the producer. Write it down.
Two: define which record survives. Latest update wins. Earliest event wins. The trusted source wins. The higher version number wins. This is a business rule, not a technical preference, and it usually needs a conversation with whoever owns the data.
Three: implement it, and only then. The common pattern is a window function.
WITH ranked_orders AS (
SELECT
*,
ROW_NUMBER() OVER (
PARTITION BY order_id
ORDER BY updated_at DESC, ingested_at DESC
) AS row_num
FROM bronze.orders
)
SELECT * EXCEPT (row_num)
FROM ranked_orders
WHERE row_num = 1Partition by the business key. Order by the rule you chose. Keep row 1.
That query is correct only if your business rule says the latest record should win. If the rule is that the first arrival wins, you flip the sort. If the rule is that the ERP beats the web form, you sort by a source priority column first. The rule matters more than the SQL, and the SQL should be readable enough that a reviewer can check the rule from the ORDER BY clause.
In PySpark, the same idea, written so each step is visible:
from pyspark.sql import functions as F
from pyspark.sql.window import Window
order_window = Window.partitionBy("order_id").orderBy(
F.col("updated_at").desc(),
F.col("ingested_at").desc(),
)
ranked_orders = bronze_orders.withColumn(
"row_num", F.row_number().over(order_window)
)
deduped_orders = ranked_orders.filter(F.col("row_num") == 1).drop("row_num")For a table you keep updating, MERGE is the natural home for this. Deduplicate the incoming batch first, then merge on the business key so a rerun updates instead of appends. We walk through that end to end in Delta Lake and again in the hands-on first data engineering project.
MERGE INTO silver.orders AS target
USING deduped_orders AS source
ON target.order_id = source.order_id
WHEN MATCHED AND source.updated_at > target.updated_at THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *Here is the habit that separates a working pipeline from a trustworthy one.
If your duplicate rate sits at around 1% every day, that is your normal, and dedup logic handles it quietly. If it jumps to 20% on a Tuesday, something changed upstream. A retry loop. A broken join. A change data capture reset that replayed history. A job that ran twice because somebody triggered it manually while the schedule was running.
Deduplication that silently absorbs a 20% spike is not protecting you. It is deleting the evidence. The goal is to protect downstream tables while keeping the problem visible.
So measure it. Track the raw row count and the deduplicated count, store the difference and the percentage, set a threshold, and alert on change rather than only on absolute value. Keep a little metadata on the surviving row, such as how many source records collapsed into it and where they came from. That single column turns a future argument into a two-minute query.
SELECT
ingest_date,
COUNT(*) AS raw_rows,
COUNT(DISTINCT order_id) AS unique_orders,
ROUND(
100.0 * (COUNT(*) - COUNT(DISTINCT order_id)) / COUNT(*), 2
) AS duplicate_percent
FROM bronze.orders
GROUP BY ingest_date
ORDER BY ingest_date DESCRun that query on your own bronze table this week. Most people are surprised the first time.
This is the same discipline as expectations and constraints in Data Quality. A check that fails loudly is worth more than a transformation that fixes things quietly.
Streaming does not change the principle. It changes the time window.
In a stream you cannot compare a new event against all of history forever, because state has to be bounded. So you need three things: a stable event ID from the producer, a watermark that tells the engine how late an event may arrive, and a deduplication window inside that boundary.
That is why late data and duplicates are usually the same conversation, and why we handled them together in the event arrived two days late. An event that shows up outside the window is not deduplicated by the stream. It has to be handled by a batch correction job, and you have to decide in advance that this is acceptable.
The rule stays the same in both worlds: define what makes an event unique before you remove anything.
Before you write dedup logic, answer three questions.
What is the grain? What does one row represent in this table?
What is the business key? Which columns identify one unique business event?
Which record should survive? And who decided that, you or the business?
If you can answer all three in one sentence each, the SQL will almost write itself. If you cannot, no amount of window functions will save the table.
Duplicates are not just extra rows. They quietly change revenue, customer counts, conversion rates, inventory levels and risk scores. They pass every schema check. They never turn a run red.
The most dangerous duplicate is not the one that crashes your pipeline. It is the one that quietly changes a business decision.
If you want the full path, from your first DataFrame to pipelines you can rerun without fear, that is what Thinking in Data Engineering with Databricks is for. Everything runs in Databricks Free Edition.