Search⌘KHide sidebarOpen menu
Switch to dark mode
Copy page content as Markdown⌘⌥C
Login
Pipelines

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:

KindExampleKeep
ExactThe same row imported twice from overlapping filesAny one copy
VersionsSeveral updates of the same orderThe 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_events

Stage 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 decoded

Stage 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 NULL

Stage 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 = 1

The 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 orders

Switch 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_latest

To 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 100

An 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.

Last updated on

On this page