Data engineering interview questions in 2026 test whether you can build pipelines that produce correct, timely, affordable data and keep producing it when sources change, jobs retry and events arrive late. Reciting what ETL stands for is not enough. This guide collects 65 high-value questions with model answers. It covers fundamentals, data modelling, SQL, Spark, lakehouse table formats, Kafka and streaming, orchestration, data quality and contracts, lineage, cost, data engineering for AI, privacy, and 13 real-world scenarios.
How to use this guide
Questions are grouped by topic and numbered continuously. Topics run from foundational to advanced, and the scenario section at the end is where senior interviews are usually won or lost. What interviewers commonly probe at each level:
- Freshers and early-career engineers: OLTP versus OLAP, star schemas, window functions, deduplication in SQL, file formats, and a clear explanation of what a Spark shuffle is.
- Mid-level engineers: incremental loads, SCD Type 2, partitioning and skew, broadcast joins, Kafka consumer groups, idempotent orchestration and dbt tests.
- Senior and lead engineers: table-format trade-offs, exactly-once semantics from end to end, data contracts, lineage, cost control, governance and privacy, and pipelines that feed AI systems.
Answer each question aloud before reading the model answer. SQL is written in a generic ANSI style; check dialect details (for example QUALIFY or MERGE syntax) for the warehouse you are interviewing on.
- Fundamentals (Q1βQ8)
- Data modelling (Q9βQ13)
- SQL deep-dive (Q14βQ19)
- Apache Spark (Q20βQ27)
- Lakehouse and table formats (Q28βQ33)
- Kafka and streaming (Q34βQ39)
- Orchestration and transformation (Q40βQ42)
- Data quality, contracts and lineage (Q43βQ45)
- Cloud warehouses, performance and cost (Q46βQ47)
- Data engineering for AI (Q48βQ50)
- Privacy and governance (Q51βQ52)
- Real-world scenario questions (Q53βQ65)
- Key takeaways
- Interview preparation checklist
- FAQ
Fundamentals
1. What is the difference between OLTP and OLAP systems?
Answer: OLTP (online transaction processing) systems run the business. OLAP (online analytical processing) systems analyse it. OLTP databases such as PostgreSQL, MySQL or Oracle handle many small, concurrent reads and writes on individual rows, such as placing an order or updating a balance. They are normalised to avoid update anomalies, use row-oriented storage and indexes, and need strict ACID transactions and low latency. OLAP systems such as cloud warehouses and lakehouse engines answer questions that scan millions or billions of rows but touch only a few columns, such as revenue by region by month. They use columnar storage, compression and massively parallel execution, and their schemas are usually denormalised (star schemas) so queries need fewer joins.
Interview tip: Explain why you do not run heavy analytics on the production OLTP database: it competes with customer transactions for locks, I/O and CPU. That is why CDC and a separate analytical store exist.
2. When would you choose batch processing over streaming?
Answer: Choose batch when the business can wait. Choose streaming when the value of the data decays in seconds or minutes. Batch processes bounded datasets on a schedule (hourly, daily). It is simpler to build, test, rerun and reason about, and it is usually cheaper. Streaming processes unbounded data continuously. You need it for fraud scoring, real-time inventory, alerting or operational dashboards, and you pay for it with harder concepts: state, event time, watermarks, late data and exactly-once delivery.
Real-world example: A bank needs card-fraud features within seconds, which means streaming. Its regulatory reports are produced daily from reconciled data, which means batch.
3. What is the difference between ETL and ELT, and why did ELT become common?
Answer: In ETL you extract data, transform it on a separate engine, then load the finished result. In ELT you load raw data into the warehouse or lakehouse first and transform it there with SQL. ELT spread because cloud warehouses separated cheap storage from elastic compute, so keeping raw data and transforming it in place became affordable. Keeping the raw layer also means you can reprocess when business logic changes, and SQL-based tools such as dbt made those transformations versioned and testable. ETL still makes sense when data must be masked or filtered before it lands, for example removing personal data before it reaches a shared platform, or when the transformation is heavy non-SQL processing such as document parsing.
4. Compare a data warehouse, a data lake and a lakehouse.
Answer: A data warehouse stores structured, modelled data in a managed engine with strong SQL, ACID transactions and governance, and storage is usually proprietary. A data lake stores files of any type cheaply in object storage (Amazon S3, Azure Data Lake Storage, Google Cloud Storage). It is flexible, but plain files give you no transactions, unreliable concurrent writes and weak schema enforcement, which is how "data swamps" happen. A lakehouse adds an open table format (Delta Lake, Apache Iceberg or Apache Hudi) on top of lake files. That gives you ACID commits, schema enforcement and evolution, time travel and efficient updates, so warehouse-style workloads can run on open files that several engines can read.
Interview tip: The boundary has blurred, so talk about the trade-off (openness and multi-engine access versus a single managed experience), not about which category "wins".
5. Why are columnar formats like Parquet preferred for analytics, and where does Avro fit?
Answer: Parquet and ORC store data column by column within row groups. A query that needs three columns out of eighty reads only those three (column pruning). Similar values sit together, so they compress well. The files also hold min/max statistics, which let engines skip row groups that cannot match a filter (predicate pushdown). Avro is row-oriented with an embedded schema and well-defined schema-evolution rules, which makes it good for writing records one at a time and for messaging. It is commonly used with Kafka and a schema registry. A common pattern: land data as Avro or JSON, store analytical tables as Parquet under a table format.
6. What does idempotency mean for a data pipeline, and how do you achieve it?
Answer: A pipeline step is idempotent if running it once or five times for the same input produces the same final state. You need this because orchestrators retry, people rerun backfills and streaming systems redeliver. Techniques:
- Overwrite a well-defined slice, such as a date partition, instead of appending blindly.
- Use
MERGEor upsert on a stable business or event key instead ofINSERT. - Derive the processing window from the scheduled interval, not from the wall clock.
- Write to a staging location and commit atomically, which table formats make cheap.
- Deduplicate on a unique event ID at the sink.
7. How should you partition a large table, and what is the small-files problem?
Answer: Partition on a column that queries commonly filter on and that has low to moderate cardinality, usually a date. Engines can then prune whole partitions. Partitioning on a high-cardinality column such as customer ID creates thousands of tiny partitions. The small-files problem is too many small files: each one adds metadata overhead, listing time, task scheduling cost and per-file open latency, so queries slow down and storage request costs rise. Common causes are streaming writes, very fine partitioning and too many Spark output tasks. Fixes include compaction jobs (OPTIMIZE in Delta, rewrite_data_files in Iceberg, compaction and clustering in Hudi), coalescing before writes, sensible target file sizes, and clustering techniques instead of deeper partition hierarchies.
8. What is the medallion (bronze, silver, gold) architecture?
Answer: It is a layering convention for lakehouses. Bronze holds raw data as received, append-only with ingestion metadata, so you can always replay. Silver holds cleaned, deduplicated, conformed data with correct types and enforced keys. Gold holds business-level aggregates and dimensional models for BI, ML features and applications. The value is not in the names. It comes from clear contracts between layers: who owns each layer, what quality checks gate promotion, and which layer each consumer is allowed to read.
Data modelling
9. Explain star and snowflake schemas. Why use surrogate keys?
Answer: A star schema has a central fact table of measurable events (order lines, payments) with foreign keys to denormalised dimension tables (customer, product, date, store). A snowflake schema normalises those dimensions into sub-tables, for example product, then category, then department. For analytics, star schemas are usually preferred: queries are simpler, BI tools understand them, and columnar engines make the redundancy cheap. Surrogate keys are warehouse-generated keys, such as integers or hashes, used instead of source system keys. They insulate the model from source key changes, let you merge customers from several systems, and are required for SCD Type 2, where one natural key has several historical rows.
10. What is the grain of a fact table, and what types of fact tables exist?
Answer: The grain is the exact meaning of one row, for example "one row per order line per shipment". Declare it before choosing any column, because mixed grains cause double counting. The three classic types:
- Transaction facts: one row per event. These are the most granular and the most flexible.
- Periodic snapshot facts: one row per entity per period, such as the daily account balance or month-end inventory.
- Accumulating snapshot facts: one row per process instance, updated as it moves through milestones, such as a loan application with submitted, approved and disbursed dates.
Facts can also be additive (sales amount), semi-additive (balances, which you cannot sum across time) or non-additive (ratios, which you should recompute from components rather than average).
11. What are slowly changing dimensions? Explain Types 1, 2 and 3.
Answer: SCDs describe how a dimension handles changing attributes. Type 1 overwrites the old value, so no history is kept. It is fine for correcting typos. Type 2 inserts a new row for each change, with a new surrogate key, valid_from and valid_to timestamps and an is_current flag, so facts join to the version that was true at the time. Type 3 keeps a limited history in extra columns, such as previous_region. A Type 2 load typically compares incoming rows to current rows using a hash of the tracked attributes, closes changed rows by setting valid_to, and inserts new versions:
-- 1. close changed current rows
UPDATE dim_customer d
SET valid_to = s.change_ts, is_current = FALSE
FROM stg_customer s
WHERE d.customer_id = s.customer_id
AND d.is_current
AND d.attr_hash <> s.attr_hash;
-- 2. insert new versions for changed and new keys
INSERT INTO dim_customer (customer_id, region,
attr_hash, valid_from, valid_to, is_current)
SELECT s.customer_id, s.region, s.attr_hash,
s.change_ts, NULL, TRUE
FROM stg_customer s
LEFT JOIN dim_customer d
ON d.customer_id = s.customer_id AND d.is_current
WHERE d.customer_id IS NULL;
Interview tip: Mention that this two-step version must run in one transaction, or as a single MERGE where the dialect supports it, otherwise a failure between the steps leaves keys with no current row. Also mention that dbt snapshots implement Type 2 for you.
12. What is Data Vault modelling, conceptually?
Answer: Data Vault is an integration-layer modelling approach designed for many changing sources and full auditability. It has three building blocks. Hubs hold unique business keys (customer number, policy number). Links hold relationships between hubs (customer holds policy). Satellites hold descriptive attributes and their history, each row stamped with a load date and record source. Loads are insert-only and can run in parallel, and adding a new source usually means adding satellites rather than redesigning tables. The trade-off is that the raw vault is hard to query directly, so teams build star-schema marts on top for consumers. It suits regulated banks and insurers that must reconstruct what was known at any point in time.
13. When would you denormalise into a wide "one big table" instead of a star schema?
Answer: A wide, pre-joined table works well when a narrow set of consumers, such as a dashboard, an ML feature set or an embedded analytics API, needs fast, simple queries and the grain is stable. Columnar engines handle wide tables efficiently, and you avoid repeated joins. The costs: every dimension change means rebuilding or backfilling the wide table, attribute definitions get copied into many tables and drift, and SCD history is awkward. A common compromise is to keep conformed dimensions and facts as the governed core and generate wide tables from them as disposable serving outputs.
SQL deep-dive
14. What is the difference between ROW_NUMBER, RANK and DENSE_RANK?
Answer: All three are window functions that number rows within a partition according to an ORDER BY. ROW_NUMBER gives unique consecutive numbers even for ties (1, 2, 3, 4), and the order among tied rows is arbitrary unless you add a tiebreaker. RANK gives ties the same number and leaves gaps (1, 2, 2, 4). DENSE_RANK gives ties the same number without gaps (1, 2, 2, 3). Use ROW_NUMBER to pick exactly one row per group, as in deduplication, and DENSE_RANK for "top N distinct values" questions.
-- top 3 products by revenue in each category
SELECT *
FROM (
SELECT category, product_id, revenue,
DENSE_RANK() OVER (PARTITION BY category
ORDER BY revenue DESC) AS rnk
FROM product_revenue
) t
WHERE rnk <= 3;
15. How do you deduplicate a table and keep only the latest record per key?
Answer: Number the rows per business key, newest first, and keep row 1. Use a deterministic tiebreaker so reruns pick the same row.
SELECT *
FROM (
SELECT o.*,
ROW_NUMBER() OVER (
PARTITION BY order_id
ORDER BY updated_at DESC, ingest_seq DESC
) AS rn
FROM raw_orders o
) t
WHERE rn = 1;
Engines that support QUALIFY let you filter on the window function directly. To delete duplicates in place, dialects differ: some allow deleting from a CTE, others need a MERGE or a rebuild into a new table. At lakehouse scale, rebuilding the affected partitions is often cheaper and safer than row-level deletes.
Interview tip: Ask what "duplicate" means first: an exact copy of the row, the same event ID delivered twice, or two versions of the same entity. Each one needs a different key and ordering column.
16. How would you compute a running total and month-over-month change in SQL?
Answer: Use window aggregates with an explicit frame for the running total, and LAG for the previous period.
SELECT month, revenue,
SUM(revenue) OVER (ORDER BY month
ROWS BETWEEN UNBOUNDED PRECEDING
AND CURRENT ROW) AS running_revenue,
revenue - LAG(revenue) OVER (ORDER BY month)
AS mom_change
FROM monthly_revenue;
Specify the frame explicitly. When an ORDER BY is present, many engines default to a RANGE frame, which treats ties on the ordering column as one peer group, and that gives surprising results when months repeat. Also check for missing months: LAG returns the previous row, not the previous calendar month, so join to a date spine first if gaps are possible.
17. How do you implement an incremental load from a source table?
Answer: Track a high-water mark, extract only rows changed since it, and merge them idempotently. The typical watermark is an updated_at timestamp or a monotonically increasing ID, stored in a control table or derived from the target.
MERGE INTO dw.orders t
USING (
SELECT * FROM staging.orders_delta
WHERE updated_at > :last_watermark
) s
ON t.order_id = s.order_id
WHEN MATCHED AND s.updated_at > t.updated_at
THEN UPDATE SET status = s.status,
amount = s.amount,
updated_at = s.updated_at
WHEN NOT MATCHED
THEN INSERT (order_id, status, amount, updated_at)
VALUES (s.order_id, s.status, s.amount, s.updated_at);
Pitfalls to mention: timestamp-based extraction misses hard deletes, which you need CDC or soft-delete flags to catch. Long-running source transactions can commit rows whose timestamps fall before the watermark you already passed, so re-read with a small overlap window and rely on the merge to stay idempotent. Deduplicate the delta before merging, because many engines fail a MERGE when several source rows match one target row. Advance the watermark only after the merge commits.
18. How would you find users who were active on three or more consecutive days?
Answer: This is the "gaps and islands" pattern. For each user's distinct active dates, subtract a row number (in days) from the date. Consecutive dates produce the same anchor value, so grouping by it gives each streak.
WITH d AS (
SELECT DISTINCT user_id, activity_date FROM events
), g AS (
SELECT user_id, activity_date,
activity_date - CAST(ROW_NUMBER() OVER (
PARTITION BY user_id ORDER BY activity_date
) AS INT) AS grp
FROM d
)
SELECT user_id, MIN(activity_date) AS start_date,
COUNT(*) AS streak_days
FROM g
GROUP BY user_id, grp
HAVING COUNT(*) >= 3;
Date arithmetic syntax varies by dialect (DATEADD, DATE_SUB, interval types); the logic does not. The DISTINCT matters, because two events on one day would otherwise break the streak.
19. What NULL-related pitfalls commonly cause wrong results in SQL?
Answer: The classic one is NOT IN with a subquery that returns a NULL. Comparing anything to NULL yields UNKNOWN, so WHERE id NOT IN (SELECT customer_id FROM orders) returns no rows at all if any customer_id is NULL. Use NOT EXISTS or a LEFT JOIN ... WHERE o.customer_id IS NULL anti-join instead. Others:
COUNT(col)skips NULLs whileCOUNT(*)does not.AVGignores NULLs, which may or may not be what the business means.- Joins on nullable keys silently drop rows.
col <> 'X'excludes NULL rows.- Concatenating with NULL returns NULL in many dialects.
Use COALESCE deliberately and test with NULLs in your fixtures.
Apache Spark
20. Explain Spark's architecture and what happens when you run an action.
Answer: A Spark application has a driver, which runs your program, builds the plan and schedules work, and executors, which are processes on worker nodes that run tasks and hold cached data. A cluster manager (YARN, Kubernetes or standalone mode) allocates the resources. Transformations such as filter, select and join are lazy and only build a logical plan. An action such as count, write or collect triggers execution. The Catalyst optimiser rewrites the logical plan (predicate pushdown, column pruning, constant folding) and picks a physical plan. The plan is split into stages at shuffle boundaries, and each stage runs as tasks, one per partition, in parallel across executor cores.
action -> logical plan -> Catalyst optimise
-> physical plan -> DAG of stages
-> tasks (one per partition) on executors
Interview tip: Cite the Spark UI as evidence ("the stage with the longest task and largest shuffle read"); it shows you have debugged real jobs.
21. What is the difference between narrow and wide transformations? Why are shuffles expensive?
Answer: In a narrow transformation (map, filter, select), each output partition depends on one input partition, so Spark pipelines these steps within a stage with no data movement. In a wide transformation (groupBy, join on non-co-partitioned data, distinct, repartition), output partitions depend on many input partitions, which requires a shuffle. A shuffle is expensive because every map task serialises and writes partitioned output to local disk, and every reduce task then fetches its blocks across the network from many executors. That costs serialisation, disk I/O, network transfer and memory pressure. It is also where skew and spills show up. Reducing shuffles means filtering early, avoiding unnecessary distinct and repartitioning, broadcasting small tables and pre-aggregating.
22. How do you choose the number of partitions? Explain repartition versus coalesce.
Answer: Aim for partitions large enough to amortise task overhead and small enough to fit in executor memory without spilling. A common rule of thumb is roughly 100β200 MB of data per partition, with at least a few times as many partitions as total executor cores, but measure rather than trust a rule. spark.sql.shuffle.partitions (default 200) controls partitions after a shuffle, and AQE can coalesce them at runtime. repartition(n) or repartition(col) performs a full shuffle and can increase or decrease the partition count, giving evenly sized or key-grouped partitions. coalesce(n) only merges existing partitions without a full shuffle, so it is cheap but can only decrease the count and may produce uneven partitions. Use coalesce to reduce output file counts and repartition to fix imbalance or to cluster data by a key before writing.
23. What is data skew in Spark, how do you detect it and how do you fix it?
Answer: Skew is when a few keys hold far more rows than the rest. After a shuffle, the tasks that own those keys run much longer than the others, so the stage waits on a handful of stragglers, and those tasks may spill or run out of memory. In the Spark UI you see it as a stage where the maximum task duration and shuffle read size are far above the median. Fixes, roughly in order of preference:
- Let AQE's skew-join handling split oversized partitions, which is on by default with AQE.
- Broadcast the smaller side if it fits in memory, which removes the shuffle entirely.
- Handle hot keys separately, for example process NULL or "unknown" keys in their own path.
- Salt the key: append a random suffix (0..N) on the large side and replicate the small side N times, then join on the salted key.
- Pre-aggregate before the join so fewer rows move.
Real-world example: In retail clickstream data, a single "guest" user ID or a NULL session ID often owns a large share of events. Handling that one key separately often beats salting.
24. Which join strategies does Spark use, and when does a broadcast join help?
Answer: The main strategies are broadcast hash join, sort-merge join, and shuffle hash join. In a broadcast hash join, the small table is sent to every executor and joined locally, with no shuffle of the large side. A sort-merge join shuffles both sides by the join key, sorts them and merges; it is the default for two large tables and is robust because it can spill to disk. A shuffle hash join skips the sort by hashing one side per partition, when those partitions fit in memory. Spark broadcasts automatically when it estimates a side is below spark.sql.autoBroadcastJoinThreshold (10 MB by default), and you can force it with a broadcast() hint. Broadcasting a side that is not really small can exhaust driver or executor memory, so check sizes after filters.
25. When should you cache or persist a DataFrame, and what can go wrong?
Answer: Cache when the same expensive intermediate result is reused by several actions in the same application, for example a cleaned dataset that feeds three aggregations, or the data of an iterative algorithm. cache() uses the default storage level (memory and disk for DataFrames), and persist() lets you choose one. Caching is lazy, so it materialises on the first action. Common mistakes:
- Caching data used only once, which adds cost with no benefit.
- Caching huge datasets, which evicts other blocks and causes recomputation or spills.
- Forgetting
unpersist()in long-running jobs. - Caching before a filter instead of after it.
For reuse across jobs, write an intermediate table instead.
26. What is Adaptive Query Execution (AQE)?
Answer: AQE lets Spark re-optimise a query during execution, using statistics collected at shuffle boundaries instead of relying only on estimates made before the query runs. It has been enabled by default since Spark 3.2. Its main features:
- Coalescing many small post-shuffle partitions into fewer, right-sized ones.
- Switching a sort-merge join to a broadcast join when one side turns out to be small after filtering.
- Splitting skewed partitions in sort-merge joins so no single task is overloaded.
AQE reduces hand-tuning of shuffle partitions, but it cannot fix a bad data model or help before the first shuffle. Configuration names and defaults change between versions, so check the documentation for your Spark version.
27. Why prefer DataFrames over RDDs, and why are Python UDFs a performance concern?
Answer: DataFrames (and typed Datasets in Scala and Java) carry a schema, so Catalyst can optimise the query and the Tungsten execution engine can use compact binary memory layouts and generated code. RDD code is opaque functions that Spark cannot optimise. Row-at-a-time Python UDFs are slow because each row must be serialised from the JVM to a Python worker process and back, and the UDF is a black box to the optimiser, which blocks pushdown. Prefer built-in functions. When you need Python logic, use vectorised pandas UDFs, which move data in Arrow batches. In Spark 4, ANSI SQL mode is on by default, so invalid casts and overflows raise errors instead of silently returning NULL.
Lakehouse and table formats
28. What problem do open table formats solve?
Answer: A folder of Parquet files is not a table. Without extra metadata you get no atomic commits, so readers can see half-written data. Two writers can corrupt each other. Updates and deletes require rewriting whole partitions. Schema changes are undocumented. Engines also have to list directories to discover files, which is slow on object storage. Open table formats add a metadata layer that tracks exactly which files make up each version of the table. That provides ACID transactions through optimistic concurrency, schema enforcement and evolution, row-level updates, deletes and merges, snapshot isolation, time travel to earlier versions, and file-level statistics for pruning. Because the format is open, several engines can read the same table.
29. How does Delta Lake work?
Answer: A Delta table is Parquet data files plus a transaction log in a _delta_log directory. Each commit is an ordered JSON file recording actions, such as files added and removed, metadata and schema changes. Periodic Parquet checkpoint files summarise the log so readers do not replay every commit. Writers use optimistic concurrency: they read a version, write new files, then try to commit the next version number, and conflicting commits are detected and retried or rejected. Features interviewers commonly ask about:
MERGEfor upserts.- Time travel by version or timestamp.
- Schema enforcement and opt-in evolution.
- Change data feed, which exposes row-level changes for downstream incremental processing.
OPTIMIZEfor compaction, with Z-ordering or liquid clustering for data layout.- Deletion vectors, which mark rows deleted without rewriting whole files.
VACUUM, which physically removes files that are no longer referenced and older than the retention period.
Feature availability depends on the Delta protocol version and on the engine, so check compatibility before you enable newer table features.
30. How does Apache Iceberg work, and what is hidden partitioning?
Answer: An Iceberg table is a tree of metadata. A catalog (for example a REST catalog, Hive Metastore or AWS Glue) stores a pointer to the current metadata file. The metadata file holds the schema, partition specs and the list of snapshots. Each snapshot points to a manifest list, which points to manifest files, which list the data files with their partition values and column statistics. A commit writes new metadata and atomically swaps the catalog pointer. Hidden partitioning means you declare partitions as transforms of columns, such as day(event_ts) or bucket(16, customer_id). Users filter on event_ts normally and Iceberg prunes for them. Partition evolution lets you change the partition spec over time without rewriting old data. Schema evolution tracks columns by ID, so renames and reorders are safe.
31. How does Apache Hudi differ, and what are Copy-on-Write and Merge-on-Read?
Answer: Apache Hudi was designed for upsert-heavy and incremental workloads, such as applying CDC streams to a lake. It organises actions on a timeline of instants (commits, compactions, cleans). Records are identified by a record key and, often, an ordering field that decides which version wins, and indexes map keys to file groups for fast upserts. It offers two table types:
- Copy-on-Write (CoW): an update rewrites the affected base Parquet files. Reads are simple and fast; writes cost more.
- Merge-on-Read (MoR): updates go to row-based log files next to the base files and are merged at read time or by later compaction. Writes are faster and fresher; reads cost more until compaction runs.
Hudi also supports incremental queries ("give me records changed since instant T"), which make chained incremental pipelines straightforward.
32. How would you choose between Delta Lake, Iceberg and Hudi?
Answer: All three provide ACID tables on object storage, so decide on ecosystem and workload rather than feature checklists:
| Consideration | What to ask |
|---|---|
| Engines | Which compute engines and warehouses must read and write the table, and how mature is each one's support? |
| Catalogue | Which catalogue will govern the tables, and does it support the format natively? |
| Workload | Mostly append and scan, or heavy upserts and CDC with tight freshness needs? |
| Platform | Is the organisation standardised on a platform that defaults to one format? |
| Operations | Who will run compaction, snapshot expiry and cleanup, and how? |
Interoperability options also exist. Delta's UniForm can generate Iceberg-compatible metadata, and Apache XTable translates metadata between formats. Capabilities move quickly, so check current documentation. In an interview, a reasoned choice with an exit path beats brand loyalty.
33. What routine maintenance do lakehouse tables need?
Answer: Four jobs, scheduled and monitored like any pipeline:
- Compaction: merge small files and, for Hudi MoR, merge log files into base files.
- Layout optimisation: clustering or sorting by common filter columns.
- Snapshot expiry and file cleanup:
VACUUMin Delta,expire_snapshotsand orphan-file removal in Iceberg, cleaning in Hudi. Without these, storage grows forever. - Statistics and metadata upkeep: for example, rewriting manifests in Iceberg.
There is also a privacy angle. Time travel means "deleted" rows still exist in older snapshots until those snapshots expire and the files are physically removed. A retention period that suits debugging may conflict with an erasure obligation, so set it deliberately and document it.
Kafka and streaming
34. Explain Kafka's core concepts: topics, partitions, offsets and consumer groups.
Answer: Kafka is a distributed, append-only log. A topic is split into partitions, and each partition is an ordered, immutable sequence of records identified by offsets. Ordering holds only within a partition, so records that must stay in order, such as all events for one account, need the same key, which the default partitioner hashes to the same partition. Partitions are replicated across brokers; acks=all with min.insync.replicas controls durability. A consumer group shares a topic's partitions among its members, with each partition read by exactly one consumer in the group, so the partition count caps a group's parallelism. Consumers commit offsets to record progress, and consumer lag (latest offset minus committed offset) is the key health metric. Recent Kafka versions run without ZooKeeper, using the built-in KRaft consensus mode, and Kafka 4.0 removed ZooKeeper support entirely.
35. What does exactly-once mean in streaming, and how is it achieved?
Answer: There are three delivery semantics. At-most-once may lose data. At-least-once may duplicate it. Exactly-once means each record affects the final result once, even through failures and retries. Within Kafka, you get exactly-once from:
- The idempotent producer, where the broker deduplicates retried sends using producer IDs and sequence numbers. It is enabled by default in modern clients.
- Transactions, which atomically write output records and the consumer's offsets together.
- Consumers set to
isolation.level=read_committed, so they ignore aborted transactions.
The key point is that exactly-once is an end-to-end property. Once data leaves Kafka for a database, a lakehouse table or an API, you need either a transactional sink that commits output and progress together (as Spark Structured Streaming and Flink do with checkpoints and supported sinks) or an idempotent sink that upserts on a unique event ID. In practice, "effectively once" usually means at-least-once delivery plus idempotent writes.
Interview tip: If you claim exactly-once, expect the follow-up "what happens if the job crashes after writing to the sink but before committing the offset?" Answer it with the mechanism.
36. What is the difference between event time and processing time? Explain window types.
Answer: Event time is when something happened, stamped at the source. Processing time is when your system sees it. They differ because of network delays, mobile devices going offline, retries and backlogs. Business metrics should almost always use event time; otherwise a 10-minute outage shifts revenue into the wrong hour. Window types:
- Tumbling: fixed-size and non-overlapping, for example every 5 minutes.
- Sliding (hopping): fixed-size and overlapping, for example a 10-minute window every 1 minute, which suits moving averages.
- Session: dynamic windows that close after a gap of inactivity per key, which suits user sessions.
Some engines also offer global windows with custom triggers.
37. What are watermarks, and how do you handle late-arriving data?
Answer: A watermark is the engine's estimate of how far event time has progressed: "I believe no more events older than T will arrive." It is usually defined as the maximum event time seen minus an allowed delay. Windows that end before the watermark can be finalised and their state dropped, which keeps memory bounded. Data that arrives behind the watermark is "late". You have four options:
- Drop it, and count the drops so the loss is visible.
- Accept it within an allowed-lateness period and emit updated results.
- Route it to a side output for separate handling.
- Correct it later in batch, for example by reprocessing the affected day's partition nightly.
The trade-off is completeness against latency and state size: a longer delay catches more stragglers but holds more state and finalises results later. In Spark Structured Streaming, withWatermark("event_ts", "30 minutes") before a windowed aggregation sets this up.
38. How do stateful stream processing and checkpointing work?
Answer: Stateful operations, such as windowed aggregations, deduplication, stream-stream joins and sessionisation, keep per-key state between records. The engine periodically checkpoints that state together with source positions (Kafka offsets) to durable storage, so after a failure it restores the state and resumes from the matching offsets. Spark Structured Streaming keeps offsets and state in a checkpoint location (incompatible query changes can invalidate it); Flink takes distributed snapshots, and its savepoints support upgrades and rescaling. Interview points: state must be bounded (watermarks, TTLs), large state needs an appropriate state backend (RocksDB-based backends are common), and checkpoint storage and intervals affect both recovery time and cost.
39. What is change data capture (CDC), and why is log-based CDC preferred?
Answer: CDC captures inserts, updates and deletes from a source database so downstream systems can stay in sync. Query-based CDC polls with WHERE updated_at > watermark. It is simple, but it misses hard deletes and intermediate updates, loads the source, and depends on reliable timestamps. Log-based CDC reads the database's transaction log (the PostgreSQL WAL, the MySQL binlog, Oracle redo logs), usually through tools such as Debezium publishing to Kafka. It captures every change in commit order, including deletes, with low source impact. Points to discuss: the initial snapshot plus the ongoing stream, schema changes in the source, ordering per primary key (key the Kafka topic by primary key), replication slot or log retention so the source does not fill its disk if the consumer stalls, and the outbox pattern, which publishes domain events reliably from application transactions.
Orchestration and transformation
40. What does an orchestrator do? Compare Airflow, Dagster and Prefect at a high level.
Answer: An orchestrator schedules and runs pipeline steps in dependency order. It handles retries, timeouts, backfills and alerting, and records run history. Apache Airflow defines workflows as Python DAGs of tasks and is the most widely deployed. Airflow 3 (released in 2025) brought a redesigned UI, DAG versioning and a task execution interface that decouples task code from the scheduler, and it renamed "datasets" to assets for data-aware scheduling. Dagster is asset-centric: you declare the tables and models (software-defined assets) and their dependencies, and Dagster gives you lineage, partitions and asset checks natively. Prefect emphasises Pythonic flows with dynamic, event-driven execution.
Interview tip: Contrast task-based orchestration ("run these steps") with asset-based orchestration ("keep these tables fresh"). The asset view maps directly onto lineage, freshness SLAs and selective backfills.
41. How do you design orchestrated tasks for retries and backfills?
Answer: Every task should be idempotent and parameterised by its data interval, the logical period it processes, not by "now". In Airflow, use the run's logical date and data interval in templates, so a rerun for 3 March processes 3 March again rather than today. Other practices:
- Keep tasks atomic, writing to staging and then swapping or merging.
- Set retries with backoff for transient failures, and fail fast on data errors.
- Put timeouts and SLAs on critical paths.
- Limit backfill concurrency so a backfill does not starve production runs or hammer a source system.
- Keep secrets in a secrets backend, not in DAG code.
- Avoid top-level code in DAG files that calls databases or APIs, because the scheduler parses those files repeatedly.
42. What role does dbt play, and how do incremental models work?
Answer: dbt handles the "T" in ELT. You write SELECT statements as models, and dbt compiles them, works out dependencies from ref(), builds them in order inside the warehouse, and generates documentation and lineage. Tests such as unique, not_null, accepted_values and relationships, plus custom tests, run as part of the build. An incremental model processes only new or changed rows on later runs. Inside an is_incremental() block you filter the source against the existing target, and with a unique_key dbt merges rather than appends. The strategy (append, merge, delete+insert, insert_overwrite) depends on the adapter. Model contracts can enforce column names and types at build time. Snapshots implement SCD Type 2. Incremental models need a full-refresh plan for logic changes, and they need protection against late data, typically a lookback window.
Data quality, contracts and lineage
43. How do you build data quality checks into a pipeline?
Answer: Test the dimensions that matter to consumers, at the boundaries where problems enter. The dimensions:
- Completeness: row counts against the source, null rates.
- Validity: types, ranges, allowed values, formats.
- Uniqueness: primary keys.
- Consistency: referential integrity, totals that reconcile.
- Timeliness: freshness against an SLA.
- Distribution: volume or value anomalies against history.
Implement checks with dbt tests, Great Expectations, Soda, Dagster asset checks or the expectations features in your platform. Decide the action per check: block promotion for critical failures such as duplicate keys in a finance table, quarantine bad rows with a reason code, or warn for soft anomalies.
Real-world example: An insurer's claims pipeline reconciles daily claim counts and paid amounts between the source system and the silver table before the gold layer builds. A mismatch blocks the dashboards and pages the owners.
44. What is a data contract, and how do you enforce one?
Answer: A data contract is an explicit, versioned agreement between a data producer and its consumers. It covers the schema (names, types, nullability), semantics (what each field means, units, keys), quality expectations, freshness SLAs, ownership, and how changes are announced. It moves breaking changes from "discovered in production" to "caught in review". Enforcement happens at several points:
- A schema registry with compatibility rules (backward, forward, full) for Kafka topics.
- CI checks that fail a producer's pull request when a contracted field is removed or retyped.
- Build-time contracts in dbt models.
- Runtime validation at ingestion, with quarantine for violations.
The hard part is organisational: producers must own the data they emit. Start with the datasets that cause the most incidents.
45. Why do lineage and data catalogues matter, and how is lineage captured?
Answer: Lineage answers "where did this number come from?" and "what breaks if I change this column?". You need it for impact analysis, debugging, audits and privacy requests. Table-level lineage shows dataset dependencies. Column-level lineage traces individual fields through transformations, which matters for PII tracking and metric definitions. Lineage is captured by parsing SQL, from dbt's dependency graph, and from runtime events emitted by engines and orchestrators. OpenLineage is an open standard for those events, with integrations for Airflow, Spark, dbt and others. A data catalogue brings together lineage, schemas, owners, descriptions, classifications, usage and access policies so people can find and trust data. Examples include open-source options such as DataHub and OpenMetadata, and platform catalogues such as Unity Catalog, AWS Glue Data Catalog, Microsoft Purview and Google Cloud's Dataplex.
Cloud warehouses, performance and cost
46. How do modern cloud data warehouses work, in general terms?
Answer: Most cloud warehouses separate storage from compute. Data lives in compressed columnar form on durable cloud storage, while independent compute clusters or serverless slots scale up and down and can be isolated per workload, so ETL does not slow down BI. Queries run as massively parallel processing (MPP) across many nodes, and engines use metadata such as micro-partition or file statistics to prune data. They also provide result caching, materialised views, clustering or sort keys, and workload management. Pricing models differ:
- Pay per data scanned.
- Pay for compute time or credits while clusters run.
- Reserved or committed capacity.
The model shapes optimisation priorities. With scan-based pricing, reducing bytes read matters most. With compute-time pricing, auto-suspend, right-sizing and workload isolation matter most.
47. What are the main levers for making analytical queries and pipelines faster and cheaper?
Answer: Read less, move less, and compute only what changed.
- Read less: partition and cluster on common filters, select only needed columns (no
SELECT *in production models), filter early, and keep healthy file sizes. - Move less: avoid unnecessary shuffles and cross-region transfer, broadcast small dimensions, and pre-aggregate.
- Compute only changes: use incremental models instead of full rebuilds, and materialised views for repeated aggregations.
- Right-size compute: auto-suspend idle warehouses, separate workloads, use spot or preemptible capacity for fault-tolerant batch jobs, and schedule heavy jobs off-peak where pricing rewards it.
- Govern spend: tag jobs to owners, set budgets and alerts, and review the most expensive queries regularly.
Our FinOps interview questions cover the organisational side of cloud cost.
Data engineering for AI
48. How is an unstructured data pipeline for AI different from a classic ELT pipeline?
Answer: The shape is the same (land, transform, publish), but the data is documents, and the transforms are parsing, chunking, enrichment and embedding instead of joins and aggregations. A data engineer brings the same discipline:
- Bronze holds raw files with source metadata.
- Silver holds parsed text, sections and tables.
- Gold holds versioned, permissioned chunks with embeddings, ready for retrieval.
The new failure modes are bad text extraction, stale documents, lost permissions and silently changed embedding models, so quality checks look at extraction quality, duplicates and chunk sanity rather than row counts. Our guide to AI for data engineers explains why data engineers end up owning GenAI answer quality, and data pipelines for RAG walks through connectors, incremental sync, deletes and permission sync stage by stage.
49. How should embeddings and vector stores be managed as part of the data platform?
Answer: Treat the chunk-and-embedding table as a governed data product. It needs a schema (chunk ID, document ID and version, text, embedding model identifier, parser version, ACL, timestamps), an owner, a freshness target and versioning. Keep the system of record in the lakehouse and treat the vector index as a rebuildable serving copy. Then changing the embedding model becomes a planned backfill into a new index plus an evaluated cutover, not an emergency. Never mix vectors from different embedding models in one index, because their similarity scores are not comparable. Many teams start with vector support in a database they already run. The pgvector RAG tutorial shows that route with PostgreSQL, and our vector database interview questions go deeper into indexing and similarity search.
50. What is a feature store, and what is point-in-time correctness?
Answer: A feature store manages ML features consistently between training and serving. An offline store (lakehouse or warehouse tables) holds historical feature values for training, and an online store (a low-latency key-value database) serves the latest values to models in production. Feast is a widely used open-source option. Point-in-time correctness means that when you build a training row for an event at time T, you join only feature values that were known at T. Using later values leaks future information, which makes offline metrics look excellent while production performance disappoints. Training-serving skew, where features are computed differently offline and online, is the other main problem a feature store aims to remove.
Interview tip: If the role leans towards ML, revise our data science interview questions too.
Privacy and governance
51. What does India's DPDP Act mean for a data engineer?
Answer: The Digital Personal Data Protection Act, 2023 applies to digital personal data processed in India, and to processing abroad connected with offering goods or services to people in India. Its Rules were notified in November 2025, with obligations phasing in over time (check the current timeline). For data engineers, the duties become engineering requirements:
- Purpose and consent: tag datasets with the purpose and consent basis they were collected under, and stop using data for incompatible purposes.
- Data minimisation: don't copy full customer tables into every sandbox.
- Erasure: when consent is withdrawn or the purpose is served, delete the data from every copy, including lakehouse snapshots, warehouse tables, backups (per policy), caches, feature stores and vector indexes.
- Security safeguards: encryption, access control, masking and logging.
- Breach response: lineage tells you whose data was affected.
Our article on the DPDP Act for AI applications maps these obligations to concrete controls. This is not legal advice; work with your legal and privacy teams.
52. How do you implement access governance on a shared data platform?
Answer: Start with classification: tag columns and tables as public, internal, confidential, PII or sensitive, automatically where possible and confirmed by owners. Then enforce policy centrally in the catalogue or warehouse rather than in each tool:
- Role-based access for broad permissions.
- Attribute- or tag-based policies, such as "only the fraud team sees unmasked PAN columns".
- Column masking and row-level security filters, for example by region or business unit.
- Tokenisation or pseudonymisation for identifiers that analytics needs to join on but not read.
Grant access to groups from the identity provider, not to individuals, review it periodically, and log every access for audit. Pipeline service accounts need least privilege too.
Real-world scenario questions
53. Your nightly pipeline's cloud cost doubled overnight. How do you investigate?
Answer: Treat it as an incident: find what changed, quantify it, contain it, then fix the cause. An overnight jump almost always has a specific trigger, not gradual growth.
What I would check:
- The billing breakdown by service, account and job tag: which job, warehouse or cluster moved?
- Deploys and configuration changes in the last 24 hours, such as new dbt models, a changed cluster size or disabled auto-suspend.
- Whether an incremental model silently fell back to full refreshes, or a filter or partition predicate was dropped so queries now scan everything.
- Input volume: did a source duplicate data, send a full re-extract or start a backfill?
- Retry storms: a failing task retrying expensive work, or a streaming job crash-looping.
- Join explosions: a dimension that lost key uniqueness and multiplies fact rows.
- Storage and request costs: a small-files explosion, or cross-region egress from a moved bucket.
Production consideration: Put budget alerts per team and per pipeline in place before this happens, plus a "cost per run" metric next to duration in pipeline monitoring.
54. After an orchestrator retry, a fact table contains duplicate records. What happened, and how do you fix it?
Answer: Almost certainly the task appended its output, failed or timed out after the write had (partly) committed, and the retry appended the same data again.
What I would check:
- Run history and logs: which attempt wrote what, and whether the first attempt's write committed.
- The scope: affected partitions or load batches, identified by ingestion metadata such as batch ID and load timestamp.
- Downstream exposure: which reports and consumers read the duplicated data, and whether they need re-running or a notice.
- Whether the source itself redelivered, which needs deduplication on event ID, not just on load batch.
The fix has two parts. Clean up by rebuilding the affected partitions or deduplicating on the business key with ROW_NUMBER. Then make the task idempotent: overwrite its partition, or MERGE on a unique key, and add a uniqueness test that blocks promotion.
Production consideration: A primary-key uniqueness check on every fact load is cheap and catches this class of bug immediately.
55. Mobile app events arrive hours or days late, and yesterday's dashboard numbers keep changing. How do you design for this?
Answer: Accept that late data is normal for mobile, because devices go offline, and design for it explicitly. Partition by event date, not arrival date.
What I would check:
- The lateness distribution: what share of events arrive within an hour, a day, a week? This sets the window.
- Whether events carry reliable event timestamps and unique IDs, since clock skew on devices is common.
- How downstream consumers use the data: do they need fast provisional numbers, final numbers, or both?
A practical design:
- Streaming aggregates with a watermark give provisional near-real-time numbers.
- A daily batch job reprocesses a rolling lookback window (for example the last 3 days of event dates) with idempotent partition overwrites or merges.
- Data older than the window is "closed". Very late events go to a correction table or a periodic restatement.
- The dashboard labels recent periods as provisional.
Production consideration: A dashboard that says "provisional until T+3" keeps trust; one that silently changes loses it.
56. An upstream team renamed a column and changed its type. Three downstream pipelines broke. How do you respond and prevent a repeat?
Answer: First restore service, then make the failure impossible to repeat silently.
What I would check:
- The blast radius, using lineage: which models, dashboards, ML features and exports depend on that column?
- Whether the change failed loudly, with errors, or silently, producing NULLs or wrong casts.
- Whether the producer can temporarily ship both the old and new columns.
Short term: add a compatibility mapping in the staging layer, which is the one place that should absorb source quirks, then rerun and backfill the affected windows. Long term:
- Agree a data contract for the source.
- Enforce it with schema-registry compatibility rules or CI checks on the producer's repository.
- Run schema-drift detection at ingestion that quarantines rather than coerces.
- Agree a deprecation process: add the new column, migrate consumers, then remove the old one.
Production consideration: Additive changes, such as new nullable columns, should flow through automatically. Removals, renames and type changes should require an explicit, versioned change.
57. You are asked to build a RAG ingestion pipeline over a company's SharePoint, Confluence and PDF policy documents. How would you design it?
Answer: Design it as a governed data pipeline whose output is a retrieval-ready, permissioned, versioned index, not a one-off embedding script.
connectors (delta sync + deletes + ACLs)
-> bronze: raw files + source metadata
-> parse (text, tables, OCR) -> quality gate
-> silver: sections + metadata
-> chunk + embed (model id recorded)
-> gold: chunks, vectors, ACLs, versions
-> vector index (rebuildable) -> eval set
What I would check:
- Source APIs and change feeds: can I sync incrementally, and are deletes and permission changes exposed?
- Permission model: how are document ACLs mapped to the identities the assistant will enforce at query time?
- Document types: scanned PDFs and tables need proper parsing (see document parsing for RAG).
- Metadata to filter on, such as effective date, region and document type, and how to handle superseded policy versions.
- Freshness target, re-embedding strategy, and an evaluation set of real questions with expected sources.
- PII in the corpus, and whether it may be indexed at all.
Production consideration: Monitor sync lag, quarantine volume, deletes propagated and retrieval quality on the evaluation set. Our RAG interview questions cover the retrieval and generation side.
58. The last few tasks of a Spark stage run for an hour while the rest finish in minutes. What do you do?
Answer: That pattern is skew or a straggler, so confirm which before tuning anything.
What I would check:
- In the Spark UI, compare the slow tasks' shuffle read size and record counts with the median. If they are much larger, the cause is data skew.
- If the sizes are similar, look at the executor: a bad node, garbage-collection pauses, spills to disk or a slow storage path.
- Find the hot keys with a quick
groupBy(key).count()ordered descending. NULLs and default values are common culprits. - Confirm AQE and its skew-join handling are enabled, and whether the join is one AQE can split.
Then apply the cheapest fix first: filter or separately process the hot key, broadcast the smaller side, pre-aggregate, and only then salt. Enable speculative execution only for genuine stragglers, because it does not help skew.
Production consideration: Record the fix and the key distribution in the job's runbook.
59. Finance says the revenue in your dashboard does not match the ERP system. How do you reconcile?
Answer: Narrow the difference systematically, layer by layer, until you find the step where the numbers diverge,.
What I would check:
- Definitions: gross or net, tax included, returns and cancellations, currency conversion date, recognition date versus order date, and time zone boundaries.
- Pick one day and one entity and compare totals at source, bronze, silver and gold to find the layer where they split.
- Row-level diffs on that slice: missing rows (late or failed loads, filters), extra rows (duplicates, test orders) and changed values (stale SCD joins, rounding).
- Join fan-out: a non-unique dimension key multiplying facts.
Production consideration: Turn the reconciliation into an automated daily check with a tolerance, and publish the metric definition in the semantic layer or catalogue so the next dispute is settled by a document, not a meeting.
60. Consumer lag on a critical Kafka topic keeps growing. What do you investigate?
Answer: Growing lag means consumers process more slowly than producers write.
What I would check:
- Is lag growing on all partitions, which points to throughput, or a few, which points to a hot key or a stuck consumer?
- Producer rate changes, such as a campaign, a backfill or a duplicate publisher.
- Consumer health: frequent rebalances, which point to processing exceeding
max.poll.interval.ms, errors and retries on poison messages, or garbage-collection pauses. - Downstream sink latency: a slow database or API call per record is the usual bottleneck. Batch writes where possible.
- Parallelism: are there more consumers than partitions (idle consumers), or too few partitions to scale out?
Fixes include batching sink writes, sending poison messages to a dead-letter topic, scaling consumers up to the partition count, and increasing partitions with care, because that changes key-to-partition mapping and ordering for keyed data.
Production consideration: Alert on lag in time terms ("how many minutes behind"), not just message counts.
61. You need to backfill two years of history after a logic change. How do you do it safely?
Answer: Backfill into a parallel target, validate, then swap, and never let the backfill compete blindly with production.
What I would check:
- That the job is idempotent and parameterised by data interval, so it can run partition by partition and resume after failures.
- Source availability and load: can the source serve two years of extraction, or should I read from the bronze layer instead?
- Cost and capacity: run in chunks with limited concurrency, preferably off-peak, and estimate cost from a one-month trial run.
- A validation plan: compare old and new outputs per partition, and explain every difference as an intended change or a bug.
Production consideration: Write the backfill to a new table version or shadow table and switch consumers atomically (a view swap, or a table format commit). Tell finance and analytics users about restated history before they notice it.
62. A customer withdraws consent and requests erasure. Their data is in the lakehouse, warehouse, feature store and a RAG vector index. What do you do?
Answer: Use lineage and identifiers to find every copy, delete or anonymise each one with the method that store needs, and record evidence of completion.
What I would check:
- An identity map: which keys (customer ID, email, phone, tokens) identify this person in each system?
- Lineage from the source tables to every derived table, extract, feature and index that carries their data.
- Store-specific deletion:
DELETEorMERGEin table formats followed by snapshot expiry and file cleanup, so old versions are physically gone. Warehouse deletes and time-travel retention. Removal from the online and offline feature stores. Deletion of the chunks and vectors derived from their documents. - Legal retention exceptions, for example records that sector rules require you to keep, decided with the legal team.
- Backups, handled according to documented policy.
Production consideration: Build erasure as a repeatable, audited pipeline, not a manual ticket.
63. A GCC team asks you to plan a migration from an on-premises Hadoop and ETL-tool estate to a cloud lakehouse. How do you approach it?
Answer: Migrate by data product in waves, with parallel runs and reconciliation.
What I would check:
- An inventory of jobs, tables, schedules, consumers and actual usage. Many old jobs have no live consumers and can be retired rather than migrated.
- Dependencies and the critical path, starting with a contained, valuable domain.
- Target standards defined first: table format, catalogue, naming, layers, orchestration, CI/CD, access model and cost tagging.
- Conversion approach: rewrite legacy transformations as SQL or Spark code under version control, with tests.
- Parallel runs with automated reconciliation before each cutover, and a decommission date for the old jobs.
- Network, security and data-residency requirements for regulated data.
Production consideration: Track migration progress by consumers moved and legacy jobs switched off, not by tables copied.
64. Users say the internal AI assistant gives outdated answers after HR updated several policies. Where in the data pipeline do you look?
Answer: Stale answers usually mean a freshness or versioning failure in ingestion, not a model problem. Trace one specific question back to the chunks that were retrieved.
What I would check:
- Did the connector pick up the updated documents? Check sync lag and last successful run per source.
- Were old versions retired? If old and new chunks coexist without version metadata, retrieval may prefer the old text.
- Did parsing or a quality gate quarantine the new documents?
- Was the index rebuilt or updated, and does the application read the current index build?
- Do retrieval filters use effective dates, so superseded policies are excluded?
Production consideration: Publish a freshness SLO for the corpus, alert on sync lag, and add "recently changed document" questions to the evaluation set so staleness shows up as a failing test. Getting enterprise knowledge ready for AI covers the content-ownership side of this.
65. Design a data platform for a retailer that needs near-real-time inventory, daily finance reporting and ML demand forecasting.
Answer: One governed lakehouse with streaming and batch paths that share tables, contracts and definitions, rather than three separate stacks.
POS / e-com DBs --CDC--> Kafka --> stream jobs
| |
v v
bronze (raw, replayable) --> silver (dedup, keys)
|
+--------------------------+---------+
v v v
inventory serving gold marts features
(low latency) (finance BI) (offline +
online)
What I would check:
- Latency needs per consumer: inventory in seconds to minutes, finance daily and reconciled, forecasting daily or weekly.
- Sources and CDC feasibility for point-of-sale, e-commerce and warehouse systems, with event IDs for idempotency.
- Late and out-of-order data from stores with flaky connectivity: watermarks for streaming, plus lookback reprocessing in batch.
- Modelling: conformed product, store and date dimensions; inventory snapshots; sales transaction facts.
- Quality gates, contracts with source teams, lineage and an access model for customer data.
- Cost: streaming only where the business needs seconds, and batch everywhere else.
Production consideration: Make finance the strictest consumer: reconciled, closed periods and documented definitions. Present it in stages, MVP first.
If you want to practise these designs on real pipelines rather than only on paper, Cloudsoft's HORIZON Data Engineering & AI program covers modern data engineering together with the data foundations that AI systems depend on.
Key takeaways
- Idempotency is the most important pipeline property: overwrite slices or merge on keys so retries and backfills are safe.
- Declare the grain before modelling, use surrogate keys, and know SCD Type 2 well enough to write it.
- Window functions handle most SQL interview problems: deduplication, top-N, running totals and gaps-and-islands.
- Spark performance comes down to shuffles, partition sizing, skew and join strategy. Read the Spark UI before tuning.
- Delta Lake, Iceberg and Hudi all bring ACID tables to object storage. Choose on engines, catalogue and workload, and plan their maintenance.
- Exactly-once is an end-to-end property. In practice it means transactional or idempotent sinks, plus watermarks and a late-data policy.
- Pipelines for AI are still data pipelines. Freshness, permissions, versioning and lineage decide answer quality.
Interview preparation checklist
- Write, from memory, SQL for deduplication, top-N per group, running totals, gaps-and-islands and an incremental
MERGE. - Model one business process as a star schema, declare the grain, and implement an SCD Type 2 dimension.
- Run a Spark job locally or on a small cluster. Find a shuffle in the UI, force a broadcast join and observe a skewed join.
- Create a Delta or Iceberg table, run an upsert, time-travel to a previous version, then compact and expire snapshots.
- Build a small Kafka pipeline with keyed messages and consumer groups, and explain what happens on a consumer crash.
- Write a streaming aggregation with a watermark and test it with deliberately late events.
- Build an orchestrated, idempotent daily pipeline with retries and a backfill, plus dbt models with tests.
- Prepare one story each about a data quality incident, a performance or cost fix, and a schema change you handled.
- Build a small document-to-embeddings pipeline that records parser and model versions and handles deletes.
- Revise privacy basics: classification, masking, row-level security and how erasure reaches every copy.
FAQ
What skills are required for a data engineer interview in 2026?
Strong SQL, Python, data modelling, and one distributed processing engine (usually Spark) form the core. Add a lakehouse table format, streaming basics with Kafka, orchestration, data quality practices, a major cloud platform, and increasingly the pipelines that feed AI and retrieval systems.
How should I prepare for a data engineering interview?
Practise SQL problems daily, build one end-to-end project with ingestion, modelling, orchestration, tests and a dashboard, and rehearse scenario answers about failures you have handled. Explain trade-offs out loud, because interviewers judge reasoning as much as the final answer.
Is SQL or Python more important for data engineering interviews?
Both are expected, but SQL is tested in almost every data engineering interview, often in a live round. Python is tested for pipeline code, data manipulation, APIs and Spark jobs. Be fluent in SQL window functions and comfortable writing clean, testable Python.
Do I need to know Spark for a data engineering role?
For most mid-level and senior roles, yes, or an equivalent distributed engine. Even warehouse-centric roles benefit from understanding partitions, shuffles, skew and join strategies, because the same ideas govern query performance in cloud warehouses.
Which table format should I learn first: Delta Lake, Iceberg or Hudi?
Learn the concepts that all three share first: transaction logs or metadata trees, snapshots, upserts, time travel and compaction. Then go deeper into whichever format your target employers use. The concepts transfer quickly between formats.
How is data engineering changing because of AI?
Data engineers now also build pipelines for documents, embeddings and vector indexes, maintain feature stores, and own the freshness, permissions and lineage that decide AI answer quality. AI coding assistants also speed up routine SQL and pipeline work, which raises expectations for review and testing skills.
Can a fresher get a data engineering job?
Yes, especially with solid SQL, Python and one well-documented project that shows ingestion, modelling, orchestration and data quality. Freshers are usually assessed on fundamentals and problem solving rather than years of experience with specific tools.
Is data engineering a good career for Indian engineers?
Data engineering skills are used across GCCs, product companies and services firms in Hyderabad, Bengaluru and other cities, in banking, retail, healthcare and technology. Combining core data engineering with cloud platforms and data for AI keeps your options broad.
How long does it take to prepare for a data engineering interview?
It depends on your background. A developer or analyst with good SQL can usually cover the core topics and build one project in a few focused weeks; a fresher should plan for longer and spend most of that time building and debugging real pipelines.
Ready to build production data pipelines and the data foundations behind enterprise AI? The HORIZON data engineering and AI course combines hands-on pipeline engineering with AI data work, in classroom sessions in Ameerpet or live online. If you want a broader path across AI, ML, cloud and security, look at the APEX AI, ML, Cloud and Cyber Security program. Call +91 96660 19191 to book a free demo.



