Design idempotent daily loads with deduping
Company: Amazon
Role: Data Scientist
Category: Data Manipulation (SQL/Python)
Difficulty: medium
Interview Round: Technical Screen
You need to load the last 7 days of orders into a large fact table from a noisy staging feed. Assume today is 2025-09-01. Requirements: idempotent reruns, late-arriving updates, and duplicate detection. Invent a complete approach (DDL + SQL/Python) that: 1) dedupes staging by (order_id, updated_at) keeping the latest; 2) MERGEs into the fact table only when the incoming row is newer; 3) detects whether a row was already loaded; 4) supports partitioning and incremental watermarks; 5) captures load audit metrics and data-quality checks (row counts, nulls, referential integrity). Use the schema and tiny sample data below, then write the exact SQL for the dedupe CTE and the MERGE (or equivalent upsert), plus a brief Python snippet (pseudo-code OK) for the orchestration and idempotent checkpointing.
Schema:
- staging_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), source_file STRING, ingest_time TIMESTAMP)
- fact_orders(order_id INT PRIMARY KEY, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), load_batch_id STRING, loaded_at TIMESTAMP)
- dim_customer(customer_id INT PRIMARY KEY, signup_date DATE)
- loads_audit(batch_id STRING PRIMARY KEY, run_start TIMESTAMP, run_end TIMESTAMP, watermark_from DATE, watermark_to DATE, staged_rows INT, deduped_rows INT, inserted_rows INT, updated_rows INT, dq_errors INT)
Sample rows (monospaced ASCII):
staging_orders (arriving 2025-09-01 for last 7 days)
| order_id | customer_id | order_date | updated_at | total_amount | source_file | ingest_time |
|---------:|------------:|-------------|----------------------|-------------:|------------------|----------------------|
| 100 | 1 | 2025-08-26 | 2025-08-26 10:00:00 | 49.00 | orders_20250826 | 2025-09-01 00:05:00 |
| 101 | 2 | 2025-08-28 | 2025-08-28 09:00:00 | 20.00 | orders_20250828 | 2025-09-01 00:05:10 |
| 101 | 2 | 2025-08-28 | 2025-08-30 12:30:00 | 22.00 | latefix_20250830 | 2025-09-01 00:05:10 |
| 102 | 3 | 2025-08-31 | 2025-08-31 21:00:00 | 15.50 | orders_20250831 | 2025-09-01 00:05:20 |
| 103 | 9 | 2025-09-01 | 2025-09-01 00:01:00 | 105.00 | orders_20250901 | 2025-09-01 00:05:30 |
fact_orders (before load)
| order_id | customer_id | order_date | updated_at | total_amount | load_batch_id | loaded_at |
|---------:|------------:|-------------|----------------------|-------------:|---------------|----------------------|
| 101 | 2 | 2025-08-28 | 2025-08-28 09:00:00 | 20.00 | batch_0828 | 2025-08-28 10:00:00 |
dim_customer
| customer_id | signup_date |
|------------:|-------------|
| 1 | 2025-05-01 |
| 2 | 2024-11-11 |
| 3 | 2025-08-01 |
Answer specifics: a) show the exact SQL to dedupe staging (window function or aggregate); b) show a MERGE that updates when s.updated_at > t.updated_at and inserts if not matched; c) propose a robust watermark (e.g., max(order_date) − 1 day) to catch late data, and where you store it; d) show how you’d compute and store a row hash to detect prior loads; e) list two data-quality assertions that should fail the batch (and how to rollback safely).
Overview: This question evaluates data engineering competencies including idempotent ETL design, deduplication and upsert semantics, handling late-arriving updates, partitioning and incremental watermarks, load auditing and data-quality checks using SQL and orchestration with Python.
Read the full Amazon Data Scientist interview experience this question came from
Deduplicate noisy staging orders for the last 7-day window
You are loading orders from a noisy staging feed into a fact table. For this exercise, treat the run date as 2025-06-01 and the 7-day load window as 2025-05-26 through 2025-06-01 (inclusive).
Write a query that deduplicates `staging_orders` to ONE row per `order_id` within the window, keeping the row with the greatest `updated_at`. If there are ties on `updated_at`, keep the row with the greatest `ingest_time`.
Return the deduped rows with columns: `order_id, customer_id, order_date, updated_at, total_amount, source_file, ingest_time` ordered by `order_id`. For the sample output, format updated_at and ingest_time as YYYY-MM-DD HH24:MI:SS.
Tables
staging_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), source_file VARCHAR(64), ingest_time TIMESTAMP)
Hints
- Use ROW_NUMBER() partitioned by order_id and ordered by updated_at desc, ingest_time desc.
- Filter to the explicit date window before deduping.
Idempotent MERGE into fact table (update only when staging is newer)
Using the same 7-day window (2025-05-26..2025-06-01 inclusive), perform an idempotent upsert into `fact_orders` from `staging_orders`:
- First dedupe staging to one row per `order_id` (latest `updated_at`, tie-breaker `ingest_time`).
- Only load rows that pass basic data quality gates:
- `total_amount` is NOT NULL
- `customer_id` exists in `dim_customer` (referential integrity)
- MERGE behavior:
- When matched on `order_id`, UPDATE only if `s.updated_at > t.updated_at`.
- When not matched, INSERT.
- For all inserted/updated rows, set:
- `load_batch_id = 'batch_0601'`
- `loaded_at = TIMESTAMP '2025-06-01 00:10:00'`
- `row_hash` computed as `MD5(CONCAT(order_id,'|',customer_id,'|',order_date,'|',updated_at,'|',total_amount))`.
Write the MERGE SQL (or equivalent upsert). After the MERGE, return the contents of `fact_orders` ordered by `order_id` (excluding `row_hash` from the output). For validation, return the final fact_orders state that the MERGE would produce, with timestamps formatted as YYYY-MM-DD HH24:MI:SS.
Tables
staging_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), source_file VARCHAR(64), ingest_time TIMESTAMP)
fact_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), row_hash VARCHAR(64), load_batch_id VARCHAR(40), loaded_at TIMESTAMP)
dim_customer(customer_id INT, signup_date DATE)
Hints
- Dedupe first, then apply DQ filters (non-null and dim join) before MERGE.
- Only update when staging.updated_at is strictly greater than fact.updated_at.
Compute an incremental watermark with a 1-day lookback overlap
To support incremental loads with late-arriving updates, compute the next watermark range for batch_id = 'batch_0601'.
Rules:
- The run date is 2025-06-01, so set `watermark_to = DATE '2025-06-01'`.
- Let `prev_watermark_to` be the maximum `watermark_to` in `loads_audit`.
- Set `watermark_from = prev_watermark_to - 1 day` (1-day lookback overlap). If there is no previous run, use `watermark_from = DATE '2025-05-26'`.
Return a single row with columns: `batch_id, watermark_from, watermark_to`.
Tables
loads_audit(batch_id VARCHAR(40), run_start TIMESTAMP, run_end TIMESTAMP, watermark_from DATE, watermark_to DATE, staged_rows INT, deduped_rows INT, inserted_rows INT, updated_rows INT, dq_errors INT)
Hints
- Use MAX(watermark_to) from loads_audit as the previous checkpoint.
- Subtract a 1-day interval to create an overlap for late-arriving updates.
Classify deduped staging rows as INSERT/UPDATE/SKIP/REJECT (idempotency planning)
For the 2025-05-26..2025-06-01 window, produce a pre-merge classification per deduped `order_id` to support idempotent reruns.
Steps:
1) Dedupe staging to one row per order_id (latest updated_at; tie-breaker ingest_time).
2) Compute `incoming_hash = MD5(CONCAT(order_id,'|',customer_id,'|',order_date,'|',updated_at,'|',total_amount))`.
3) LEFT JOIN to `fact_orders` on `order_id` and LEFT JOIN to `dim_customer` on `customer_id`.
4) Assign an `action` for each deduped row:
- 'REJECT_NULL_TOTAL_AMOUNT' if total_amount IS NULL
- 'REJECT_DIM_CUSTOMER' if customer_id not found in dim_customer
- 'INSERT' if not in fact_orders
- 'UPDATE' if in fact_orders AND staging.updated_at > fact.updated_at
- 'SKIP_ALREADY_LOADED' if in fact_orders AND staging.updated_at = fact.updated_at AND values are identical (use hash comparison; if fact.row_hash is NULL, compute it from fact columns)
- 'SKIP_NOT_NEWER' otherwise
Return: `order_id, action` ordered by order_id.
Tables
staging_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), source_file VARCHAR(64), ingest_time TIMESTAMP)
fact_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), row_hash VARCHAR(64), load_batch_id VARCHAR(40), loaded_at TIMESTAMP)
dim_customer(customer_id INT, signup_date DATE)
Hints
- Compute incoming hash and a comparable fact hash (use COALESCE for missing fact row_hash).
- Order the CASE checks so DQ rejects happen before deciding insert/update/skip.
Data-quality checks: nulls and referential integrity failures in the batch window
For the batch window 2025-05-26..2025-06-01 (inclusive), write SQL that outputs all data-quality failures that should cause the batch to fail.
Rules:
- First dedupe staging to one row per order_id (latest updated_at; tie-breaker ingest_time).
- Then emit one row per failing deduped record with:
- `order_id`
- `dq_issue`
Include at least these two assertions:
1) `total_amount` is NULL -> dq_issue = 'NULL_TOTAL_AMOUNT'
2) `customer_id` not found in dim_customer -> dq_issue = 'MISSING_DIM_CUSTOMER'
Return all failures ordered by `order_id, dq_issue`.
Tables
staging_orders(order_id INT, customer_id INT, order_date DATE, updated_at TIMESTAMP, total_amount DECIMAL(10,2), source_file VARCHAR(64), ingest_time TIMESTAMP)
dim_customer(customer_id INT, signup_date DATE)
Hints
- Run DQ on the deduped dataset, not on raw staging.
- Use a LEFT JOIN to detect missing dimension keys.