Data Engineer interview questions
100 real questions with model answers and explanations for Senior candidates.
See a Data Engineer resume example →Practice with flashcards
Spaced repetition · Hunter Pass
Questions
I would use one Kafka intake and separate Flink and Spark paths into an Iceberg lakehouse.
- Kafka gets 384 partitions at roughly 6 MB/s each, keyed by account_id, with 24-hour retention and replication factor 3 across zones.
- Flink commits dashboard tables to Iceberg every 30 seconds, while Spark compacts and rebuilds the certified daily facts before 06:00 UTC.
- Raw Parquet stays immutable for 90 days, so both paths can be replayed without maintaining a second ingestion contract.
Why interviewers ask this: The interviewer is testing whether the candidate can separate latency classes while preserving one replayable source of truth.
I would build a partition-pruned Spark job and size it from measured throughput rather than node count folklore.
- A 2-hour window requires at least 2.8 GB/s of useful scan throughput, so I would load-test 5% of the data and provision 30% headroom.
- Iceberg manifests select only the processing date, and Spark writes 512 MB Parquet files with 128 MB row groups to limit listing and shuffle overhead.
- The first run uses 80 r7gd.4xlarge workers with dynamic allocation; spot capacity can cover 60% because failed stages are retryable.
Why interviewers ask this: The interviewer is evaluating capacity math, physical layout choices, and a concrete reliability trade-off.
I would keep alert computation in Flink and retain Kafka long enough to make replay an ordinary operation.
- Kafka uses device_id as the key across 192 partitions and keeps 72 hours, which covers the 48-hour replay plus a 24-hour safety margin.
- Flink emits provisional alerts within 10 seconds, closes event-time windows at a 2-minute watermark, and checkpoints to S3 every 30 seconds for recovery under 5 minutes.
- Alerts go to a compacted Kafka topic with alert_id as the idempotency key, while raw events land in hourly Iceberg partitions.
Why interviewers ask this: The interviewer is checking whether low latency, state recovery, and bounded replay are designed together.
I would share storage and metadata but enforce tenant isolation in both authorization and compute.
- Iceberg tables partition by event_date, not tenant_id, then sort by tenant_id and event_time to avoid 200 tiny partition trees.
- Lake Formation row filters and separate IAM roles enforce tenant_id at query time, with 20 automated negative-access tests in each release.
- The five heavy tenants get dedicated Spark queues and a 40% concurrency cap each, preventing one tenant from exhausting the shared cluster.
Why interviewers ask this: The interviewer is assessing whether the design balances physical efficiency, security, and noisy-neighbor control.
I would make the stream provisional and the daily batch the audited authority over the same immutable order events.
- Flink computes 30-second operational aggregates from Kafka and labels them provisional with the source offset range.
- Spark rebuilds the daily revenue fact from Iceberg raw events plus CDC adjustments, then reconciles totals to the payment ledger within $0.01.
- Both outputs use the same metric SQL in dbt macros, but only the daily table is exposed to finance as certified.
Why interviewers ask this: The interviewer is testing whether the candidate can combine fast feedback with an explicit correctness authority.
I would standardize managed extraction and keep transformation and storage in our own platform.
- Airbyte Cloud or Fivetran handles the 120 connectors only if a 30-day volume trial stays below the $30,000 cap, including history syncs.
- Every source lands unchanged into a source/date Iceberg table with connector cursor, extraction time, and raw payload for replay.
- Airflow creates one parameterized DAG per source class, limiting concurrent syncs to 20 so vendor APIs and warehouse slots are not saturated.
Why interviewers ask this: The interviewer is evaluating a constrained ingestion design rather than an abstract build-versus-buy opinion.
I would compute features once from event time and publish versioned online and offline views.
- Flink updates Redis or Feast online keys within the 2-second SLA and records feature_name, version, and event timestamp with every value.
- The same feature definitions write append-only Iceberg tables partitioned by event_date and bucketed by entity_id for 180-day point-in-time joins.
- A daily Spark parity job samples 1 million entities and fails publication if online and offline values differ by more than 0.1%.
Why interviewers ask this: The interviewer is checking point-in-time correctness and online-offline consistency under explicit latency and history constraints.
I would keep raw and identifiable data regional and replicate only approved aggregates globally.
- Each region has its own Kafka, object store, Iceberg catalog, and KMS keys, with EU bucket policies denying replication outside eu-west-1.
- Regional Spark jobs produce 15-minute aggregates stripped of direct identifiers and publish them through cross-region Kafka topics.
- A global Iceberg catalog reads only those aggregates, while a daily 100-row residency probe verifies that EU identifiers never appear outside the region.
Why interviewers ask this: The interviewer is testing whether residency is an architectural boundary rather than a documentation promise.
I would normalize transport, not business schemas, by putting every CDC record through Kafka and into raw Iceberg tables.
- Debezium runs one connector per database, and topics retain 7 days so the 1-minute consumers and failed sink jobs can replay independently.
- A Flink sink writes Iceberg every minute with source LSN, operation, and schema version, preserving deletes instead of flattening them away.
- Six-hour Spark jobs read Iceberg snapshots, not the databases, which removes 40 production systems from the batch critical path.
Why interviewers ask this: The interviewer is assessing whether one ingestion backbone can serve two latency classes without coupling them to source databases.
I would size from replicated bytes, partition throughput, and failure headroom before choosing broker count.
- Ingress is about 1.08 GB/s and 93 TB per day raw, so replication factor 3 requires roughly 280 TB before the 70% disk limit.
- At 10 MB/s per partition I would start near 120 partitions, then use 180 to keep traffic balanced when one broker is unavailable.
- Twelve brokers with 40 TB usable each provide 480 TB, and load testing must confirm each broker sustains about 300 MB/s during rebalance.
Why interviewers ask this: The interviewer is evaluating quantitative capacity planning and failure-aware Kafka partitioning.
I would partition by day and sort within files by customer_id and event_time.
- Daily partitions keep a 7-day query to about 840 TB before pruning and avoid one partition per highly skewed customer.
- Iceberg hidden partitioning and sort order on customer_id let manifests and Parquet statistics skip files without exposing layout columns to users.
- I would target 1 GB files and rewrite only partitions with more than 10% small files, trading compaction spend for predictable scans.
Why interviewers ask this: The interviewer is checking whether partitioning handles both query shape and extreme key skew.
I would optimize for column pruning and bounded scan units rather than maximum file size.
- Files would target 512 MB to 1 GB, while 128 MB row groups keep one decompressed working set comfortably below 16 GB.
- ZSTD level 3 is my default because it typically cuts storage and scan bytes more than Snappy at an acceptable batch CPU cost.
- I would order by the two most selective filter columns and verify with a 1 TB benchmark that Parquet min/max statistics skip at least 80% of row groups.
Why interviewers ask this: The interviewer is evaluating whether the candidate connects Parquet internals to memory, CPU, and scan cost.
I would prevent tiny files at the writers and run bounded compaction as part of the table service.
- Flink buffers each partition to 256 MB or 5 minutes, whichever comes first, instead of committing one file per task checkpoint.
- Iceberg rewrite_data_files compacts files below 64 MB into 512 MB targets every hour with a 30-minute lookback.
- The service caps compaction at 20% of cluster capacity and alerts when a partition exceeds 1,000 files, preserving query p95 without starving ingestion.
Why interviewers ask this: The interviewer is testing whether small-file control is built into steady-state operations with explicit resource limits.
I would use a star schema with an order-line fact and conformed dimensions, then materialize only the proven wide access paths.
- The fact keeps numeric measures and surrogate keys, partitioned by order_date and clustered by customer_key for the 6-year history.
- The 200 attributes stay in product, customer, and channel dimensions so a corrected label does not rewrite 10 billion fact rows.
- For the 30 dashboards, dbt builds daily aggregate tables and one selective wide mart, with freshness tests at 06:15 UTC.
Why interviewers ask this: The interviewer is assessing whether the model balances semantic consistency, update cost, and dashboard performance.
I would turn ordered CDC changes into non-overlapping validity intervals and merge only affected keys.
- The Iceberg table uses bucket(256, customer_id) and stores customer_key, valid_from, valid_to, is_current, source_lsn, and a hash of tracked attributes.
- A Flink or Spark stage sorts the 20 million daily changes by customer_id and source_lsn, dropping updates whose attribute hash did not change.
- Iceberg MERGE rewrites buckets containing changed customer keys, and a test rejects any key with overlapping intervals or more than one current row.
Why interviewers ask this: The interviewer is checking temporal correctness, CDC ordering, and write amplification at large dimension scale.
I would partition by event_date, cluster by account_id, and enforce scan budgets at the query boundary.
- require_partition_filter blocks full-table scans, while authorized dashboard views default to 30 days and cap routine reads at 180 TB before clustering.
- Clustering on account_id is retained only if dry-run bytes show at least a 50% reduction on the top 20 production queries.
- Custom query quotas total 120 TiB per day, about $22,500 monthly at $6.25 per TiB, leaving 10% of the $25,000 budget for exceptions.
Why interviewers ask this: The interviewer is evaluating whether physical design and cost governance are tied to measured query behavior.
I would choose Iceberg because the engine mix and lock-in constraint matter more than Spark-only convenience.
- Iceberg has first-class Spark, Trino, and Flink connectors, so all three engines can commit through the same open catalog.
- I would benchmark 10 representative MERGE and scan workloads on 50 TB before committing, because connector maturity differs by engine version.
- The trade-off is giving up some Delta-specific Databricks optimizations, which is acceptable for a mandated 5-year multi-engine path.
Why interviewers ask this: The interviewer is testing a table-format decision against explicit engines and organizational constraints.
I would use field-ID-aware schema evolution and a two-release compatibility window.
- Iceberg renames preserve column field IDs, so readers using IDs avoid a 40 TB rewrite; I would still test every engine for name-based behavior.
- Release 1 exposes old and new views for 14 days and logs which of the 60 jobs still reads an old name.
- Release 2 removes aliases only after usage reaches zero, with an atomic view swap that stays inside the 5-minute outage budget.
Why interviewers ask this: The interviewer is assessing safe schema evolution across heterogeneous consumers without unnecessary data rewrites.
I would derive the first cluster from a throughput test and make shuffle reduction the main cost lever.
- The target is 5.6 GB/s useful throughput, so I would benchmark a 1.5 TB sample and start with 100 r7gd.4xlarge workers plus 25% headroom.
- Adaptive Query Execution uses 256 MB advisory partitions, and map-side partial aggregation cuts the 30 TB shuffle before the wide stage.
- At roughly $1.20 per worker-hour, a 125-worker 90-minute run is about $225 before service markup, leaving ample room under $2,000 for retries.
Why interviewers ask this: The interviewer is checking throughput math, Spark execution knowledge, and cost verification.
I would isolate the hot key and use a different join strategy for it instead of salting every row.
- Spark metrics first confirm the 35% key and its two straggler partitions; the remaining 65% uses an ordinary sort-merge join.
- The hot customer is split into 64 deterministic salt buckets, and its dimension row is expanded only 64 times, not across the full 600 GB table.
- AQE skew splitting stays enabled at a 256 MB threshold, and a 5% production sample must show p95 task time below 3 minutes before the full run.
Why interviewers ask this: The interviewer is evaluating targeted skew handling under a fixed runtime rather than generic repartitioning.
Locked questions
- 21
A Spark pipeline joins 8 TB of events to a 6 GB reference table on executors with 16 GB memory; decide whether to broadcast and state the safeguards.
joinsmemoryci-cd - 22
Process a 4 TB daily fact increment when 8% of events arrive up to 72 hours late, while corrected dashboards must appear within 30 minutes.
concurrency - 23
A Flink fraud job holds 15 TB of keyed state, must checkpoint every 60 seconds, and has a 10-minute recovery objective; design state and checkpointing.
design - 24
Join clicks and purchases at 800,000 events per second, where a purchase may follow a click by 6 hours and Flink state cannot exceed 8 TB.
joins - 25
Choose Spark Structured Streaming or Flink for 300,000 events per second, p99 latency under 3 seconds, 6-hour keyed sessions, and exactly-once Iceberg writes.
sessionslatencystreaming - 26
Design Airflow for 2,000 DAGs and 100,000 task runs per day, with scheduler failover under 2 minutes and p95 task-queue delay under 30 seconds.
designdata-structuresairflow - 27
Backfill 730 daily partitions totaling 400 TB in 5 days while live pipelines keep 70% of a 1,000-core Spark pool.
partitioningbackfillci-cd - 28
Design an idempotent daily pipeline that reads 2 TB, writes 12 partitions, and may be retried 3 times after any task boundary.
partitioningidempotencydesign - 29
An Airflow DAG must process 50,000 customer files every hour, but the executor allows only 2,000 concurrent tasks and scheduler delay must stay below 20 seconds.
concurrencyairflow - 30
Fifty producer DAGs feed 120 consumer DAGs, outputs arrive between 02:00 and 05:00 UTC, and consumers must start within 5 minutes of real readiness.
health-checks - 31
An Airflow task calls a billing API for 10 million rows and may retry 4 times, but duplicate charges are forbidden.
resilienceapiairflow - 32
Airflow has 300 critical tasks due by 07:00 and 20,000 noncritical tasks, with only 5,000 worker slots during the 05:00 peak; design admission control.
schedulingdesignairflow - 33
Implement data contracts for 80 producing services and 300 consuming pipelines, with incompatible schema changes blocked before deployment and review under 10 minutes.
deploymentci-cdschema - 34
Design data quality coverage for 1,000 tables when only 100 are business-critical and the validation cluster budget is 500 core-hours per day.
qualitydesigncoverage - 35
A revenue table must be fresh by 06:15 UTC on 99.9% of days and complete within 0.05%; define the SLA and its enforcement path.
- 36
A shared Kafka topic has 45 consumers in 6 languages; schemas change weekly and no consumer may break without 30 days of notice.
schemakafka - 37
Validate a daily 4-billion-row revenue fact against 3 payment processors, with a maximum $0.01 discrepancy per transaction and a 20-minute test budget.
transactionsvalidationconcurrency - 38
A dbt project has 400 models, but pull-request validation must finish in 10 minutes and production data is 200 TB.
validationdbt - 39
A pipeline receives 2 billion records per day and up to 0.2% may fail validation; good records must publish within 15 minutes and bad records must remain replayable for 30 days.
validationci-cd - 40
Capture 20,000 PostgreSQL transactions per second from a 12 TB database into Iceberg with p95 lag under 60 seconds and no source downtime.
databasetransactionspostgres - 41
Bootstrap CDC from a 15 TB PostgreSQL database while writes continue at 5 TB per day and retained WAL may not exceed 1 TB.
databasepostgres - 42
Replicate 60 PostgreSQL tables into a lakehouse for 5 years, preserving deletes and 2 schema changes per week without rewriting more than 1 TB per change.
postgresschemareplication - 43
Replicate orders and order_items from PostgreSQL at 50,000 row changes per second, and guarantee consumers never see an item before its order for more than 5 seconds.
postgresreplication - 44
A read replica feeds analytics at 30,000 writes per second; production requires replication lag below 20 seconds and analytics may consume at most 25% of replica I/O.
replication - 45
A CDC stream delivers duplicates after failover, but a 6-billion-row Iceberg current-state table must contain exactly one version per primary key within 2 minutes.
primary-keys - 46
Events arrive at 600,000 per second; 95% are within 2 minutes, 4.9% within 30 minutes, and 0.1% within 24 hours, while dashboards may revise for only 1 hour.
- 47
Deduplicate 1 million events per second using event_id, with duplicates possible for 7 days and Flink state capped at 12 TB.
queries - 48
Design exactly-once processing from Kafka through Flink into Iceberg at 400,000 events per second, with recovery under 8 minutes and no duplicate business keys.
kafkadesignconcurrency - 49
A 1 PB Iceberg table serves 2,000 queries per day, but the monthly engine bill must fall from $180,000 to $90,000 while p95 stays below 12 seconds.
queries - 50
A platform processes 300 TB per day; 10% needs results in 1 minute and 90% within 6 hours, with annual compute spend capped at $4 million.
concurrency - 51
At 04:00, a daily Spark pipeline is 3 hours behind its 06:00 freshness SLA after input grew from 8 TB to 13 TB; what do you do?
ci-cd - 52
An Airflow DAG has been waiting 70 minutes for one of 240 hourly source partitions, and the consumer SLA expires in 20 minutes; how do you respond?
partitioningairflow - 53
After a release, a Spark aggregation that normally takes 95 minutes is projected to take 6 hours and miss its SLA by 4 hours; walk through your actions?
aggregation - 54
An S3-backed ingestion job suddenly drops from 9 GB/s to 2 GB/s and is 2 hours behind while 18,000 new files arrive each minute; what do you check and change?
aws - 55
A vendor API cuts your quota from 20,000 to 5,000 requests per hour at noon, putting a 4-hour ingestion SLA at risk for 60 customers; what do you do?
procurementapi - 56
At 05:20, a dbt model misses its 06:00 SLA because an upstream team added 14 columns and runtime rose from 8 to 55 minutes; how do you recover?
schemadbt - 57
A Kubernetes Spark job has restarted 38 executors for OOM in 25 minutes and will miss a 2-hour SLA; what do you do before raising memory?
memorykubernetes - 58
One tenant grows from 4% to 46% of rows overnight, leaving 12 Spark tasks running 90 minutes after 3,000 others finish; how do you recover and prevent recurrence?
- 59
A 7 TB Spark read fails on 23 corrupt Parquet files after 96% of processing, with a 45-minute deadline; what do you do?
estimationconcurrency - 60
Twenty Flink writers produce 140 Iceberg commit conflicts in 10 minutes, freshness is already 35 minutes over a 15-minute SLA; how do you stabilize it?
- 61
A checkout dashboard jumps 18% at 09:00, and a quality check finds 42 million duplicate orders in the last 6 hours; what do you do?
- 62
Finance reports 7% fewer invoice rows than the billing system two hours before month close; how do you scope and repair the incident?
incidentssystem-design - 63
Revenue is 4.6% too high across 11 executive dashboards after a dbt deployment; what do you do in the first 30 minutes?
deploymentdbt - 64
Nulls in customer_country rise from 0.2% to 12% in one hour, but blocking publication would delay 40 dashboards by 3 hours; what is your decision?
- 65
A dimension join turns 900 million fact rows into 2.7 billion and inflates six metrics 3 times; how do you debug and roll back?
monitoringrollbackjoins - 66
After a daylight-saving change, hourly demand is shifted by 2 hours for 9 million European events and staffing decisions start in 90 minutes; what do you do?
- 67
A product dashboard is wrong by 9% because a customer dimension stopped updating 36 hours ago while every Airflow task stayed green; how do you respond?
airflow - 68
A CDC sink omitted deletes for 14 hours, leaving 2.4 million canceled accounts active in analytics; what do you do?
- 69
A statistical quality check blocks a 5 TB batch 25 minutes before SLA, but the 22% distribution shift may be a real promotion; what do you do?
distributionsbatch - 70
Sales says a dashboard is 12% low, finance says it is correct, and the board deck is due in 3 hours; how do you decide what ships?
- 71
A backfill accidentally overwrites 180 million valid rows across 12 daily Iceberg partitions; how do you repair them?
partitioningbackfill - 72
A retry appended a 1.2 billion-row backfill twice, doubling 45 days of customer events; what do you do without rewriting a 9 PB table?
resiliencebackfill - 73
A 400 TB backfill fails after 63% when spot capacity disappears, and it must finish in 36 hours without restarting completed work; what do you do?
capacitybackfill - 74
You discover that a 90-day backfill used the wrong tax rule and changed 640 million rows consumed by finance; how do you unwind it?
backfill - 75
A backfill consumes 82% of a shared Spark cluster and makes 17 production pipelines miss freshness by 50 minutes; what do you do?
backfillci-cd - 76
A Kafka consumer group falls 12 million messages behind in 18 minutes while producer throughput remains 700,000 events per second; how do you diagnose it?
kafkathroughput - 77
Kafka lag reaches 9 million on 3 of 96 partitions while the other 93 stay near zero; what do you do under a 20-minute alert SLA?
partitioningkafkaalerting - 78
Flink checkpoints grow from 40 seconds to 7 minutes, fail 6 times, and threaten a 10-minute recovery objective; what do you do?
- 79
A consumer group rebalances 28 times in 15 minutes, throughput drops 65%, and 4 million events queue up; how do you stabilize it?
throughputdata-structures - 80
A producer retry storm creates 6.5 million duplicate payment events in 12 minutes; dashboards must recover within 30 minutes; what do you do?
resilience - 81
Eight percent of ad conversions arrive 5 hours late after a mobile outage, but campaign dashboards permit revisions for only 2 hours; what do you do?
- 82
CDC updates for 14 million orders arrive out of order after failover, and 3% of current-state rows regress to older statuses; how do you repair them?
- 83
A schema bug sends 2.2 million events to a dead-letter topic in 25 minutes and its 3-hour retention is nearly full; what do you do?
retentionschema - 84
An event producer changes price from cents as int to dollars as string, breaking 17 consumers and building 5 million messages of lag; how do you restore service without downtime?
- 85
Source counters show 1.8 billion events, but Kafka and the lake contain 1.76 billion after a broker incident; how do you investigate the missing 40 million?
incidentskafka - 86
Snowflake spend jumps from $85,000 to $260,000 in 9 days while query volume rises only 12%; what do you do this week?
queriessnowflake - 87
A new BigQuery dashboard scans 900 TB per day and adds $5,600 daily cost for 1,200 views; how do you fix it under a 24-hour deadline?
estimationbigquery - 88
An EMR Spark pipeline costs 3.4 times more after autoscaling expands from 60 to 240 workers but runtime improves only 8%; what do you change?
scalingci-cd - 89
Iceberg storage grows 140 TB in one week although new business data is only 18 TB; how do you stop the cost explosion safely?
- 90
Kafka cross-region egress rises from $12,000 to $74,000 per month after 6 consumers are deployed, and finance demands a 40% cut in 2 weeks; what do you do?
kafkadeployment - 91
A rename of 4 columns breaks 27 production jobs, and you have 45 minutes to restore them with no consumer deployment window; what do you do?
schemadeployment - 92
An order_id changes from 32-bit to 64-bit, 0.4% of new values overflow in 9 consumers, and lag reaches 3 million; how do you repair it?
- 93
A Protobuf producer makes a field required, 11 older consumers reject 2.8 million messages, and the incident SLA is 30 minutes; what do you do?
incidents - 94
A dbt column keeps the same name but changes from gross to net revenue, shifting 23 dashboards by 6%; the release is already live; what do you do?
schemadbt - 95
A Parquet writer upgrade changes timestamp encoding and 8 of 35 Trino jobs fail over 6 TB of new files; how do you restore zero-downtime reads?
- 96
During a 45-minute freshness incident, a mid-level engineer wants to scale Spark from 100 to 400 workers without checking a single stuck task; how do you mentor in the moment?
mentoringproblem-solvingincidents - 97
A junior engineer's backfill duplicates 320 million rows 2 hours before a finance deadline; how do you handle the repair and the engineer?
soft-skillsestimationbackfill - 98
Product demands an unverified 48-hour backfill for a launch in 6 hours, but a 2% sample differs from production by 7%; what do you decide?
backfill - 99
Your on-call engineer has handled 14 pages in 3 hours and misses a second pipeline failure that puts data 90 minutes stale; what do you do?
on-callci-cd - 100
A data engineer has caused 3 schema incidents in 6 weeks across 22 consumers despite ordinary code reviews; what concrete mentorship plan do you use?
mentoringcode-reviewincidents