Interview framing
You are designing a batch ETL system in a data engineering interview. Time box is about 45 minutes. The interviewer cares about the clock, correctness under retries, and whether BI can trust what they see at 6 AM. They do not care whether you memorized every Airflow operator name.
The prompt:
Design a batch pipeline that processes about 2 million orders per day for an e-commerce company and lands curated warehouse tables by 6:00 AM local time for daily reports.
Background from first principles
Batch ETL means you process data in scheduled chunks (often one business day), not as each row arrives. That is different from streaming. At 2M orders/day you are not in firehose territory. If an order JSON is about 5 KB, raw orders alone are near 10 GB/day. A modest Spark or warehouse ELT job can transform that in minutes. The hard parts are dependencies, the definition of yesterday, idempotency, and atomic publish.
Think of three promises:
1. Completeness: for business date D, the warehouse has all orders that belong to D within an agreed late window. 2. Correctness: rerunning the job does not double-count. 3. Visibility: either date D is published and safe for BI, or it is clearly not published and someone is alerted.
Finance treating 6 AM as hard is normal. A missing day with an alert beats wrong revenue that leadership trusts.
Also clarify extract versus load versus transform. Extract copies source data into a raw zone you control. Transform builds analytics shapes. Load or publish makes those shapes visible to BI. Mixing those stages without clear boundaries is how half-written mornings happen.
Expanded problem and constraints
- Volume: 2M orders/day; assume about 5 line items average so about 10M item rows/day.
- SLA: critical marts ready by 06:00; page by 05:00 if the run is behind.
- Sources: orders, order_items, payments, refunds for the completed business day.
- Load pattern: prefer incremental by date with ability to re-run a single day safely.
- Quality: checks must pass before BI is allowed to read the new day.
- Publish: atomic. No half-visible days.
- Cost: ephemeral compute / spot where reasonable; partition by order_date as history grows.
What good looks like
You define business day and timezone first. You estimate GB and runtime so you do not overbuild a streaming stack. You draw raw zone to transform to staging to atomic publish. You explain what happens when a job fails at 4 AM. You name idempotent keys and a success signal BI can trust.
Strong candidates also name which tables are on the critical path and which can finish later without blocking finance.
Clarifying questions
- Source of truth: OLTP primary, read replica, CDC topic, or nightly dump?
- Timezone for business day, and how late orders after midnight are handled?
- Which reports sit on the critical 6 AM path vs can land later?
- SCD needs for customers and products?
- Full refresh vs incremental by date?
- Who owns schema changes on source tables?
- Are refunds same-day or can they adjust prior days?
Scale and estimation prompts
- 2e6 orders x 5 KB is about 10 GB/day orders payload (plus items, payments).
- Transform runtime: minutes to low tens of minutes on a modest cluster if designed well. The risk is waiting on source readiness and retries, not raw compute.
- Batch window: if day closes late (stores until 2 AM), extract cannot start at midnight. Count remaining hours to 6 AM.
- Historical growth: partition facts by order_date; avoid rewriting all history nightly.
- Rough cost: object storage cents per GB; warehouse storage grows with years of partitions; compute is the dial you turn for retries.
Out of scope for v1
- Sub-minute freshness for every metric (that is a later CDC or streaming add-on).
- Rebuilding the entire multi-year history every night.
- Perfect SCD2 for every attribute on day one.
- Multi-region active-active warehouses.
- Replacing the finance mart with a real-time dashboard as system of record.