Design an idempotent churn ETL pipeline
Company: Intuit
Role: Data Scientist
Category: Data Manipulation (SQL/Python)
Difficulty: medium
Interview Round: HR Screen
You must build a daily pipeline that produces month-end churn metrics (logo churn, gross revenue churn, net revenue retention) from streaming subscription events with late arrivals (up to T+3 days). Requirements: idempotent runs, backfills for the past 12 months, slowly changing dimensions (SCD Type 2) for plan changes, data quality checks, and reproducibility.
1) Outline the DAG (tasks, dependencies) from raw events to curated snapshots to metrics. Specify partitioning, clustering, and keys.
2) Provide SQL for an upsert/merge that constructs a monthly snapshot table from event-level changes, correctly handling late-arriving cancels/reactivations and preventing double-application on reruns.
3) Describe a strategy to recompute only affected months when late data arrives (e.g., incremental backfill windows, watermarks) and how you’d validate the recomputation via invariants.
4) Show how you’d version metric definitions so historical reports remain interpretable when definitions evolve (e.g., semantic layer with versioned views).
Overview: This question evaluates a candidate's competence in designing robust ETL pipelines, covering idempotent processing, handling late-arriving streaming subscription events, SCD Type 2 slowly changing dimensions, incremental backfills, data quality checks, upsert/merge semantics, and metric versioning.
Read the full Intuit Data Scientist interview experience this question came from
Idempotent monthly subscription snapshot with late events (MERGE)
You maintain a streaming event log of subscription activity (`subscription_events`) and a slowly-changing plan dimension in SCD Type 2 form (`plans_dim_scd2`). You must build the content of a **monthly subscription snapshot** as of pipeline run date **2025-06-01**, producing exactly **one row per `(account_id, month_end_date)`** for the five month-end dates **2025-01-31, 2025-02-28, 2025-03-31, 2025-04-30, 2025-05-31** and for every account that appears in `subscription_events` (3 accounts x 5 months = 15 rows).
**Active rule.** For a given `month_end_date`, find the **latest** event for that account whose `event_effective_date <= month_end_date` **and** `event_ingested_at <= DATE '2025-06-01'` (this models late-arriving events that are only counted once ingested by the run date). If two events tie on `event_effective_date`, break the tie by the larger `event_id`. The account is **active** at that month-end if that latest event's `event_type` is **not** `'CANCEL'`. If the account has no qualifying event at all for that month-end, it is **inactive**.
**MRR rule (SCD2).** When active, look up the `mrr` and `plan_sk` from `plans_dim_scd2` for the latest event's `plan_id`, choosing the SCD2 row whose `[effective_start_date, effective_end_date]` interval contains `month_end_date` (i.e. `month_end_date BETWEEN effective_start_date AND effective_end_date`).
**Output columns** (in this exact order): `account_id`, `month_end_date`, `plan_id`, `plan_sk`, `mrr`, `is_active`, `last_event_id`, where:
- `is_active` is the boolean active flag computed above.
- `last_event_id` is the `event_id` of that latest qualifying event (or `NULL` if none).
- When **inactive** (latest event is a `CANCEL`, or no qualifying event): `plan_id` and `plan_sk` must be `NULL` and `mrr` must be `0.00`.
Order the result by `account_id`, then `month_end_date` ascending.
Write a single `SELECT` that produces these 15 snapshot rows. (In production this result set is the `USING` source of an idempotent `MERGE` keyed on `(account_id, month_end_date)` so reruns update existing rows and insert missing ones rather than duplicating; here you only need to produce the snapshot rows.)
Tables
subscription_events(event_id INT, account_id INT, event_type VARCHAR(20), event_effective_date DATE, event_ingested_at DATE, plan_id INT, quantity INT)
plans_dim_scd2(plan_sk INT, plan_id INT, plan_name VARCHAR(50), mrr DECIMAL(10,2), effective_start_date DATE, effective_end_date DATE, is_current BOOLEAN)
monthly_subscription_snapshot(account_id INT, month_end_date DATE, plan_id INT, plan_sk INT, mrr DECIMAL(10,2), is_active BOOLEAN, last_event_id INT)
Hints
- Cross-join all accounts with the five hard-coded month-end dates first, so accounts with no event still get an (inactive) row per month.
- Per (account, month_end), rank qualifying events with ROW_NUMBER() OVER (ORDER BY event_effective_date DESC, event_id DESC) and keep rn = 1 as the 'latest' event.
Compute monthly churn and net revenue retention from snapshots
Write a PostgreSQL query. You now have a curated monthly subscription snapshot with one row per `(account_id, month_end_date)` and columns indicating whether the account is active and its MRR at that month-end.
Using the `monthly_subscription_snapshot` table defined below, write a SQL query that produces **month-end churn metrics** for each month from **2025-02-28** through **2025-05-31**:
- `starting_accounts`: number of accounts that were active at the **previous** month-end.
- `ending_accounts`: number of accounts active at the **current** month-end.
- `logo_churn_accounts`: number of accounts that were active at the previous month-end but inactive at the current month-end.
- `starting_mrr`: sum of MRR at the previous month-end across accounts that were active at that previous month-end.
- `contraction_mrr`: for accounts that were active at the previous month-end, the sum of MRR decreases (including going to zero on cancel) from previous to current month-end.
- `expansion_mrr`: for accounts that were active at the previous month-end, the sum of MRR increases from previous to current month-end.
- `gross_revenue_churn`: `contraction_mrr / starting_mrr`.
- `net_revenue_retention`: `(starting_mrr - contraction_mrr + expansion_mrr) / starting_mrr`.
Notes:
- Use only accounts that were active at the **previous** month-end when computing `starting_mrr`, `contraction_mrr`, and `expansion_mrr` (i.e., exclude new or reactivated accounts from the expansion base).
- Do not output metrics for the first month in the series (2025-01-31) because there is no previous month to compare.
- Assume the snapshot already reflects late-arriving events as of 2025-06-01.
Tables
monthly_subscription_snapshot(account_id INT, month_end_date DATE, plan_id INT, plan_sk INT, mrr DECIMAL(10,2), is_active BOOLEAN, last_event_id INT)
Hints
- Use window functions (LAG) over (account_id, ordered by month_end_date) to bring previous-month activity and MRR onto each row.
- Aggregate across accounts by month_end_date, computing contractions and expansions as positive differences in MRR where prev_is_active is true.
Identify months to incrementally backfill on late-arriving events
Your pipeline processes subscription events daily and keeps a watermark of the last fully processed ingestion date for the churn snapshot. Late-arriving events (up to T+3 days after their effective date) may arrive with `event_ingested_at` **after** the watermark.
To avoid recomputing all history, you want to recompute only the **affected month-end snapshots** and downstream churn metrics for the last 12 months (from **2024-06-30** to **2025-05-31**), given a pipeline run date of **2025-06-01**.
Using the tables below, write a SQL query that returns the list of `month_end_date` values that must be recomputed when new late events arrive after the current watermark for the `churn_snapshot` pipeline.
Logic to implement:
- Read `last_processed_ingested_at` from `pipeline_watermarks` for `pipeline_name = 'churn_snapshot'`.
- Identify all events in `subscription_events` with `event_ingested_at > last_processed_ingested_at` and `event_effective_date` between **2024-06-01** and **2025-05-31**.
- For those late events, compute the **month-end date** of their `event_effective_date`.
- Let `first_affected_month_end` be the **earliest** such month-end date.
- Return all month-end dates in the 12-month window (2024-06-30 to 2025-05-31) that are **greater than or equal to** `first_affected_month_end`, ordered ascending.
Assume your SQL engine supports `DATE_TRUNC('month', ...)` and `INTERVAL` arithmetic similar to PostgreSQL.
Tables
subscription_events(event_id INT, account_id INT, event_type VARCHAR(20), event_effective_date DATE, event_ingested_at DATE, plan_id INT, quantity INT)
pipeline_watermarks(pipeline_name VARCHAR(50), last_processed_ingested_at DATE)
Hints
- First compute late events by comparing event_ingested_at to the watermark, then convert their effective dates to month-end dates.
- Use the earliest affected month-end as a lower bound and return all month-ends from a 12-month spine greater than or equal to that date.
Join versioned metric definitions to computed churn metrics
You want churn metrics to remain interpretable even as definitions evolve over time. You store versioned metric definitions in a semantic-layer table and record which version was used for each metric run.
Using the tables below, write a SQL query that returns, for every row in `computed_metrics`:
- `metric_run_id`, `metric_name`, `metric_version`, `month_end_date`, `value`, and
- the corresponding `definition_sql` from `metric_definitions`.
The join must ensure that each metric run is associated with the **exact definition version** that was used at computation time, so that later changes to the definition do not retroactively alter historical interpretations.
Tables
metric_definitions(metric_name VARCHAR(50), version INT, definition_sql VARCHAR(4000), effective_start_date DATE, effective_end_date DATE)
computed_metrics(metric_run_id INT, metric_name VARCHAR(50), metric_version INT, month_end_date DATE, value DECIMAL(10,4), computed_at TIMESTAMP)
Hints
- Join the two tables on both metric_name and the version number recorded with the metric run.
- You do not need to use effective date ranges here because the metric_version column already encodes which version was used.