These 26 data engineer interview questions cover what senior data engineers own: idempotent batch and streaming pipelines, orchestration in Airflow and Dagster, dbt modeling, Spark tuning, lakehouse tables on Iceberg and Delta Lake, change data capture and data quality. At Ryz, senior data engineers are evaluated the same way: recruiters source people who have run pipelines in production, candidates complete structured NTRVSTA AI interviews on design and debugging, and recruiters review every candidate before and after. The AI scores are advisory, and people make the decisions.
Pick six to eight questions per interview and match them to the level you are hiring. Fundamentals show whether someone can own a single pipeline; the lakehouse and senior sections show whether they can own a platform other teams depend on. Ask follow-ups from their history: "When did a backfill go wrong, and what changed afterward?"
Listen for reasoning about failure, cost and freshness rather than tool names. A senior data engineer explains what happens when a job reruns, data arrives late or an upstream column is renamed. Pair the conversation with the practical exercise below.
ETL transforms data before loading it into the target. ELT loads raw data first and transforms it inside the warehouse or lakehouse, usually with SQL and dbt. ELT is the default today because warehouse compute is elastic and keeping raw data makes reprocessing possible.
ETL still makes sense when data must be masked before it lands (PII, card data) or when the payload is huge and mostly discarded.
What a strong answer shows: They treat this as a decision about compliance, cost and reprocessing, not a fashion choice, and they mention keeping an immutable raw layer.
Every run should produce the same result for the same input window, no matter how many times it runs. The usual patterns are overwriting the partition the run owns, or using MERGE on a stable business key instead of a blind append.
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
(daily_orders
.write
.mode("overwrite")
.insertInto("analytics.orders_daily")) # replaces only partitions present in daily_orders
What a strong answer shows: They tie each run to a fixed data interval, never to "now", and they know that dynamic partition overwrite or MERGE is what makes retries safe.
Facts record events or measurements (an order line, a payment, a page view) with numeric measures and foreign keys. Dimensions describe the context: customer, product, date. The grain is what one row of the fact table represents, such as "one row per order line per day".
Declaring the grain prevents double counting: join a per-order shipping fee onto a per-line fact and the fee is summed once per line.
What a strong answer shows: They state grain before columns and can explain a fan-out bug they have seen in a join.
Type 1 overwrites the old value, so history is lost. Type 2 keeps a new row per change with validity columns, so facts can be joined to the attribute as it was at the time. Use Type 2 when reports must reflect history, such as sales by the customer's region at purchase time.
{% snapshot customers_snapshot %}
{{ config(
target_schema='snapshots',
unique_key='customer_id',
strategy='timestamp',
updated_at='updated_at'
) }}
select * from {{ source('crm', 'customers') }}
{% endsnapshot %}
What a strong answer shows: They know dbt snapshots add valid-from and valid-to columns, and they mention the join condition on the event timestamp that Type 2 requires.
Parquet stores values column by column with per-column compression and encoding, and keeps min and max statistics per row group. Queries read only needed columns and skip row groups their statistics rule out; CSV and JSON force a full parse.
What a strong answer shows: They mention column pruning, predicate pushdown via statistics and that sorting data on common filter columns makes those statistics useful.
This is the small files problem. It comes from streaming micro-batches, over-partitioning or too many writer tasks. Fix it by compacting to files of roughly 128 MB to 1 GB, partitioning on a coarser key and repartitioning before the write.
What a strong answer shows: They find the cause in the write path rather than only compacting after the fact.
Start from the freshness the consumer actually needs and what it costs to deliver. A finance close report needs correctness by morning, so batch is simpler. Fraud scoring may need seconds, which justifies Kafka and Flink or Spark Structured Streaming and their operational load.
What a strong answer shows: They ask who consumes the data and what decision it drives before choosing a technology.
Retries, backfills and catch-up runs all load "today" instead of the interval they were scheduled for. Use the data interval Airflow passes to every run so a rerun of last Tuesday loads last Tuesday.
import pendulum
from airflow.sdk import dag, task
@dag(schedule="@daily", start_date=pendulum.datetime(2026, 1, 1, tz="UTC"), catchup=False)
def orders_daily():
@task
def load(**context):
start = context["data_interval_start"]
end = context["data_interval_end"]
run_sql("load_orders.sql", start=start, end=end)
load()
orders_daily()
What a strong answer shows: They know Airflow 3 moved DAG authoring to the Task SDK and that logical dates and intervals, not wall-clock time, make runs reproducible.
On incremental runs is_incremental() is true and you filter to new or changed rows, which dbt merges using unique_key. Late updates are handled with a lookback window rather than a strict high-water mark.
{{ config(materialized='incremental', unique_key='order_id', incremental_strategy='merge') }}
select order_id, status, amount, updated_at
from {{ source('shop', 'orders') }}
{% if is_incremental() %}
where updated_at > (select max(updated_at) from {{ this }}) - interval '3 days'
{% endif %}
What a strong answer shows: They explain the lookback trade-off (more reprocessing versus missed updates) and when to schedule a periodic full refresh.
Airflow is task-centric with a huge operator ecosystem, and Airflow 3 added assets and better backfills. Dagster is built around software-defined assets: you declare the tables that should exist, and lineage, freshness checks and partitioned backfills follow. Dagster suits teams that think in data assets; Airflow suits many heterogeneous jobs or existing investment.
What a strong answer shows: A preference backed by migration cost, team skills and observability needs, not a feature checklist.
Query-based extraction polls with updated_at > last_run, missing hard deletes and rows whose timestamp was not updated. Log-based CDC with Debezium reads the write-ahead log through logical replication, capturing inserts, updates and deletes in commit order with low source impact.
The cost is operational: an unconsumed replication slot retains WAL and can fill the source disk.
What a strong answer shows: They raise deletes, replication slot monitoring and initial snapshots without prompting.
That pattern is data skew: one join or group-by key holds a disproportionate share of rows, often a null or a default value. Confirm via shuffle read size per task in the Spark UI. Then handle null keys separately, rely on Adaptive Query Execution's skew join handling, broadcast the smaller side or salt the hot key.
What a strong answer shows: They check the Spark UI before changing configs and know that AQE splits skewed partitions automatically in Spark 3.x but cannot fix every case.
A broadcast join ships the small table to every executor and avoids shuffling the large one. It is ideal for small lookup tables and a mistake once the "small" table grows and causes out-of-memory errors.
from pyspark.sql.functions import broadcast
enriched = events.join(broadcast(countries), on="country_code", how="left")
What a strong answer shows: They know the threshold setting exists, that hints should be revisited as data grows and that AQE can switch strategies at runtime.
Check at ingestion (schema, row counts, freshness), at transformation (uniqueness, not-null, accepted values, referential integrity via dbt tests) and at publication (reconciliation against source totals). Block on checks that make downstream numbers wrong, such as duplicate keys or a collapse in row count. Warn on softer drift.
What a strong answer shows: They distinguish severity levels, mention write-audit-publish or staging tables and route alerts to an owner.
They add a metadata layer that turns a folder of files into a table: atomic commits, snapshot isolation for readers, schema evolution by column ID, time travel to earlier snapshots and row-level updates and deletes.
What a strong answer shows: They explain that the commit is an atomic swap of a metadata pointer, which is why concurrent readers see consistent snapshots.
Iceberg partitions by transforms of columns, such as days of a timestamp, so queries filter on the timestamp itself and the engine prunes partitions. When volume grows you can change the spec without rewriting old files.
CREATE TABLE lake.analytics.events (
event_id string, user_id bigint, ts timestamp, payload string
) USING iceberg
PARTITIONED BY (days(ts), bucket(16, user_id));
ALTER TABLE lake.analytics.events REPLACE PARTITION FIELD days(ts) WITH hours(ts);
What a strong answer shows: They contrast this with Hive-style partitioning, where a query that filters on the wrong column scans everything.
Compaction to merge small files, snapshot expiry to drop old metadata and unreferenced data files, orphan file cleanup and, for Delta, OPTIMIZE and VACUUM. Without it, planning slows, storage grows and delete files pile up.
CALL lake.system.rewrite_data_files(table => 'analytics.events');
CALL lake.system.expire_snapshots(table => 'analytics.events',
older_than => TIMESTAMP '2026-09-01 00:00:00', retain_last => 10);
CALL lake.system.remove_orphan_files(table => 'analytics.events');
What a strong answer shows: They link snapshot retention to time travel and audit needs, and schedule maintenance as part of the platform rather than as a cleanup chore.
Copy-on-write rewrites every data file that contains an updated row, so writes are slow but reads are clean. Merge-on-read writes small delete or change files and applies them at read time, so writes are fast but reads do extra work until compaction. CDC tables with frequent small updates usually favor merge-on-read plus compaction.
What a strong answer shows: They choose per table based on the write and read pattern and know both Iceberg and Delta (via deletion vectors) support the faster write path.
Table formats use optimistic concurrency: each writer prepares files and tries to commit against the snapshot it started from. If another commit changed overlapping data, the commit is retried or rejected. Design so writers own disjoint partitions and compaction avoids partitions a stream is writing.
What a strong answer shows: They understand retries are safe only if the job is idempotent, and they coordinate maintenance with ingestion schedules.
Use log-based CDC (Debezium on Kafka Connect, or a managed equivalent) into Kafka topics per table, with schemas in a registry. A streaming job writes raw change events to an append-only bronze table, then a merge job applies them to current-state silver tables every few minutes. Handle the initial snapshot separately, monitor replication lag per source and publish freshness per table.
What a strong answer shows: They cover backfill, schema changes, deletes, monitoring and who gets paged, not only the happy-path arrows.
Introduce data contracts for critical sources: an agreed schema with owners, versioning and a deprecation window. Enforce them in the producer's CI and at ingestion, and keep a modeled layer between raw sources and dashboards so a rename is absorbed in one place.
What a strong answer shows: They solve the people problem (ownership and notice) as well as the technical one.
Attribute cost from query history by warehouse, user, dbt model and schedule. A few culprits usually dominate: full refreshes that should be incremental, dashboards refreshing every minute, unpartitioned scans, idle oversized compute. Fix those, add query tags and set budget alerts.
What a strong answer shows: They measure before optimizing and build attribution so cost has an owner.
True end-to-end exactly-once is about effects, not delivery. Make the sink idempotent: commit offsets atomically with the table commit (Spark Structured Streaming checkpoints or Flink two-phase commit sinks) and MERGE on an event ID so replays do not duplicate.
What a strong answer shows: They distinguish delivery guarantees from outcome guarantees and name where deduplication happens.
Build the new version alongside the old one, backfill it in partition batches with controlled concurrency, then validate row counts and key metrics against the old table. Switch consumers atomically through a view or table swap, and keep the old version until sign-off.
What a strong answer shows: They plan validation and rollback, communicate the metric change and protect production compute during the backfill.
Row-level deletes on Iceberg or Delta remove the rows logically, but the bytes stay in older files until snapshots are expired and files are rewritten or vacuumed. A complete process covers all layers, derived tables and backups, and expires snapshots within the required window. Some teams avoid the problem by tokenizing PII and deleting the key.
What a strong answer shows: They know time travel conflicts with erasure and design retention to satisfy both.
Buy (Fivetran, Airbyte Cloud or similar) for standard SaaS sources where a vendor maintains API changes. Build when volume makes usage pricing expensive, the source is internal or you need control over sensitive data.
What a strong answer shows: They count maintenance hours as a real cost and do not build out of pride.
Give a three-hour take-home: a CSV snapshot of an orders table plus a JSON-lines file of CDC events (inserts, updates, deletes, some duplicated or out of order). The candidate builds a pipeline in Python with DuckDB or PySpark, optionally dbt, producing a current-state orders table and a daily revenue fact, plus a README on reruns, late events and schema changes.
Ryz can introduce senior data engineers who have built CDC pipelines, dbt projects and lakehouse platforms in production. They are the top 1% of the candidates we interview, they work on your team, repos and standups, and they keep hours within ±1h of US time zones. You can read how we screen in our vetting process, and use our data engineer job description as a starting point for the role.
Six to eight in a 60-minute interview. Use two fundamentals to calibrate, then go deep on three or four that match your stack. Depth beats breadth.
A data engineer interview focuses on reliability: idempotency, orchestration, storage layout, CDC and cost. A data scientist interview focuses on statistics, experiment design and modeling. Both need SQL, but data engineer questions should center on what happens when the pipeline fails.
Yes, briefly. A 20-minute SQL exercise with window functions, deduplication and a join that can fan out catches gaps quickly. Seniors should also explain how it performs on billions of rows.
Questions we didn't answer? Email info@ryzlabs.com.