Skip to content

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

design

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.

designbatchci-cd

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.

alertingci-cddesign

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.

designstreamingbatch

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.

cloud

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.

designci-cd

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.

design

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.

designstreamingbatch

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.

retentionreplicationkafka

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.

queries

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.

queriesschemamodeling

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.

queriesdesign

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.

queries

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.

design

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.

queriesbigquery

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.

procurementaws

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.

schema

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.

aggregation

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.

joinsmodeling

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