Interview questions

Data engineer interview questions for senior hires (2026)

Practical questions on pipelines, orchestration, dbt, Spark and lakehouse tables, with what separates a senior data engineer from a mid-level one.

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.

How to use these questions

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.

Fundamentals

What is the difference between ETL and ELT, and when would you still choose ETL?

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.

A nightly job was rerun after a failure and the revenue table now has duplicate rows. How do you make the pipeline idempotent?

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.

Explain fact and dimension tables. What is the grain of a fact table and why does it matter?

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.

What are slowly changing dimensions? When do you use Type 1 versus Type 2?

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.

Why are columnar formats like Parquet faster for analytics than CSV or JSON?

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.

A table has 400,000 files of a few kilobytes each and queries are slow. What happened and how do you fix it?

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.

How do you choose between batch and streaming for a new pipeline?

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.

Intermediate

An Airflow DAG uses datetime.now() to decide which day to load. What goes wrong and how do you fix it?

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.

How do dbt incremental models work, and how do you handle late-arriving updates?

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.

When would you choose Dagster over Airflow, or the other way around?

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.

Compare log-based change data capture with query-based extraction from a Postgres database.

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.

One Spark task runs for 40 minutes while the other 199 finish in a minute. How do you diagnose and fix it?

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.

When does Spark use a broadcast join, and when is it a mistake?

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.

Where do you put data quality checks, and which ones should block a pipeline?

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.

Lakehouse and table formats

What do Apache Iceberg and Delta Lake add on top of plain Parquet files in object storage?

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.

What is hidden partitioning in Iceberg, and how does partition evolution work?

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.

What maintenance does a lakehouse table need, and what happens if nobody runs it?

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.

Explain copy-on-write versus merge-on-read for row-level updates.

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.

Two jobs write to the same Iceberg table at once and one fails with a commit conflict. Why, and how should you design around it?

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.

Senior and architecture

Design ingestion from 50 operational Postgres databases into a lakehouse with 15-minute freshness.

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.

An upstream team renamed a column and broke six dashboards overnight. How do you stop this from happening again?

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.

The warehouse bill doubled in one quarter. How do you investigate?

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.

How do you get exactly-once results from Kafka into a lakehouse table?

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.

Business logic changed and two years of history must be recomputed. How do you backfill without disrupting consumers?

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.

A customer requests deletion under GDPR. How do you remove their data from an immutable lakehouse?

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.

When would you buy a managed ingestion tool instead of building connectors?

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.

Red flags to watch for

A practical exercise

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.

Hire senior data engineers vetted with these questions

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.

FAQ

How many of these questions should I use in one interview?

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.

How is a data engineer interview different from a data scientist interview?

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.

Should I test SQL separately?

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.

Explore Ryz Labs

Staff augmentationDedicated development teamsAI pod teamsForward deployed engineersNearshore software developmentAI engineering teamsHire engineers by roleRyz Labs vs competitorsAlternatives guidesBuyer guidesCase studiesHow we vet engineers
Ryz Labs

Senior engineers in your time zone. AI pod teams that ship.

Tell us who you need. You'll get a scoped plan, a price and the names of the people who would do the work.

Start a conversation →