Deduplicate records
Keeping one row per entity in a pipeline stage.
Source data often holds the same record more than once. A sender retries a webhook after a timeout and the same event arrives twice. A monthly export overlaps the previous one by a few days, so a File Set that holds both files contains those rows twice. A system sends a new version of an order every time its status changes, and you only want the current one.
This recipe builds a pipeline that turns such a source into a dataset with one row per key. It uses a Webhook dataset as the example, because its JSON body has to be unpacked before you can tell which deliveries describe the same thing. For a CSV or Parquet source, skip the first 2 stages.
Decide what a duplicate is
Pick the key that identifies one entity, such as order_id, and decide which kind of duplicate
you have:
| Kind | Example | Keep |
|---|---|---|
| Exact | The same row imported twice from overlapping files | Any one copy |
| Versions | Several updates of the same order | The newest version per key |
Webhook rows are never exact duplicates as they land. Each delivery gets its own message_id and
ingested_at, so 2 deliveries of the same event differ in those columns even when the body is
identical. Compare the fields inside payload instead of the envelope.
The pipeline
The pipeline has 1 source, the webhook dataset with the alias orders_events. Each stage below is
one stage in the builder and reads the previous one by its stage name.
Stage decoded turns the base64 payload into JSON text.
SELECT
message_id,
ingested_at,
CAST(decode(payload, 'base64') AS VARCHAR) AS body
FROM orders_eventsStage orders extracts the fields you need into typed columns, using the typed getters listed
under JSON functions.
SELECT
json_get_str(body, 'order_id') AS order_id,
json_get_str(body, 'status') AS status,
TRY_CAST(json_as_text(body, 'total') AS DECIMAL(12, 2)) AS total,
to_timestamp(json_get_str(body, 'updated_at')) AS updated_at,
ingested_at,
message_id
FROM decodedStage ranked numbers the versions of each order, newest first.
SELECT
order_id,
status,
total,
updated_at,
ROW_NUMBER() OVER (
PARTITION BY order_id
ORDER BY updated_at DESC NULLS LAST, ingested_at DESC, message_id DESC
) AS rn
FROM orders
WHERE order_id IS NOT NULLStage latest_orders keeps the first row per key. It is the last stage, so it feeds the
destination.
SELECT order_id, status, total, updated_at
FROM ranked
WHERE rn = 1The ORDER BY inside ROW_NUMBER decides which row survives. updated_at picks the newest
version, ingested_at breaks ties between versions with the same updated_at, and message_id
settles an event delivered twice in the same instant. Without a tiebreak that is unique per row, 2
rows can share first place and either may win, so the same input can give different output on
different runs.
In descending order NULL sorts first by default, which is why NULLS LAST is there: a version
without updated_at would otherwise beat every dated one.
Exact duplicates only
If every copy of a row is identical and there is no version to choose, SELECT DISTINCT over the
business columns is enough. Leave out columns that differ per copy, such as message_id:
SELECT DISTINCT order_id, status, total, updated_at
FROM ordersSwitch to ROW_NUMBER as soon as 2 rows can share a key and differ in any other column, because
DISTINCT keeps both of them.
Write mode
The pipeline above reads the whole source on every run, so its last stage always holds the complete deduplicated table. Write it with Overwrite, full table. Each run replaces the destination with a fresh result, and a version that arrives late still ends up in the right place.
Use Merge with order_id as the merge key when a run reads only part of the source, for example
one day with WHERE dt = '{{ run.dt }}'. Merge updates the rows whose key is already in the
destination and inserts the rest. Keep the ranked and latest_orders stages, because merge does
not deduplicate the stage result and one slice can hold several versions of the same order.
A merge replaces a matched row with whatever the current run wrote. If an older version of an order can arrive after a newer one, it overwrites the newer one. When that can happen, read the full history and overwrite.
See Write modes and partitioning.
Check the result
After a run, query the destination by name in Query. With a destination called
orders_latest, the row count should equal the number of distinct keys:
SELECT COUNT(*) AS row_count, COUNT(DISTINCT order_id) AS key_count
FROM orders_latestTo list keys that still appear more than once:
SELECT order_id, COUNT(*) AS copies
FROM orders_latest
GROUP BY order_id
HAVING COUNT(*) > 1
ORDER BY copies DESC
LIMIT 100An empty result means every key is unique. Before you deploy, you can check the same thing by
previewing the ranked stage and looking for rows with rn greater than 1.
Pitfalls
Rows with a NULL key never match anything in a merge, so every run inserts them again. Filter them
out, as ranked does, or send them to a separate dataset where you can inspect them.
Keys that differ only in whitespace or case count as different keys, so 'A-100' and 'a-100 '
both survive. Normalize with lower(trim(...)) in the extraction stage when the source is
inconsistent. trim accepts only strings, so cast a numeric key first:
trim(CAST(order_no AS VARCHAR)).
Collecting rows into an array, taking its first element and joining that back to the source gives
an unstable answer. Use ROW_NUMBER with a full tiebreak.
Timestamps stored as text sort correctly only when every value has the same format and time zone.
Convert them with to_timestamp before you order by them.
Related
Last updated on