Skip to content
Tech Interview Prep home

Top 100 Data Engineer Interview Questions and Answers

The questions most likely to actually come up in your Data Engineer interview, ranked by likelihood — with detailed, senior-level answers covering what an interviewer is really listening for.

Curated: · Written: · Reviewed:

Reviewed 78Review pending 22
QA-1The source dump is 1.8 TB. Where do you transform, and what breaks if you pick the wrong side?(show answer)

The first thing I would establish about ETL versus ELT is which write actually landed, not which task box turned green.

ETL transforms before the warehouse sees rows. ELT lands raw bytes first and transforms in the warehouse or lake engine. The choice is about where CPU, schema enforcement, and PII redaction happen, not about which scheduler box is prettier.

Concretely, extract with a documented schema, land immutable raw objects when ELT, push heavy joins and scd merges into the engine that already holds the data, and keep source-side ETL only for redaction or volume the warehouse cannot ingest.

The reason for that specificity is a failure I have seen: A Spark ETL cluster hashed 1.8 TB on 48 executors for 6 hours, then copied 1.8 TB again into the lake. The warehouse still re-joined the same keys. Slot-hours were 3.4× an ELT path that landed Parquet once and transformed in 41 minutes.

Same 1.8 TB dump, two transform locations.

PathBytes written to object storageWall timeSlot-hours
ETL hash then copy3.6 TB (twice)6 hours plus copy3.4×
ELT land Parquet then SQL1.8 TB41 minutes transform1×

I would not consider it settled without evidence: time a 1.8 TB extract with transform-in-cluster versus land-then-SQL, and require the chosen path to name where PII is stripped before any analyst role can SELECT.

Transform location is a cost and a privacy boundary, not a branding choice.

Curated: · Written: · Reviewed:

QA-2ELT landed JSON. The transform still runs as a Python UDF on the cluster. What did you fail to use?(show answer)

I would start warehouse SQL after landing from the destination checksum and the failed partition, not from the DAG run.

Once the bytes are in a table engine, the default transform is set-based SQL or native Spark, not a row Python loop. Landing then UDFing every record is ETL wearing an ELT landing zone.

Concretely, register the landed files as a table, express filters and joins in SQL or DataFrame native functions, and confine Python to code the engine cannot express.

The reason for that specificity is a failure I have seen: 12 billion JSON rows were exploded in a Python UDF at 2,400 rows/s per core. On 200 cores that is 480,000 rows/s, so 12 billion rows take about 6.9 hours, not 19. The same explode in Spark SQL finished in 28 minutes on the same 200 cores.

Python explode versus Spark SQL explode.

TransformRowsCoresWall time
Python UDF explode12,000,000,000200~6.9 hours
Spark SQL explode12,000,000,00020028 minutes

I would not consider it settled without evidence: replace the UDF with native explode on a 10 million row sample and require wall time to drop by at least 10× before the production DAG is allowed to keep Python.

Landing is wasted if every row still visits CPython.

Curated: · Written: · Reviewed:

QA-3Debezium emits 40,000 changes an hour. The warehouse also gets a 02:00 dump. How do you apply both without inventing a third history?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

A CDC stream is an ordered change log. A snapshot is a complete cut. Applying the dump as inserts on top of CDC duplicates keys; applying CDC as a replace of the table deletes every key the hour's file did not mention.

Concretely, tag each run with mode, use the dump only to rebuild an as-of table or to bootstrap, apply CDC with sequenced LSNs onto a current table, and refuse a job that cannot name which operator it is.

The reason for that specificity is a failure I have seen: Tuesday's 02:00 dump was merged as inserts onto a CDC table. Row count jumped from 9.4 million to 18.8 million. COUNT(DISTINCT customer_id) stayed 9.4 million because every key already existed.

Dump merged as CDC inserts.

Operator usedRowsDistinct customer_id
dump as insert onto CDC table18,800,0009,400,000
dump as snapshot replace9,400,0009,400,000
CDC apply of 40,000 changes9,400,000 ± net9,400,000 ± net

I would not consider it settled without evidence: feed one dump and one CDC file into the loader in staging and require two different operators, with row count matching the dump only on the snapshot path.

A dump is not a change log, and a change log is not a dump.

Curated: · Written: · Reviewed:

QA-4You turn on Debezium against a 14-year orders table. What must finish before streaming CDC is allowed to merge?(show answer)

My answer to CDC snapshot bootstrap begins with where extract, load, and transform actually ran, and on which clock.

Streaming CDC from "now" without a consistent snapshot leaves a hole for every row that did not change after the connector started. Bootstrap is a snapshot at an LSN plus CDC from that LSN, not a hope that old rows will be touched.

Concretely, take a snapshot that records the binlog/LSN, land it as the base table, start streaming from that LSN, and block MERGE until the snapshot job's row count matches the source count at that LSN.

The reason for that specificity is a failure I have seen: Streaming started with 0 snapshot. 61 million unchanged orders never arrived. Finance closed on 38 percent of GMV for 11 days while the DAG stayed green because the stream task had no lag.

Stream without a snapshot.

StateOrders in warehouseSource orders at start LSNGMV captured
stream only0 of unchanged61,000,00038 percent
snapshot then CDC from LSN61,000,00061,000,000100 percent

I would not consider it settled without evidence: compare source COUNT(*) at the snapshot LSN with the landed base table before enabling the stream MERGE.

CDC from now is not a backfill of history.

Curated: · Written: · Reviewed:

QA-5The extract uses WHERE updated_at > last_run. Why do 0.7 percent of Saturday's rows still vanish?(show answer)

I would treat incremental watermark column as a claim about every partition and every late window, not about the one that was demonstrated.

A high-water mark on updated_at misses rows whose clock is behind the extractor, rows updated in a transaction that commits after the stamp, and rows whose updated_at is timezone-naive. The watermark is a cursor, not a completeness proof.

Concretely, use a strictly increasing extract cursor the source actually commits (LSN, incrementing id, or commit timestamp), overlap the previous window by the observed clock skew, and reconcile nightly with a full-count check.

The reason for that specificity is a failure I have seen: A replica lagged 9 minutes. 18,400 Saturday updates committed with timestamps already behind last_run. Those rows never extracted. Monday's reconcile found 0.7 percent of GMV missing, $410,000.

Replica lag versus a timestamp cursor.

CursorReplica lagSaturday rows missedGMV missed
updated_at > last_run9 minutes18,400$410,000 (0.7 percent)
LSN cursor plus 15-minute overlap9 minutes0$0

I would not consider it settled without evidence: inject 9 minutes of replica lag and require the extract either to overlap by more than the lag or to key off LSN rather than updated_at.

updated_at is a hint; the cursor has to be monotonic on commit.

Curated: · Written: · Reviewed:

QA-6Max event time seen is 10:07. The out-of-orderness bound is 20 minutes. What is the watermark, and what does it claim?(show answer)

The useful question for event-time watermark arithmetic is what a second consumer would read from the same sink after a retry.

A streaming watermark is the maximum observed event time minus an out-of-orderness bound: Flink's forBoundedOutOfOrderness delay, or Spark's withWatermark threshold. At 10:07 with a 20-minute bound the watermark is 09:47. Flink's allowedLateness is a separate knob that holds window state open past the window end; it is not what the watermark subtracts.

Concretely, set watermark = max_event_time_seen - out_of_orderness_bound, close windows that end at or before that watermark, keep allowedLateness as a separate retraction budget, and send still-earlier events to a documented late path rather than dropping them without a counter.

The reason for that specificity is a failure I have seen: The job used now plus 20 minutes as a watermark. Windows never closed. State grew to 180 GB in 4 hours, checkpoints took 11 minutes, and the 09:00 rollup was 0 rows while 2.4 million events sat in timers.

Watermark at max event time 10:07.

FormulaWatermark09:00 windowState
max event time minus the 20-minute bound09:47closablebounded
wall clock plus 20 minutes~10:27never closed180 GB in 4 hours

I would not consider it settled without evidence: feed events whose max timestamp is 10:07 with a 20-minute out-of-orderness bound and require the watermark to read 09:47, not 10:27.

Minus the out-of-orderness bound, not plus, and not wall clock.

Curated: · Written: · Reviewed:

QA-7The freshness SLA is 15 minutes. Allowed lateness is 2 hours. What did you just tell finance?(show answer)

I would settle allowed lateness versus SLA by replaying the failed path against the contract, not against a green Airflow square.

Allowed lateness is how long a window stays open for retractions after the watermark has already passed its end; the out-of-orderness bound is what holds the watermark back in the first place. Either way the hour is not final until that budget expires. The freshness SLA is when a consumer is allowed to read a closed number. If the lateness budget runs longer than the SLA, the published hour is still open when the board meeting starts.

Concretely, set the lateness budget to the observed tail of out-of-order delays, publish a separate closed-hour flag only when the watermark plus that budget has passed the hour, and do not advertise a 15-minute SLA on a 2-hour lateness job.

The reason for that specificity is a failure I have seen: The DAG marked 09:00 success at 09:12. Watermark did not close 09:00 until 11:07. The board pack used 09:12 numbers that later moved 4.8 percent when late events landed.

SLA clock versus watermark close.

Clock09:00 marked09:00 actually closedLater move
DAG success09:12not closed4.8 percent
watermark 2-hour lateness—11:070 after close

I would not consider it settled without evidence: measure p99 event delay for 7 days and require advertised SLA ≥ watermark close time for that hour.

A green hour at minute 12 is not a closed hour at minute 12.

Curated: · Written: · Reviewed:

QA-8An event arrives with timestamp 08:51 after the watermark has passed 09:00. Where does it go?(show answer)

The judgement in late events after watermark close is which operator wrote the files and whether that operator is idempotent.

Late-beyond-watermark behaviour is engine-specific, not a universal exclude. Flink may drop, side-output, or still update under allowedLateness; Spark may send a late-records path or ignore the event depending on output mode. Dropping without a counter is still a silent completeness hole.

Concretely, name what this engine does after the watermark, emit late events to a repair topic keyed on the same grain when it will not retract, MERGE them into the warehouse with the same idempotent key, and increment a late_after_close counter the on-call pages on.

The reason for that specificity is a failure I have seen: 84,000 checkout events arrived 26 minutes late. They were dropped. Conversion for 08:00 stayed 1.9 points low for the rest of the day with no DLQ depth because there was no DLQ.

08:51 event after 09:00 watermark.

HandlingLate events08:00 conversion errorPage
drop84,0001.9 pointsnone
repair MERGE84,000 applied0optional

I would not consider it settled without evidence: inject events 26 minutes after close and require either a repair MERGE that restores the grain or an explicit drop metric that pages.

Late is whatever that engine documents, plus a path you can measure.

Curated: · Written: · Reviewed:

QA-9One Kafka partition goes silent for 40 minutes. The other 11 keep moving. Why did every window freeze?(show answer)

Where candidates lose the interview on idle-source watermarks is calling the DAG success the load.

A watermark that waits for the minimum across all partitions stalls on an idle source. Idle-source handling excludes the idle partition from the watermark calculation so the others can close; it does not advance that partition's own watermark. Silence may still be a real producer death, so page it.

Concretely, enable idleness with a timeout shorter than the SLA so idle partitions are skipped in the min watermark, page if a partition is idle past that timeout without a producer heartbeat, and never wait forever for a silent partition that still holds a low event time.

The reason for that specificity is a failure I have seen: Partition 7 idled. The global watermark stuck at 07:14. 11 healthy partitions buffered 3.1 million events. The 08:00 rollup was empty for 40 minutes while Airflow still showed the stream task running.

One idle Kafka partition.

SettingWatermarkBuffered events08:00 rollup
wait for all partitionsstuck 07:143,100,000empty 40 minutes
idle timeout 2 minutesidle partition excluded from calcflushedon time

I would not consider it settled without evidence: stop producing on one partition and require windows on the others to close within the idle timeout, plus an alert if idle exceeds it.

Silence is either idleness or an incident; it is not a reason to stall the fleet.

Curated: · Written: · Reviewed:

QA-10You windowed on processing time because event time was messy. What did the 09:00 bucket actually mean?(show answer)

I would answer processing-time windows by separating delivery, transform location, and the row set the sink is allowed to keep.

A processing-time window groups by when the job saw the record. Replay, lag, and retries move rows between buckets. Event-time windows group by the business clock and need watermarks; they do not move on replay of the same events.

Concretely, use event time for any metric a human will compare across days, keep processing time only for operational flush, and refuse a finance hour built on processing time.

The reason for that specificity is a failure I have seen: A 14-minute consumer lag shifted 620,000 orders from the 09:00 processing-time bucket into 09:15. Replay the next day put them back in 09:00. Two board packs disagreed by $6.1 million.

Lag then replay of the same offsets.

Window clock09:00 after 14-minute lag09:00 after replay
processing timemissing 620,000 ordersrestored
event timestablestable, same $

I would not consider it settled without evidence: replay the same Kafka offsets after a 14-minute pause and require event-time totals to match while processing-time totals are allowed to move.

Processing time is when the job woke up, not when the purchase happened.

Curated: · Written: · Reviewed:

QA-11The stream says exactly-once. The sink is an S3 PUT plus a later Iceberg commit. What guarantee do you actually have?(show answer)

The engineering content of exactly-once sink transactions is the merge key and the quality gate that would refuse the write, not the scheduler.

Exactly-once end-to-end needs a single atomic commit that covers the sink files and the source offsets. Iceberg data files and Kafka offsets are not one atomic transaction. A crash after PUT and before the table commit leaves orphan files that are in no snapshot; a retry writes and commits a new copy, it does not commit both copies into the table.

Concretely, commit Iceberg (or Delta) and the streaming checkpoint in one two-phase protocol the engine supports, remove failed pre-commit orphans, or make the sink idempotent on event id so a replay merges to the same row set.

The reason for that specificity is a failure I have seen: PUT succeeded, then the job died before the Iceberg commit. The 2.2 million rows sat as orphans, not in a snapshot. Restart wrote new files and committed those. Table COUNT(*) did not double; COUNT(DISTINCT order_id) in the table was one copy, with the first files left on disk until orphan cleanup.

PUT then crash before table commit.

RestartTable COUNT(*)On-disk orphans
retry after failed commitone committed copyfirst PUT files
two-phase commit of files and offsetsone committed copycleaned

I would not consider it settled without evidence: kill the job after object PUT and before table commit, restart, and require one committed snapshot plus orphans outside it, not two committed copies of the same keys.

Failed pre-commit files are orphans, not a second snapshot.

Curated: · Written: · Reviewed:

QA-12Kafka is at-least-once. How do you keep the fact table from growing on every retry?(show answer)

Before calling at-least-once plus MERGE done I would write down the partition, consumer, or late window nobody checked.

At-least-once delivery is the usual bus. Uniqueness is the sink's MERGE on a deterministic event id. Believing the broker adjective is not a uniqueness constraint.

Concretely, carry event_id from the producer, MERGE on that id, and treat offset commits as progress, not as proof of a single row.

The reason for that specificity is a failure I have seen: A 3-hour outage retried the same batch 9 times into an append sink. Extra rows were 8× the original 1.4 million. AVG(order_value) did not change because the duplicates were exact copies; SUM(order_value) did.

Nine append retries of 1.4 million events.

SinkRowsSUMAVG
append12,600,000×9unchanged
MERGE on event_id1,400,000×1unchanged

I would not consider it settled without evidence: replay one batch 9 times and require 1 row per event_id and an unchanged AVG of an exact duplicate payload.

Retries duplicate rows; exact copies do not change the average.

Curated: · Written: · Reviewed:

QA-13The consumer commits offsets before the warehouse MERGE returns. What happens on a crash?(show answer)

The first thing I would establish about Kafka offset versus sink write is which write actually landed, not which task box turned green.

Committing offsets first marks the messages as done while the sink may still be empty. Those offsets will not be reread. Committing after a successful idempotent MERGE can reread, which is safe if the MERGE key holds.

Concretely, write the sink, confirm the commit, then checkpoint offsets, and never ack Kafka before the table snapshot exists.

The reason for that specificity is a failure I have seen: Offsets advanced, then the MERGE timed out. 960,000 payments were acked and missing. The DAG was green. Finance found a $2.8 million hole on Monday.

Offset commit before MERGE.

Order of operationsPayments after crashDAG
commit offsets, then MERGE timeout960,000 missinggreen
MERGE then checkpoint offsets960,000 present (or retried merge)retry

I would not consider it settled without evidence: kill the worker after offset commit and before MERGE, and require the test to fail that design; the passing design still has the rows after restart.

Ack after the table commit, not before.

Curated: · Written: · Reviewed:

QA-14A Spark job uses a transactional Kafka producer. The consumer isolation is read_uncommitted. What did the consumer just see?(show answer)

I would start Kafka transactional producer from the destination checksum and the failed partition, not from the DAG run.

Kafka transactions are invisible to consumers that use read_uncommitted. Aborted transactions still appear as records. Exactly-once production requires read_committed on every consumer that must not see aborted batches.

Concretely, set isolation.level=read_committed on warehouse consumers, fence zombie producers, and refuse a job that reads the topic with the default uncommitted isolation while claiming EOS.

The reason for that specificity is a failure I have seen: An aborted 7-minute transaction still delivered 410,000 cancelled orders to a read_uncommitted sink. Those orders were later aborted in Kafka but stayed in the lake. GMV was high by $9.4 million until a manual delete.

Aborted Kafka transaction.

isolation.levelAborted orders in lakeExtra GMV
read_uncommitted410,000$9.4 million
read_committed0$0

I would not consider it settled without evidence: abort a transactional batch and require 0 of those keys in a read_committed consumer; require the uncommitted consumer to be banned from the warehouse path.

Transactions need a committed reader.

Curated: · Written: · Reviewed:

QA-15You MERGE on (order_id). The source emits order_id plus a revision. Why do you still lose updates?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Idempotence needs the key that distinguishes two writes that must both survive. If revisions are real, the merge key is (order_id, revision) or a version comparison. Merging only on order_id keeps the first or last writer depending on the clause, not the latest revision.

Concretely, include the version in the ON clause or use WHEN MATCHED AND source.rev > target.rev, and test a lower revision arriving after a higher one.

The reason for that specificity is a failure I have seen: Revision 3 arrived, then a delayed revision 2 overwrote it because MERGE matched only order_id. 27,000 orders reverted. Refund state was wrong for 6 hours, $1.1 million.

Late lower revision.

ON clauseAfter rev 3 then rev 2Refund dollars wrong
order_id onlyrev 2$1,100,000
order_id and source.rev > target.revrev 3$0

I would not consider it settled without evidence: apply rev 3 then rev 2 and require the table to keep rev 3.

The merge key is the identity of a write, not only the business id.

Curated: · Written: · Reviewed:

QA-16Reloading 2 March should be safe. When is INSERT OVERWRITE the right idempotent operator?(show answer)

My answer to partition replace idempotency begins with where extract, load, and transform actually ran, and on which clock.

A partition overwrite is idempotent when the input is the complete row set for that partition. It is a delete of good data when the input is a partial hour, a filtered sample, or a stream micro-batch. Completeness of input is the contract; the operator itself is valid.

Concretely, overwrite a day only from a job whose input contains every row that should survive for that day, use MERGE for partials, and fail the overwrite if the input count is below the source count minus a documented tolerance.

The reason for that specificity is a failure I have seen: An overwrite of day=2026-03-02 ran from a 1-hour retry file. Iceberg replaced the day's files with that hour. 23 hours of facts, 19.4 million rows, vanished. A complete-day overwrite would have been a legal reload.

Day overwrite from a 1-hour file.

InputRows after overwriteValid operator?
1 incomplete hour1 hour kept, 19,400,000 goneno
complete 2 March extractfull dayyes

I would not consider it settled without evidence: attempt a day overwrite from a 1-hour file and require the job to fail until the input is the full day.

Overwrite is legal; incomplete input is the bug.

Curated: · Written: · Reviewed:

QA-17The DAG is green. Finance says 2 March is short. What did the green square not measure?(show answer)

I would treat Airflow success versus table checksum as a claim about every partition and every late window, not about the one that was demonstrated.

Orchestrator success means every task returned exit 0. It does not mean the destination partition matches the source, the contract, or yesterday's grain. The check is a checksum, a count, or a dbt test against the table, not the scheduler UI.

Concretely, after the load task, compare source count and checksum to the destination partition, fail the DAG on mismatch, and never page only on task state.

The reason for that specificity is a failure I have seen: Spark exited 0 after writing 0 files because the input glob was empty. Airflow was green. 2 March had 0 new rows. Source had 2.1 million. The hole lasted 14 hours until a human queried.

Empty glob, green DAG.

Signal2 March rowsSource rows
Airflow success02,100,000
checksum gatefail2,100,000

I would not consider it settled without evidence: run a successful Spark job with an empty glob and require the DAG to fail on destination count 0 versus source 2.1 million.

Exit 0 is not a row count.

Curated: · Written: · Reviewed:

QA-18dbt run succeeded. unique tests were skipped because the model was ephemeral. What is in the table?(show answer)

The useful question for task success versus row-count contract is what a second consumer would read from the same sink after a retry.

A transform success is a compile-and-execute success. Contracts live on the relation consumers read. Skipping tests on ephemeral models, or testing a view that is not the served table, leaves the served table unchecked.

Concretely, run not-null, unique, and accepted-values on the served relation after it is built, fail the DAG on test failure, and forbid skipping tests because the model is ephemeral if a table was still materialized.

The reason for that specificity is a failure I have seen: Ephemeral staging was tested; the incremental table was not. 44,000 duplicate keys landed. unique on staging passed. The served table's COUNT(*) was 44,000 above COUNT(DISTINCT key).

Tests on ephemeral, duplicates in incremental.

Relationunique testDuplicate keys
ephemeral stagingpassed0
served incrementalskipped44,000

I would not consider it settled without evidence: materialize incremental, skip tests on it, and require the review to fail; passing means tests bind to the served relation.

Test the table people query.

Curated: · Written: · Reviewed:

QA-19When should dbt test run relative to the Iceberg commit that consumers will read?(show answer)

I would settle dbt tests as a load gate by replaying the failed path against the contract, not against a green Airflow square.

A test after consumers have already queried a bad snapshot is a report, not a gate. The contract has to refuse the snapshot (or swap the pointer) before the serving name moves.

Concretely, write to a staging snapshot, run tests, publish the snapshot id to the serving pointer only on pass, and leave the previous snapshot serving on fail.

The reason for that specificity is a failure I have seen: Tests ran 26 minutes after commit. Two BI extracts read 12 million null foreign keys. The tests then failed. The DAG was marked failed after the extracts had already loaded.

Test after the serving pointer moved.

OrderConsumer reads of bad snapshotTest result
commit, serve, then test2 extracts, 12,000,000 null FKsfail too late
test staging, then pointer swap0fail, old snapshot remains

I would not consider it settled without evidence: commit a violating snapshot and require 0 consumer reads of that snapshot id.

Gate the pointer, not the Slack message.

Curated: · Written: · Reviewed:

QA-20The producer bumped the Avro subject. The Spark job uses a frozen schema file in Git. Who wins?(show answer)

The judgement in schema registry in the job is which operator wrote the files and whether that operator is idempotent.

The job that decodes bytes must use a reader schema compatible with those bytes. A Git file that is older or newer than the registry subject without a compatibility check is a decode incident waiting for a backfill.

Concretely, fetch the reader schema from the registry at start, pin a subject version, run compatibility against the oldest retained files, and fail fast on decode error rather than writing null columns.

The reason for that specificity is a failure I have seen: Git schema lagged 4 versions. 3.6 million records decoded with missing fields filled from an outdated reader. channel was always "". Attribution was wrong for 9 days.

Frozen Git reader versus newer writer.

ReaderRecordschannel
Git, 4 versions behind3,600,000always ""
registry-compatible reader3,600,000producer values

I would not consider it settled without evidence: consume a newer writer schema with the frozen Git reader and require the job to fail compatibility, not to silently default.

The registry is the job's reader, not a wiki.

Curated: · Written: · Reviewed:

QA-21You add country to the Avro reader so old files still decode. Where does the default belong?(show answer)

Where candidates lose the interview on Avro reader-schema defaults is calling the DAG success the load.

Avro fills a field the writer never stored from the reader schema's default. Putting the default only on a new writer schema does not help a new reader decode old files that lack the field. A required reader field with no default cannot read history.

Concretely, add country on the reader with a default the business signed, prove the newest Spark job reads pinned historical files, and reject a required reader field with no default.

The reason for that specificity is a failure I have seen: country was required on the new reader with no default. Replay of 11 days, 5.2 million files, failed deserialize. A retry path wrote JSON in parallel and duplicated 310,000 facts.

Adding country for a replay.

Reader changeHistorical decodeDuplicate facts
required country, no reader default0 of 5,200,000310,000 via JSON path
country with reader default "UNSET"5,200,0000

I would not consider it settled without evidence: point the new reader at historical files without country and require decode success only when the reader default exists.

History is saved by the reader default.

Curated: · Written: · Reviewed:

QA-22A producer reuses Protobuf field 4 for a new string because the old int32 died. What does the Spark job decode from old files?(show answer)

I would answer Protobuf numbers in the pipeline by separating delivery, transform location, and the row set the sink is allowed to keep.

Protobuf field numbers are the contract. Reusing a number at a new wire type leaves old bytes as unknown fields, not as the new string. Reusing at the same wire type silently reinterprets values. The pipeline must reserve dead numbers.

Concretely, fail the job if the descriptor reuses a reserved number, decode a golden old file in CI, and never bind a warehouse column to a recycled number.

The reason for that specificity is a failure I have seen: Field 4 changed from int32 amount_cents to string sku. Old files produced empty sku and dropped amounts as unknown. Revenue printed $0 for 2 days on 7.4 million rows while row counts passed.

Field 4 recycled.

DescriptorOld int32 filesRevenue
field 4 now string skuamounts unknown, sku empty$0 for 2 days
field 4 reserved, sku is field 9amounts still int32 via old readercorrect

I would not consider it settled without evidence: decode a golden int32 file with the new descriptor and require the compatibility gate to fail before Spark is scheduled.

Reserve the number; do not recycle it.

Curated: · Written: · Reviewed:

QA-23Spark wrote 400 files then died. The next list_dir shows the files. Are they in the table?(show answer)

The engineering content of Iceberg commit atomicity is the merge key and the quality gate that would refuse the write, not the scheduler.

Iceberg readers see a snapshot. Orphan data files are not rows until a commit metadata pointer includes them. Listing the data directory is not a table read.

Concretely, read through the catalog snapshot, run remove_orphan_files on a delay that matches job retries, and never SUM files on disk as the table.

The reason for that specificity is a failure I have seen: On-call counted 400 parquet files in s3 and told finance the load landed. The snapshot still pointed at the previous list. Query returned yesterday. The DAG later committed a second try and duplicated nothing, which confused the incident review.

Files on disk versus snapshot.

ObservationRows a SELECT sees
400 parquet files in the prefixyesterday
Iceberg snapshot after commityesterday plus today's input

I would not consider it settled without evidence: kill Spark after file write and before commit, then SELECT count through Iceberg and require yesterday's count, not 400 extra files.

The snapshot is the table; the prefix is a junk drawer.

Curated: · Written: · Reviewed:

QA-24Is overwriting an Iceberg day partition inherently unsafe?(show answer)

Before calling Iceberg day overwrite completeness done I would write down the partition, consumer, or late window nobody checked.

A day overwrite is a valid reload when the input is the complete day. It is unsafe when a partial batch replaces every file in that day's slice. The format does not forbid overwrite; the incomplete input does.

Concretely, gate overwrite on source_count versus input_count for that day, use MERGE for partials, and test both a complete overwrite and a rejected partial.

The reason for that specificity is a failure I have seen: After hour-level files existed, a partial overwrite from hour=14 replaced the whole day. 41 million facts dropped. Time travel still had snapshot 1904. A complete-day overwrite the next night would have been the correct repair.

Partial versus complete day overwrite.

InputFacts afterRecoverable in snapshot 1904
1 hour41,000,000 missingyes
complete daymatch sourcen/a

I would not consider it settled without evidence: overwrite a day from a 1-hour input and require a fail; overwrite from a complete extract and require counts to match source.

Complete overwrite is a tool; partial overwrite is a delete.

Curated: · Written: · Reviewed:

QA-25Spark lists parquet in a Delta table path and unions them. Why is that not reading Delta?(show answer)

The first thing I would establish about Delta transaction log is which write actually landed, not which task box turned green.

Delta's transaction log names which files are live. Unioning every parquet in the directory includes removes, aborted jobs, and vacuum candidates. That is a replay of junk, not a table.

Concretely, read via Delta scan, vacuum on the retention clock, and forbid Spark glob of the table path in production jobs.

The reason for that specificity is a failure I have seen: A glob job unioned 1,140 remove-pending files. Duplicate keys were 2.2 million. SUM doubled. AVG stayed the same because duplicates were exact.

Glob versus Delta scan.

ReadDuplicate keysSUMAVG
glob parquet2,200,000doubledunchanged
Delta scan0correctcorrect

I would not consider it settled without evidence: remove a file in Delta then glob the path, and require the glob job to be rejected in review.

The log is the membership; the folder is not.

Curated: · Written: · Reviewed:

QA-26The table is partitioned by dt. The query filters DATE(event_ts)=dt. Why did Spark still scan 400 days?(show answer)

I would start Hive-style partition pruning from the destination checksum and the failed partition, not from the DAG run.

Partition pruning needs a predicate on the partition column itself. Wrapping the column, or filtering a different timestamp column that happens to equal dt, often disables prune and reads every directory.

Concretely, filter dt = '2026-03-02' (or a range on dt), generate dt at write time from event_ts, and EXPLAIN to confirm partitionFilters is non-empty.

The reason for that specificity is a failure I have seen: DATE(event_ts) = '2026-03-02' scanned 412 partitions, 18 TB, for 74 minutes. dt = '2026-03-02' scanned 1 partition, 44 GB, in 3 minutes.

Predicate on event_ts versus dt.

FilterPartitions scannedBytesWall time
DATE(event_ts) = day41218 TB74 minutes
dt = day144 GB3 minutes

I would not consider it settled without evidence: require EXPLAIN to show a single partition for the daily job before merge to main.

Prune the column you partitioned on.

Curated: · Written: · Reviewed:

QA-27A streaming job writes every 30 seconds into a daily partition. What happens to the next day's batch?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Each commit can create files smaller than a row group. Thousands of tiny files make listing and opening dominate scan time. Compaction is part of the write path's SLO, not a weekend hobby.

Concretely, target 128–512 MB files, compact on a schedule tighter than the read SLO, and use a streaming trigger that batches enough records to hit the size.

The reason for that specificity is a failure I have seen: 24 hours of 30-second commits left 2,880 files in one day, not 5,760, median 2.1 MB. The morning scan opened files for 22 minutes before any aggregate. After binpack to 128 MB, about 47 files, scan was 90 seconds.

30-second commits versus compacted day.

StateFilesMedian sizeScan
24 hours uncompacted2,8802.1 MB22 minutes open
binpack 128 MB~47~128 MB90 seconds

I would not consider it settled without evidence: count files and p50 size in the partition and require compaction before the dashboard pack if file count exceeds 200 per day.

Tiny files are a write SLO miss.

Curated: · Written: · Reviewed:

QA-28MOR deletes pile up. When must compaction run relative to the 07:00 pack?(show answer)

My answer to compaction SLO begins with where extract, load, and transform actually ran, and on which clock.

Merge-on-read makes writes cheap and reads pay. If the pack queries before compaction, it pays yesterday's CDC as delete-file fanout. Compaction is on the critical path of the read SLO, not after it.

Concretely, compact until delete files per partition are under a budget, then run the pack, and page if compaction misses the pack by more than 10 minutes.

The reason for that specificity is a failure I have seen: Compaction was scheduled 07:30. The pack at 07:00 scanned 3,900 delete files. Slot-hours were 8× and the pack missed the meeting.

Pack versus compaction order.

OrderDelete files in pack querySlot-hours
pack 07:00, compact 07:303,9008×
compact then packtens1×

I would not consider it settled without evidence: after 24 hours of CDC, query the pack query and require delete-file count under budget before 07:00.

Compact before the readers you promised.

Curated: · Written: · Reviewed:

QA-29spark.sql.shuffle.partitions is still 200 on a 12 TB join. What are you waiting for?(show answer)

I would treat Spark shuffle partitions as a claim about every partition and every late window, not about the one that was demonstrated.

Each shuffle partition is a task. 200 tasks on 12 TB is ~60 GB per task, which spills and times out. 4,000 partitions is 12 TB / 4000 ≈ 3 GB each, which is not a 128–256 MB shuffle block; that target needs about 49,000–98,000 partitions. Too many partitions on a small job wastes scheduling. The number follows data size, not the default.

Concretely, size partitions toward 128–256 MB shuffle blocks (~49k–98k on 12 TB), enable AQE where it actually coalesces, and measure spill bytes, not just wall time.

The reason for that specificity is a failure I have seen: 200 partitions spilled 4.8 TB to disk. The join ran 6.2 hours. At 4,000 partitions each task still held ~3 GB, spill dropped to 110 GB, and wall time to 48 minutes — better, but still not the 128–256 MB target.

200 versus 4,000 versus a 128–256 MB target on 12 TB.

shuffle.partitionsApprox blockSpillWall time
200~60 GB4.8 TB6.2 hours
4,000~3 GB110 GB48 minutes
~49,000–98,000128–256 MBsized to targetsized to block

I would not consider it settled without evidence: report shuffle spill for the join and require it under a budget before calling the DAG healthy.

200 is a default, not a cluster size.

Curated: · Written: · Reviewed:

QA-30One customer_id has 8 percent of events. The join hangs on two tasks. What is the actual problem?(show answer)

The useful question for Spark join skew is what a second consumer would read from the same sink after a retry.

Skew is a key that hashes to a partition far larger than the rest. Salting, skew join hints, or broadcasting the small side fix it. Adding executors does not split a single skewed partition.

Concretely, measure max partition size versus median, salt the hot key or enable skew join, and never scale the cluster as the first response to two slow tasks.

The reason for that specificity is a failure I have seen: Two tasks ran 11 hours on a 90-minute job. The hot key had 8 percent of 2.4 billion events. Doubling executors left those two tasks unchanged.

Hot customer_id.

PlanSlow tasksWall time
2× executors2 at 11 hours11 hours
salt hot key090 minutes

I would not consider it settled without evidence: show the key histogram and a salted plan that finishes within 1.3× the median task.

More boxes do not split one key.

Curated: · Written: · Reviewed:

QA-31The dimension is 2.1 GB. autoBroadcastJoinThreshold is 10 MB. Why did the fact join shuffle 9 TB?(show answer)

I would settle broadcast versus shuffle join by replaying the failed path against the contract, not against a green Airflow square.

Broadcast ships a copy of the small side to every executor. Above threshold Spark shuffles both sides. A 2.1 GB dimension may still be worth broadcasting if executors have RAM; it is not free, and OOM is the failure mode.

Concretely, broadcast when the small side fits in executor memory with headroom, set the hint explicitly, and watch peak executor memory.

The reason for that specificity is a failure I have seen: Shuffle join of 9 TB facts with 2.1 GB dim ran 4 hours. Broadcast of the dim finished in 22 minutes. A later 18 GB dim broadcast OOM-killed 40 executors; that one needed a shuffle.

2.1 GB dimension.

JoinWall timeExecutor RAM
shuffle both sides4 hoursfine
broadcast 2.1 GB22 minutes+2.1 GB
broadcast 18 GBOOM 40 executorsexceeded

I would not consider it settled without evidence: compare wall time and peak RAM for broadcast versus shuffle on the actual dim size.

Broadcast is a memory bet, not a default.

Curated: · Written: · Reviewed:

QA-32You enabled AQE and the job got slower. What did coalescing do to your skewed join?(show answer)

The judgement in Spark AQE is which operator wrote the files and whether that operator is idempotent.

Adaptive Query Execution can coalesce small partitions and switch join strategies mid-flight. It can also hide a bad static partition count until runtime, or coalesce away parallelism you needed. It is a tool to measure, not a piety.

Concretely, compare spill, task times, and join strategy with AQE on and off for the job, and pin a plan when AQE regresses.

The reason for that specificity is a failure I have seen: AQE coalesced a 4,000-partition shuffle down to 280 after a filter. The remaining join was skewed. Runtime went from 51 minutes to 3.4 hours.

AQE coalesce into skew.

AQEShuffle partitions after filterWall time
off4,00051 minutes
on, coalesced2803.4 hours

I would not consider it settled without evidence: run the same input with AQE on and off and keep the faster plan in the DAG.

Adaptive is a measurement, not a slogan.

Curated: · Written: · Reviewed:

QA-33A Python note does df.collect() on 14 million rows to "take a look." What dies?(show answer)

Where candidates lose the interview on driver collect OOM is calling the DAG success the load.

collect pulls the dataset onto the driver JVM/Python process. The cluster can hold 14 million rows while the driver cannot. take, toPandas with a limit, or write-then-sample are the inspection tools.

Concretely, forbid collect in production DAGs, cap interactive collects, and write samples to a table instead.

The reason for that specificity is a failure I have seen: collect of 14 million rows, 2.2 KB each, needed ~30 GB on an 8 GB driver. The DAG failed after 3 hours of successful executors. Airflow retried 4 times, same OOM.

collect 14 million rows.

ActionDriver RAMResult
collect 14e6 × 2.2 KB~30 GBOOM on 8 GB driver, 4 retries
write sample 10,000smallsucceeds

I would not consider it settled without evidence: replace collect with a 10,000 row sample write and require the driver heap to stay under 50 percent.

The driver is not a warehouse.

Curated: · Written: · Reviewed:

QA-34pandas.read_sql pulls 9 million rows into one DataFrame on the extractor box. What is the production pattern?(show answer)

I would answer pandas chunked extract by separating delivery, transform location, and the row set the sink is allowed to keep.

A single in-memory frame is bounded by the extractor's RAM. Chunked or Spark JDBC reads stream rows to files. Production extracts do not need a global DataFrame.

Concretely, read in chunks or via Spark, write Parquet as you go, and cap the extractor's RSS.

The reason for that specificity is a failure I have seen: One frame of 9 million × 18 columns peaked at 41 GB. The box had 32 GB. The extract OOM'd at 92 percent, retried hourly, and never landed Saturday.

One pandas frame versus chunks.

ReadPeak RSSSaturday landed
one DataFrame41 GBno
chunks of 100,0003.1 GByes

I would not consider it settled without evidence: chunk at 100,000 rows and require peak RSS under 4 GB for the same extract.

If it does not fit in RAM, it was never one frame.

Curated: · Written: · Reviewed:

QA-35A Python UDF parses timestamps because to_timestamp "looked fussy." What did you pay?(show answer)

The engineering content of Spark UDF versus native is the merge key and the quality gate that would refuse the write, not the scheduler.

Python UDFs disable whole-stage codegen, serialize to Python, and run single-row. Native timestamp parsing stays in the JVM. Fussy formats still belong in native functions or a compiled UDF, not CPython per row.

Concretely, use to_timestamp / unix_timestamp, and measure rows/s native versus Python before keeping a UDF.

The reason for that specificity is a failure I have seen: Python strptime on 6.8 billion rows ran 11 hours. Native to_timestamp ran 17 minutes. The DAG SLA was 2 hours.

strptime UDF versus to_timestamp.

ParserRowsWall time
Python strptime UDF6,800,000,00011 hours
native to_timestamp6,800,000,00017 minutes

I would not consider it settled without evidence: benchmark 10 million rows both ways and require native unless the UDF is the only parser that exists.

Per-row Python is a last resort.

Curated: · Written: · Reviewed:

QA-36ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY ts) is used to "dedupe." Why can two runs keep different current rows?(show answer)

Before calling SQL window grain done I would write down the partition, consumer, or late window nobody checked.

ROW_NUMBER assigns exactly one rn=1 per partition even when timestamps tie. The failure is nondeterministic selection: without a unique tiebreaker, which tied row gets rn=1 is undefined and can change between runs. QUALIFY rn=1 without a total order is a coin flip, not two rn=1 values.

Concretely, partition by the business key, order by event_ts, event_id, and assert 1 current row per key after the filter, with the same survivor on rerun.

The reason for that specificity is a failure I have seen: 14 percent of users had tied timestamps. Each run kept exactly one rn=1, but a later rerun swapped 640,000 survivors. The serving table drifted. Downstream unique tests failed only on a 1 percent sample and stayed green.

Tied timestamps.

ORDER BYrn=1 per userSurvivor across reruns
ts onlyexactly 1can swap 640,000
ts, event_idexactly 1stable

I would not consider it settled without evidence: insert two rows with identical (user_id, ts) and require the same survivor via event_id on two runs.

Ties still yield one rn=1; the question is which row.

Curated: · Written: · Reviewed:

QA-37orders.items is an array. You explode then SUM(amount). What grain did you just change?(show answer)

The first thing I would establish about JSON array explode is which write actually landed, not which task box turned green.

Explode fans one parent row into N child rows. Additive parent measures repeat N times unless you keep them only on the parent or allocate. The grain after explode is the item, not the order.

Concretely, explode to item grain, sum item amounts, and never sum a parent amount after explode without dividing by array size.

The reason for that specificity is a failure I have seen: Parent order_total was summed after explode. Orders with 4 items quadrupled GMV. Phantom GMV was $22 million in a week, 11 percent of reported revenue.

Summing order_total after explode.

Measure after explode4-item order contributionWeekly phantom GMV
SUM(order_total)×4$22 million
SUM(item_price)×1$0

I would not consider it settled without evidence: explode a 4-item order and require SUM(order_total) to be rejected; SUM(item_price) should equal the parent once.

Explode changes the grain; parent measures do not automatically shrink.

Curated: · Written: · Reviewed:

QA-38A new optional nested field appears in 8 percent of JSON. Spark inferred the schema from the first file. What happens to the rest?(show answer)

I would start nested schema drift from the destination checksum and the failed partition, not from the DAG run.

Schema-on-read from a sample freezes a schema that later files can violate or silently drop. Merge schema or a registry contract has to include the nested field before those files land in a typed table.

Concretely, merge schema on read or enforce a declared schema, fail on unknown required paths, and land leftovers in a variant column with a quota.

The reason for that specificity is a failure I have seen: First-file inference omitted payment.network. 8 percent of files dropped the field. Card mix was wrong for 6 days, 410,000 rows, while Spark jobs succeeded.

Inferred schema from file 1.

Readpayment.network presentRows affected
infer file 10 in later files410,000
merged/declared schemayes0 missing

I would not consider it settled without evidence: add a nested field in file 2 and require it to appear in the table or the job to fail, not to disappear.

The first file is not the contract.

Curated: · Written: · Reviewed:

QA-39How does a pipeline apply type 2 without rewriting the whole dimension every night?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Type 2 closes the current row when tracked attributes change, inserts a new current row, and leaves unchanged members untouched. A MERGE that only expires via WHEN MATCHED never inserts the successor: matched keys take the expire branch, and WHEN NOT MATCHED INSERT fires only for brand-new natural keys. Successors must be staged as extra source rows or inserted in a second statement.

Concretely, stage both the expire image and the new current row, MERGE (or MERGE then INSERT) so successors land, and test that unchanged members keep their surrogate and dates.

The reason for that specificity is a failure I have seen: MERGE WHEN MATCHED expired 18,000 current rows. WHEN NOT MATCHED inserted only 210 brand-new customers. The 18,000 successors never landed. A panic nightly reload then assigned new surrogates to all 4.1 million customers. Facts still pointed at old surrogates. 96 percent of facts joined unknown for 2 days.

Expire-only MERGE versus staged successor insert.

JobSuccessors for 18,000 changesFacts joining unknown
MERGE expire only0 (not-matched never fires)members have no current row
staged expire + insert18,000 new versions0
reload all surrogatesn/a96 percent for 2 days

I would not consider it settled without evidence: change 1 member, rerun MERGE, and require 4,099,999 members to keep surrogate and valid_from.

Type 2 is a delta of versions, not a nightly rebirth.

Curated: · Written: · Reviewed:

QA-40Email corrections must land tonight. Why is a type-1 MERGE still a pipeline decision?(show answer)

My answer to SCD type 1 overwrite in the job begins with where extract, load, and transform actually ran, and on which clock.

Type 1 overwrites the current attribute and restates history for every fact that joins current. The pipeline must name that restatement, run it off a correction feed, and not mix it with type-2 attributes in the same UPDATE.

Concretely, split correction columns into a type-1 MERGE, keep historical attributes on type-2, and log how many facts will reread the new email.

The reason for that specificity is a failure I have seen: A type-1 email fix ran in the same UPDATE as region. Region restated 14 months. Bonus calculations moved $3.2 million between regions with 0 moves in the source.

Email correction bundled with region.

UPDATEEmailRegion historyDollars moved
mixedcorrectedrestated 14 months$3.2 million
type-1 email onlycorrectedheld$0

I would not consider it settled without evidence: correct email only and require region versions to stay put.

One UPDATE is not one SCD type.

Curated: · Written: · Reviewed:

QA-41Raw landing is "temporary." Analysts still have s3:GetObject. What is the incident?(show answer)

I would treat PII in the landing zone as a claim about every partition and every late window, not about the one that was demonstrated.

Landing is a copy. If it holds national ids in plaintext, it is a production PII store. Temporary is not a control. Redact or tokenise on ingest, or lock the prefix to pipeline roles only.

Concretely, strip or tokenise identifiers in the ingest job, encrypt the prefix, deny analyst GetObject, and scan landing with the same classifier as curated.

The reason for that specificity is a failure I have seen: Landing held 22 million plaintext national ids for 18 days "until the curated job caught up." An intern sync'd the prefix to a laptop. Legal counted 22 million records in scope.

Plaintext landing.

ControlNational ids in landingDaysScope
temporary prefix, analyst read22,000,0001822 million
tokenise on ingest, deny GetObject0 plaintext00

I would not consider it settled without evidence: attempt analyst GetObject on landing and require deny, plus a classifier count of 0 plaintext ids.

Landing is in production the moment it has PII.

Curated: · Written: · Reviewed:

QA-42DELETE FROM customers WHERE subject_id=? returned 1. Why can legal still find the email?(show answer)

The useful question for GDPR erasure in the lake is what a second consumer would read from the same sink after a retry.

A current-row DELETE does not rewrite historical files, expire snapshots, or remove orphan data files. Erasure in a lakehouse is rewrite without identifiers, snapshot expiry so AS OF cannot return them, and physical file removal.

Concretely, rewrite files without the identifiers, expire snapshots that referenced the old files, delete orphans, and prove AS OF the old snapshot id fails expired, not masked, by decoding remaining data files rather than listing object keys.

The reason for that specificity is a failure I have seen: DELETE hit 1 current row. Snapshot 441 still served the email for 11 days. Orphan files held 7 prior versions. Legal listed prefixes, saw fewer keys, and declared erasure; a decode of remaining Parquet still returned the email.

DELETE versus rewrite plus expiry plus remove.

StepEmail queryable
DELETE currentyes in snapshot 441 and orphans, 11 days
rewrite, expire, remove filesno, AS OF expired

I would not consider it settled without evidence: after the erasure job, AS OF pre-erasure snapshot must error expired, and a decode/scan of remaining data files must show 0 copies of the email. Listing object keys cannot prove file contents.

DELETE current is not erasure.

Curated: · Written: · Reviewed:

QA-43One malformed JSON blocks a Spark micro-batch of 200,000. Where should that record go?(show answer)

I would settle dead-letter poison messages by replaying the failed path against the contract, not against a green Airflow square.

A poison record should land in a DLQ with offset, bytes, and error, and the batch should commit the rest. Failing the whole micro-batch retries the poison forever and stalls the watermark.

Concretely, catch parse errors per record, write DLQ, continue, page on DLQ depth, and replay DLQ after a schema fix with the same event ids.

The reason for that specificity is a failure I have seen: One truncated JSON failed the batch. Spark retried 14 days. Lag grew to 61 million. The DAG was "running." Downstream freshness SLA missed every hour.

One bad JSON in 200,000.

HandlingSink rowsDLQLag after 14 days
fail the micro-batch0061,000,000
per-record DLQ199,99910 extra

I would not consider it settled without evidence: inject one bad record into 200,000 and require 199,999 in the sink and 1 in DLQ.

Poison is a record, not a reason to stall the fleet.

Curated: · Written: · Reviewed:

QA-44The SLA is "hourly." The job starts at minute 50 and runs 40 minutes. Did you meet it?(show answer)

The judgement in freshness SLA is which operator wrote the files and whether that operator is idempotent.

Freshness is when the destination partition is queryable with the contract, not when the DAG started. A job that starts at :50 and finishes at :30 past the next hour missed the hour it claimed.

Concretely, measure time from event-time close to serving pointer update, keep headroom for runtime p99, and page on pointer lag, not on DAG start.

The reason for that specificity is a failure I have seen: Hour 09 was served at 10:31. The SLA said 10:00. 14 of 24 hours missed. Airflow still showed success because the task had no sla_miss callback on the pointer.

Start versus serve.

Hour 09ClockSLA 10:00
DAG start09:50—
serving pointer10:31miss
hours missed in a day14 of 24—

I would not consider it settled without evidence: track serving pointer timestamps and require p95 pointer lag inside the SLA.

Started on time is not served on time.

Curated: · Written: · Reviewed:

QA-45Row count is 99.2 percent of source. The freshness SLA is green. Can you publish?(show answer)

Where candidates lose the interview on completeness SLA is calling the DAG success the load.

Freshness without completeness is a fast hole. The contract names a count or checksum tolerance. Publishing at 99.2 percent when the tolerance is 99.9 percent is a miss, even if the hour arrived on time.

Concretely, compare source and destination counts per partition, fail below tolerance, and do not let freshness alone flip the serving pointer.

The reason for that specificity is a failure I have seen: 09:00 arrived at 09:08 with 99.2 percent of rows. Missing 0.8 percent was $1.6 million. The freshness tile was green.

Fresh and short.

09:00Count vs sourceDollars missingFreshness tile
published99.2 percent$1.6 milliongreen
heldwait for 99.9$0 published holeyellow

I would not consider it settled without evidence: hold the pointer until count ≥ 99.9 percent or a documented exception.

On time and incomplete is still incomplete.

Curated: · Written: · Reviewed:

QA-46Incremental is hard, so the team rebuilds 4 years every night. What is the bill you should put on the design review?(show answer)

I would answer cost of full refresh by separating delivery, transform location, and the row set the sink is allowed to keep.

A full refresh rereads and rewrites history the incremental job would not touch. Cost is slot-hours, object-storage writes, and the lock on partitions that incremental also needs. It is justified for a true rebuild, not as a substitute for a cursor.

Concretely, price 4-year rewrite versus incremental plus weekly reconcile, and require a cursor before nightly full refresh of multi-year facts.

The reason for that specificity is a failure I have seen: Nightly full refresh of 4 years scanned 2.1 PB, $18,400/day. Incremental plus 0.4 percent reconcile would have been $610/day. The DAG was "simpler."

Nightly 4-year rebuild.

ModeScanCost/day
full refresh2.1 PB$18,400
incremental + weekly reconcile~8 TB$610

I would not consider it settled without evidence: show 7-day cost of full versus incremental on the same table before approving the simpler DAG.

Simple and 2.1 PB is not cheap.

Curated: · Written: · Reviewed:

QA-47The source has no updated_at. You extract WHERE id > last_max_id. Which rows do you never see?(show answer)

The engineering content of incremental extract key is the merge key and the quality gate that would refuse the write, not the scheduler.

A monotonic id cursor misses in-place updates to old ids and late inserts with ids below the cursor if the source reuses or backfills ids. It only captures append-only increasing ids.

Concretely, use id cursors only on append-only tables, otherwise CDC or a dump, and test an update to id=12 after the cursor passed 12.

The reason for that specificity is a failure I have seen: Prices were updated in place on old product ids. The extractor never refetched them. 16,200 SKUs stayed at last year's price for 11 weeks, $4.7 million of wrong margin.

In-place update behind the id cursor.

ExtractOld id updatesSKUs staleWeeks
id > last_maxmissed16,20011
CDC or dumpcaught00

I would not consider it settled without evidence: update an old id and require the pipeline to pick it up via CDC or a dump, not via max(id).

max(id) is an append cursor.

Curated: · Written: · Reviewed:

QA-48dt is computed as DATE(event_ts) in the Spark session TZ. Executors are UTC. Airflow is US/Eastern. Which day owns 23:30 Eastern?(show answer)

Before calling timezone of partition columns done I would write down the partition, consumer, or late window nobody checked.

Partition dates must be computed in a named timezone, usually the business calendar, from a stored UTC timestamp. Under a UTC session the offset only bites in the Eastern evening: 23:30 Eastern is 04:30 UTC the next calendar date, while 00:30 Eastern is 05:30 UTC on the same date and looks correct.

Concretely, store event_ts in UTC, compute dt with from_utc_timestamp(..., 'America/New_York') or equivalent, and unit-test the 23:30 Eastern boundary, not only 00:30.

The reason for that specificity is a failure I have seen: 23:30 Eastern in January landed in UTC dt the next calendar date. 1.1 million orders sat in the wrong finance day for a month, $8.2 million shifted between closes. The suite only asserted 00:30 Eastern, which is on the same UTC date and passed.

23:30 Eastern in January.

dt computationPartitionOrders misplaced
DATE(event_ts) UTC sessionnext UTC date1,100,000
DATE in America/New_Yorkbusiness date0

I would not consider it settled without evidence: convert a stamped 23:30 America/New_York instant and require dt equal to that local date, not the UTC date.

The partition is a calendar, not the executor's default.

Curated: · Written: · Reviewed:

QA-49On the spring-forward night, the 02:00 hour does not exist. Your hourly job still has a 02:00 partition. What do you write?(show answer)

The first thing I would establish about DST around midnight loads is which write actually landed, not which task box turned green.

DST gaps and overlaps are calendar facts. A job keyed on local hour must skip the gap and disambiguate the overlap (usually with UTC). Inventing a 02:00 Eastern in March is an empty partition or a duplicate hour in November.

Concretely, schedule on UTC, label partitions with UTC or with a DST-aware local hour, and test both transition dates.

The reason for that specificity is a failure I have seen: The 02:00 Eastern job in March wrote 0 rows and marked success. In November the overlapping 01:00 ran twice and duplicated 84,000 events.

US spring-forward and fall-back.

Night02:00 local jobEvents
spring-forwardsuccess, 0 rowshole
fall-back naive 01:00 twice2 writes84,000 dupes
UTC hourly1 write per UTC hour0 extra

I would not consider it settled without evidence: run the DAG on both transition timestamps in a clock test and require 0 invented gap hours and 1 write per UTC hour in the overlap.

Local hours are not a uniform timeline.

Curated: · Written: · Reviewed:

QA-50The vendor file grows new columns every quarter without notice. Your Spark job selects * into a typed table. What happens in April?(show answer)

I would start slowly changing source extracts from the destination checksum and the failed partition, not from the DAG run.

SELECT * from a drifting file maps extra columns into the table only if merge schema is on; otherwise Spark drops or fails. Typed sinks need an allow-list plus a remainder column, not star.

Concretely, project named columns, land unknown columns in a variant remainder with a size quota, and alert on remainder growth.

The reason for that specificity is a failure I have seen: April added discount_bps. SELECT * with a frozen schema dropped it. Margin reports ignored 9 percent discounts for 7 weeks, $6.4 million.

New vendor column.

Readdiscount_bpsMargin error
SELECT * frozen schemadropped$6.4 million, 7 weeks
named plus remainderretained0 silent

I would not consider it settled without evidence: add a column in a fixture file and require it in remainder or in a versioned schema, not silence.

Star is how columns disappear.

Curated: · Written: · Reviewed:

QA-51The source 429s. Airflow retries every task every 30 seconds with 16 parallelism. What did you just do to the source and to your sink?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Unbounded retries amplify load and, without idempotent sinks, duplicate writes. Backoff, jitter, and a shared rate limit are the extract policy. Retrying harder is an attack on your own vendor.

Concretely, exponential backoff with jitter, cap parallelism, idempotent MERGE, and a circuit breaker after N 429s.

The reason for that specificity is a failure I have seen: 16 tasks retrying every 30 seconds produced 32 attempts/minute, not 8,400. The vendor still 429'd under catch-up. Catch-up then appended 3 copies of Friday, 12.6 million extra rows. COUNT(DISTINCT order_id) did not triple.

Retry storm then catch-up append.

BehaviourRequests/minExtra rowsDistinct ids
16× every 30s then append32 then extra copies12,600,000unchanged
backoff plus MERGEunder cap0unchanged

I would not consider it settled without evidence: simulate 429s and require request rate under the vendor cap and 1 sink row per id after catch-up.

Retries need backoff and a merge key.

Curated: · Written: · Reviewed:

QA-52catchup=True after the DAG was paused for 9 days. What runs, and in what order relative to today's incremental?(show answer)

My answer to Airflow catchup begins with where extract, load, and transform actually ran, and on which clock.

Catchup schedules every missed interval. Those runs can overlap today's writer on the open partition, replay old intervals out of CDC order, or stampede the source. Catchup is a backfill plan, not a default.

Concretely, leave catchup off for incremental CDC, backfill with a bounded date list and partition locks, and never let 9 days of intervals start in parallel on the same table.

The reason for that specificity is a failure I have seen: Catchup launched 216 hourly tasks. They raced today's job. 31 percent of keys duplicated on the open day; 9-day-old CDC applied after newer CDC and reverted status.

Nine days of catchup.

ModeParallel intervalsOpen-day duplicatesStatus reverted
catchup=True21631 percentyes
bounded backfill with lock10no

I would not consider it settled without evidence: pause 9 days then enable catchup in staging and require a lock or catchup=False plus a planned backfill.

Missed intervals are a backfill, not a fork bomb.

Curated: · Written: · Reviewed:

QA-53A sensor waits for s3://.../dt=today/_SUCCESS. The upstream wrote 0-byte _SUCCESS because the extract was empty. What do you load?(show answer)

I would treat sensor on empty partition as a claim about every partition and every late window, not about the one that was demonstrated.

A success marker is not a row count. Sensors that only wait for a file will unblock a zero-row partition. The sensor has to wait for the contract (count, checksum) or the downstream must fail the empty-when-source-nonempty case.

Concretely, sensor on a manifest with counts, or fail the load if destination is empty while source count > 0.

The reason for that specificity is a failure I have seen: _SUCCESS landed at 02:01 with 0 objects. Downstream loaded 0. Source had 3.4 million. The hole lasted until noon, $5.9 million of Saturday missing.

Empty _SUCCESS.

MarkerObjectsSource rowsLoaded
0-byte _SUCCESS03,400,0000
manifest with countmust match3,400,0003,400,000 or fail

I would not consider it settled without evidence: write _SUCCESS with 0 files while source count is 3.4 million and require the downstream DAG to fail.

_SUCCESS is a file, not a census.

Curated: · Written: · Reviewed:

QA-54A 14-month backfill will run all weekend. When do uniqueness and not-null tests run?(show answer)

The useful question for backfill quality gates is what a second consumer would read from the same sink after a retry.

Gates after the backfill has overwritten 14 months are an autopsy. Each partition (or a rolling window) has to pass tests before the serving pointer for that partition moves. Deferring tests until Monday is how you publish 14 months of duplicates.

Concretely, test per partition before pointer swap, stop the backfill on first failed partition, and keep serving the previous snapshot for failed months.

The reason for that specificity is a failure I have seen: Tests ran Monday. 6 months had duplicate keys, 88 million extra rows. SUM doubled on those months; AVG did not, because duplicates were exact. The weekend DAG was green each night.

Weekend backfill, Monday tests.

When tests ranMonths duplicatedExtra rowsAVG
Monday688,000,000unchanged
per partition before pointer0 published0—

I would not consider it settled without evidence: fail a uniqueness test on month 2 of a staged backfill and require months 3–14 to not overwrite serving.

Do not wait for the whole history to finish before you look.

Curated: · Written: · Reviewed:

QA-55You have 12,000 changed keys in a 40 million row day. Which operator, and why is overwrite still not free?(show answer)

I would settle MERGE versus INSERT OVERWRITE by replaying the failed path against the contract, not against a green Airflow square.

MERGE touches changed files (or rewrite those files). Overwrite rewrites the whole partition. For 12,000 keys, MERGE is the cost model unless the engine would rewrite almost every file anyway.

Concretely, choose MERGE when the change ratio is small, overwrite when the input is a complete partition cheaper than merge, and measure files rewritten.

The reason for that specificity is a failure I have seen: Overwrite rewrote 40 million rows, 1.2 TB, to apply 12,000 updates. Slot-hours were 19× a MERGE that rewrote 4 files.

12,000 updates in 40 million rows.

OperatorRows rewrittenFilesSlot-hours
INSERT OVERWRITE day40,000,000all19×
MERGE~12,000 plus file rewrite41×

I would not consider it settled without evidence: compare files rewritten for 12,000 keys under both operators on a clone.

Overwrite is for complete days, not for 12,000 keys.

Curated: · Written: · Reviewed:

QA-56Spark dynamic partition overwrite is on. The input only has dt=2026-03-02 and dt=2026-03-03. What happens to 2026-03-01?(show answer)

The judgement in dynamic partition overwrite is which operator wrote the files and whether that operator is idempotent.

Dynamic overwrite replaces only partitions present in the input. Static overwrite can replace the whole table. Mixing them, or assuming dynamic will clear a missing day, leaves stale partitions that look like they were "reloaded."

Concretely, name the overwrite mode, list partitions in the input, and delete stale days with an explicit operation if that is the intent.

The reason for that specificity is a failure I have seen: A job intended to replace all of March sent only two days. Dynamic overwrite left 29 stale days. Finance summed a mix of new and 3-week-old partitions, $14 million off.

Two days in a dynamic overwrite.

Mode2026-03-01Finance error
dynamic, input 02 and 03stale$14 million
explicit drop of stale March daysgone or replaced0

I would not consider it settled without evidence: run dynamic overwrite with two days and require 2026-03-01 to still be the old snapshot unless an explicit drop ran.

Absent from the input is not deleted unless you said so.

Curated: · Written: · Reviewed:

QA-57Two Spark jobs MERGE the same Iceberg table. One fails with a commit conflict. What should the loser do?(show answer)

Where candidates lose the interview on Iceberg commit conflicts is calling the DAG success the load.

Iceberg serialises snapshot commits. A conflict means the loser's snapshot base is stale. Retry the MERGE against the new snapshot; do not INSERT OVERWRITE the day to "win," which deletes the winner's rows if the input is stale.

Concretely, retry MERGE with exponential backoff, isolate writers by partition when possible, and forbid overwrite-as-conflict-resolution.

The reason for that specificity is a failure I have seen: The loser overwrote the day from an 08:00 snapshot after the winner had merged 11:00 CDC. 140,000 midday orders disappeared.

Loser overwrites after conflict.

Loser actionMidday CDCOrders lost
overwrite from 08:00 snapshotdeleted140,000
retry MERGEkept0

I would not consider it settled without evidence: run two MERGEs and require the loser to retry MERGE, with both change sets present.

Conflict is a retry, not an overwrite.

Curated: · Written: · Reviewed:

QA-58Executors die with java.lang.OutOfMemoryError during a groupBy. The driver is fine. Where do you look first?(show answer)

I would answer Spark executor OOM by separating delivery, transform location, and the row set the sink is allowed to keep.

Executor OOM is partition size, nested explosion, or broadcast of something huge. Driver OOM is collect. Raising driver memory does not fix executor death.

Concretely, inspect max partition size, explode cardinality, and broadcast sizes, then repartition or filter before aggregate.

The reason for that specificity is a failure I have seen: On-call raised driver memory from 8 GB to 32 GB. Executors still died on a 6 GB partition after JSON explode to 40 fields × 200 array items. The job never finished. Wall clock wasted: 7 hours.

groupBy after explode.

ChangeWho OOM'dFinished
driver 8→32 GBexecutorsno, 7 hours wasted
cap explode / repartitionneitheryes

I would not consider it settled without evidence: show the dying executor's partition size and a plan that caps explode or repartitions before the aggregate.

The driver heap is the wrong knob.

Curated: · Written: · Reviewed:

QA-59Disk spill is 2.4 TB and the job still completes. Is that success?(show answer)

The engineering content of shuffle spill is the merge key and the quality gate that would refuse the write, not the scheduler.

Completion with multi-TB spill is a cost and an SLA risk. Spill is the engine surviving a bad partition size. The fix is partition count, skew, or filter pushdown, not celebrating exit 0.

Concretely, alert on spill bytes per job, budget spill as a fraction of shuffle write, and fail CI if a golden job spills above the budget.

The reason for that specificity is a failure I have seen: Nightly join spilled 2.4 TB, cost $1,100 extra, and finished at 06:50, 10 minutes before the pack. A 2× partition change cut spill to 80 GB and finish to 05:10.

2.4 TB spill.

Shuffle sizingSpillFinishExtra cost
default2.4 TB06:50$1,100
2× partitions80 GB05:10~0

I would not consider it settled without evidence: report spill on the golden job and require it under budget.

Spilled success is still a mis-sized job.

Curated: · Written: · Reviewed:

QA-60You replayed Friday's Kafka topic into an append table. COUNT(*) went up 3×. What should COUNT(DISTINCT user_id) do if users were unchanged?(show answer)

Before calling COUNT DISTINCT under replay done I would write down the partition, consumer, or late window nobody checked.

Exact replay of the same users into an append table multiplies rows, not distinct users. COUNT(DISTINCT user_id) stays put. Using COUNT(*) as a people metric is the bug; the distinct aggregate is doing its job.

Concretely, use MERGE on event or user grain for facts, and never treat COUNT(*) movement on replay as a user-growth signal.

The reason for that specificity is a failure I have seen: On-call saw COUNT() ×3 and announced 3× users. Distinct users stayed 2.1 million. A growth dashboard that used COUNT() shipped a false 200 percent spike.

Triple append of Friday.

MetricAfter 3× replay
COUNT(*)×3
COUNT(DISTINCT user_id)2,100,000 unchanged

I would not consider it settled without evidence: replay the same events 3× and require COUNT(DISTINCT user_id) unchanged and COUNT(*) explained as duplicates.

Replay does not mint people.

Curated: · Written: · Reviewed:

QA-61Duplicates doubled every row. Finance asks why AVG(order_value) did not halve. What do you say?(show answer)

The first thing I would establish about AVG under exact duplicates is which write actually landed, not which task box turned green.

AVG is SUM/COUNT. Exact duplicates scale SUM and COUNT by the same factor, so AVG is unchanged. Halving would require extra rows with different values, zeros, or a distinct-based average.

Concretely, detect duplicates with COUNT(*) versus COUNT(DISTINCT grain), fix with MERGE, and do not "correct" AVG by dividing by two.

The reason for that specificity is a failure I have seen: An analyst divided AVG by 2 to "undo duplicates." The true average was already correct. Reported AOV was $27 instead of $54. Pricing changed for 4 days.

Every row copied once.

AggregateAfter exact 2×Naive "fix"
SUM×2—
COUNT×2—
AVGunchanged $54wrongly $27

I would not consider it settled without evidence: duplicate every row exactly and require AVG unchanged and SUM doubled.

Exact copies cancel in the average.

Curated: · Written: · Reviewed:

QA-62incremental unique_key is order_id. The source emits two lines per order. What does the incremental merge do?(show answer)

I would start dbt incremental unique_key from the destination checksum and the failed partition, not from the DAG run.

unique_key is the merge grain dbt uses for incremental MERGE. If it is not unique in the source, the merge is undefined: engines may error, keep an arbitrary line, or insert both. It does not reliably keep one line and delete the other. The model still lies about grain.

Concretely, set unique_key to a key that is actually unique, or pre-aggregate to that grain on purpose, and test two lines surviving when line grain is intended.

The reason for that specificity is a failure I have seen: unique_key=order_id met two source lines. Warehouse A errored the incremental. Warehouse B inserted both lines. Neither reliably kept a single line. The warehouse that silently picked one line dropped 1.9 million others and item count was 48 percent low, while order_id unique passed.

Two lines, unique_key order_id.

unique_keyLines keptItem count
order_id (not unique)error, both, or arbitrary 148 percent low or fail
(order_id, line_id)2correct

I would not consider it settled without evidence: load a 2-line order and require 2 rows when grain is line, or 1 aggregated row when grain is order.

unique_key is the grain you keep.

Curated: · Written: · Reviewed:

QA-63A dbt snapshot uses check strategy on 3 columns. The source also changes a 4th. What history do you have?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

check strategy versions when listed columns change. Unlisted columns can change silently with no new snapshot row. timestamp strategy versions on an updated_at you trust. The missing column is the missing history.

Concretely, list every attribute that must freeze, or use a trusted timestamp, and test a change to an unlisted column.

The reason for that specificity is a failure I have seen: status was unlisted. 210,000 orders flipped status with no snapshot row. Downstream SCD stayed "pending" for 9 days after ship.

status not in check list.

ColumnIn check listSnapshot rows on change
amountyesnew version
statusno0, 210,000 silent

I would not consider it settled without evidence: change an unlisted column in a fixture and require either a new snapshot row or an explicit decision that the column is type-1.

Unchecked columns are type-1 by accident.

Curated: · Written: · Reviewed:

QA-64Great Expectations runs after Spark writes the partition. The suite fails. What is already true for readers?(show answer)

My answer to quality checks before write begins with where extract, load, and transform actually ran, and on which clock.

A failing suite after write means the bad partition is already visible unless you wrote to staging. Quality after publish is a ticket. Quality before pointer swap is a gate.

Concretely, write staging, run the suite, swap to serving on pass, and leave serving on the previous partition on fail.

The reason for that specificity is a failure I have seen: Suite failed on null keys 18 minutes after write. Three streaming consumers had already read 2.2 million bad rows into caches that lived 1 hour.

GE after serving write.

OrderBad rows consumedServing on fail
write serving, then GE2,200,000the bad partition
GE on staging, then swap0previous

I would not consider it settled without evidence: fail a suite in staging and require serving still on the previous snapshot.

Validate, then publish.

Curated: · Written: · Reviewed:

QA-65Landing allows extra JSON fields. Curated is typed. Who is allowed to SELECT landing?(show answer)

I would treat landing versus curated contracts as a claim about every partition and every late window, not about the one that was demonstrated.

Landing is for replay and for the pipeline. Curated is for consumers. If analysts query landing, they bypass every contract you put on curated and freeze accidental fields into dashboards.

Concretely, revoke analyst SELECT on landing, grant pipeline roles only, and document replay as a pipeline operation.

The reason for that specificity is a failure I have seen: A dashboard pointed at landing because curated lagged 2 hours. It parsed a field that later changed type. 6 weeks of "revenue" was a string concat, then a fail, $0 rendered.

Dashboard on landing.

SourceContractOutcome
landingnonestring concat then $0
curatedtypedlagged but correct

I would not consider it settled without evidence: attempt analyst SELECT on landing and require deny.

Landing is not a faster curated.

Curated: · Written: · Reviewed:

QA-66The load lists s3://bucket/dt=2026-03-02/*.parquet. A late file arrives after the list. What did you miss?(show answer)

The useful question for manifest versus listing is what a second consumer would read from the same sink after a retry.

A list at T0 is a snapshot of names. Files that land at T0+n are invisible until the next list. Manifests from the producer are the membership contract; listings race writers.

Concretely, read a producer manifest with checksums, or use a table format commit, and do not glob a prefix that is still receiving PUTS.

The reason for that specificity is a failure I have seen: List at 02:00 missed 14 files that finished at 02:04. 220,000 rows never loaded. The DAG was green. The files sat in the prefix looking "already processed."

Late PUT after list.

MembershipFiles missedRows
glob at 02:0014220,000
closed manifest00 missing

I would not consider it settled without evidence: PUT a file after list start and require the job to wait for a closed manifest, not a glob.

A glob is a race.

Curated: · Written: · Reviewed:

QA-67You paginate ListObjects. During the pagination the producer overwrites a key. What set did you load?(show answer)

I would settle object-store listing races by replaying the failed path against the contract, not against a green Airflow square.

List pagination is not a transaction. Overwrites and deletes during the walk can yield a mix of old and new, or skip keys. A table commit or a versioned manifest is the consistent set.

Concretely, do not glob a live prefix; load from a committed snapshot or a complete manifest written after the producer finishes.

The reason for that specificity is a failure I have seen: Mid-list overwrite mixed 8 old files with 11 new. Checksums failed 3 days later. Mixed GMV was $2.1 million off.

Overwrite during ListObjects.

MethodFile setGMV error
paginated listmixed 8 old + 11 new$2.1 million
post-complete manifestone generation$0

I would not consider it settled without evidence: overwrite a key during a mock list and require the production job to use a manifest instead.

List is not isolation.

Curated: · Written: · Reviewed:

QA-68Consumer lag is 0. The warehouse hour is still 4 percent short. What was lag not measuring?(show answer)

The judgement in Kafka lag versus completeness is which operator wrote the files and whether that operator is idempotent.

Lag is how far the consumer offset is behind the log end. It is not a count of decoded, merged, passing-contract rows. Decode drops, DLQ, and failed MERGE can leave lag 0 and a short table.

Concretely, compare merged row counts to produced counts (or end-to-end checksums), and treat lag as a delay metric only.

The reason for that specificity is a failure I have seen: Lag was 0. A schema error sent 4 percent to nowhere (no DLQ). The hour was short $900,000. The lag tile was green.

Zero lag, short hour.

SignalValueTable
Kafka lag0—
merged vs produced96 percent$900,000 short

I would not consider it settled without evidence: drop 4 percent in decode with lag 0 and require a completeness page.

Caught up is not complete.

Curated: · Written: · Reviewed:

QA-69The job uses processing-time windows to "meet the SLA." Event time is 40 minutes behind. What number did you publish?(show answer)

Where candidates lose the interview on watermark versus wall-clock SLA is calling the DAG success the load.

Processing-time windows close on the processing clock. There is no processing-time watermark: watermarks bound event-time lateness. If event time lags, a processing-time hour can close while those events have not arrived, then they land in later processing-time buckets.

Concretely, window and watermark on event time, advertise SLA on event-time watermark close, and do not close event-time hours on the processing clock to look fast.

The reason for that specificity is a failure I have seen: A processing-time window published 09:00 at 09:05 missing 40 minutes of events, 1.8 million rows. They appeared in 09:40 processing buckets. Two hours both wrong.

40-minute event-time lag.

Clock09:00 published atMissing rows
processing-time window09:051,800,000
event time minus latenesswhen events arrive0 in the wrong hour

I would not consider it settled without evidence: lag event time 40 minutes and require the event-time 09:00 window to stay open.

Fast on the wrong clock is a wrong hour.

Curated: · Written: · Reviewed:

QA-70The source deletes a customer. Debezium emits a tombstone. Your MERGE only handles upserts. What remains?(show answer)

I would answer CDC tombstones by separating delivery, transform location, and the row set the sink is allowed to keep.

A tombstone is a delete of that primary key, not a null upsert. Ignoring it leaves the curated row forever. Applying it as an upsert with null columns may create a zombie that still joins. The MERGE needs a delete branch when the op is d or the record is a Kafka tombstone.

Concretely, when op=d or the value is null then DELETE the key, when op=c/u then upsert, and test that a tombstone leaves 0 rows for that key in curated.

The reason for that specificity is a failure I have seen: 12,400 deleted customers stayed in curated. A GDPR extract still listed them 19 days later.

Tombstones ignored.

MERGEDeleted customers remainingGDPR days
upsert only12,40019
delete on tombstone00 extra

I would not consider it settled without evidence: emit a tombstone and require 0 rows for that key in curated.

No delete branch means no deletes.

Curated: · Written: · Reviewed:

QA-71The source sets deleted_at. Your incremental copies WHERE updated_at > :cursor and never looks at deleted_at. What do you keep?(show answer)

The engineering content of soft-delete incremental is the merge key and the quality gate that would refuse the write, not the scheduler.

A soft delete is an update. If you copy the row as a living fact because it still has a recent updated_at, you keep a tombstone as a customer. The job must apply deleted_at as a delete or an is_current=false.

Concretely, treat deleted_at IS NOT NULL as a delete in MERGE, and test a newly stamped deleted_at.

The reason for that specificity is a failure I have seen: 8,900 soft-deleted users remained current. Email sends hit 8,900 retired addresses for 3 weeks.

Soft deletes copied as live.

HandlingRetired users still currentWeeks of mail
copy on updated_at8,9003
MERGE delete on deleted_at00

I would not consider it settled without evidence: stamp deleted_at and require is_current=false or row gone.

deleted_at is an operator, not a comment.

Curated: · Written: · Reviewed:

QA-72The OLTP primary's clock jumped back 4 minutes. Your extractor uses updated_at > last_run. What hole did you punch?(show answer)

Before calling clock skew on updated_at done I would write down the partition, consumer, or late window nobody checked.

A backward clock makes new commits look older than the cursor. Those rows are skipped until a reconcile. NTP and LSN cursors are the engineering response, not a hope that clocks are monotonic.

Concretely, prefer LSN, overlap the extract by more than observed skew, and alert on source clock steps.

The reason for that specificity is a failure I have seen: Clock stepped back 4 minutes. 6,200 orders had timestamps behind last_run. They never extracted. Hole: $180,000 until the Sunday dump.

OLTP clock stepped back 4 minutes.

CursorOrders missedDollars
updated_at > last_run6,200$180,000
LSN or 10-minute overlap0$0

I would not consider it settled without evidence: simulate a -4 minute step and require overlap or LSN to still capture the orders.

Clocks step; LSNs do not.

Curated: · Written: · Reviewed:

QA-73Ops dropped a full dump in the CDC folder. The job applied it as inserts. What should have happened?(show answer)

The first thing I would establish about full dump applied as CDC is which write actually landed, not which task box turned green.

Folder is not mode. A dump in a CDC path is still a dump. The loader has to read a header or filename contract and choose replace versus apply, or refuse the file.

Concretely, require a mode header, refuse dumps in the CDC path, and alert on row-count jumps over 2×.

The reason for that specificity is a failure I have seen: Dump-as-insert doubled 11 million rows. Distinct keys stayed 11 million. Unique tests sampled 800 rows and passed.

Dump in the CDC folder.

OperatorRowsDistinct keysSample unique
insert22,000,00011,000,000800 rows, pass
refuse or replace11,000,00011,000,000full

I would not consider it settled without evidence: place a dump in the CDC path in staging and require a refuse or a snapshot replace, not insert.

The path is not the operator.

Curated: · Written: · Reviewed:

QA-74Streaming MERGE and a batch overwrite both own today. Who wins?(show answer)

I would start dual writers on the open partition from the destination checksum and the failed partition, not from the DAG run.

Two writers without a lock produce lost updates or duplicates. The open partition has one writer, or they are coordinated transactions. "They usually don't overlap" is not a lock.

Concretely, lock today's partition to the stream, run batch on closed days, or use a single Iceberg writer queue.

The reason for that specificity is a failure I have seen: Overlap window 14 minutes. Stream inserted 90,000; overwrite replaced the day from a 10-minute-old dump and dropped those 90,000.

Stream and overwrite on today.

Control14-minute overlapStream rows kept
noneboth ran0 of 90,000
lockserial90,000

I would not consider it settled without evidence: run both in staging on one partition and require a lock that serialises them.

Today has one owner.

Curated: · Written: · Reviewed:

QA-75The job writes to a path and then renames _temporary to the partition. A reader globbed mid-rename. What did they see?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Rename-based commits are not atomic for readers who glob. They can see a mix of old files and partial new files. A table format snapshot or a copy-on-complete _SUCCESS after an atomic replace is the reader contract.

Concretely, serve through Iceberg/Delta/Hive commit, or stage then atomic replace the partition prefix in a way readers do not glob mid-flight.

The reason for that specificity is a failure I have seen: A glob mid-rename mixed 40 old files and 12 new. Query failed schema merge, then succeeded on a mix. One dashboard showed $0 for 8 minutes then a spike.

Glob during _temporary rename.

ReadFiles seenDashboard
glob mid-rename40 old + 12 new$0 then spike, 8 minutes
Iceberg snapshotprevious or next wholestable

I would not consider it settled without evidence: read during a rename commit and require the production path to use a snapshot isolation table.

Rename is not a snapshot.

Curated: · Written: · Reviewed:

QA-76A 40-minute dashboard query started on snapshot 88. A compaction committed 120 times during it. Which files should the query use?(show answer)

My answer to Iceberg snapshot isolation for readers begins with where extract, load, and transform actually ran, and on which clock.

A scan pins a snapshot. Compaction creating new snapshots must not make the query mix file generations. If the engine re-plans mid-query against latest, isolation is gone.

Concretely, pin snapshot id at query start, and verify long scans do not pick up mid-flight files.

The reason for that specificity is a failure I have seen: A query that refreshed metadata every 5 minutes mixed compacted and uncompacted files, double-counted 6.2 million rows, SUM ×1.4.

Metadata refresh during compaction.

IsolationRowsSUM
refresh latest every 5 minutes+6,200,000 mixed×1.4
pin snapshot 88stablestable

I would not consider it settled without evidence: compact during a long SELECT and require counts to match the start snapshot.

Pin the snapshot you started.

Curated: · Written: · Reviewed:

QA-77A bad load committed snapshot 200. You roll back serving to 199. What must the next incremental not do?(show answer)

I would treat time travel for job rollback as a claim about every partition and every late window, not about the one that was demonstrated.

Rolling back the pointer does not un-consume Kafka or un-extract CDC. The next incremental must not apply offsets that were already in 200 as if they were new, and must not skip them if 200 is discarded. Offsets and snapshots have to rewind together.

Concretely, rewind the streaming checkpoint (or CDC cursor) to the snapshot 199 corresponding offsets, then serve 199.

The reason for that specificity is a failure I have seen: Serving rolled to 199. Kafka stayed at the 200 offsets. Those 1.4 million events never reapplied. Hole stayed until a manual replay.

Snapshot 199, Kafka still at 200.

CursorEvents in 200 not in 199After rollback
Kafka stays1,400,000missing
rewind with snapshot1,400,000reapplied

I would not consider it settled without evidence: rollback snapshot and require the cursor to match 199's offsets before enabling the stream.

Table rollback without cursor rollback is a hole.

Curated: · Written: · Reviewed:

QA-78Legal wants the email gone from Parquet that snapshot 12 still references. Is a SQL DELETE enough if you also expire snapshot 12's metadata?(show answer)

The useful question for GDPR file rewrite not DELETE is what a second consumer would read from the same sink after a retry.

Expiring snapshot metadata while leaving data files is not erasure. SQL DELETE may only add a delete file (MOR) or rewrite some files (COW) and still leave copies in older files until those files are gone. Erasure is rewrite, expire, remove.

Concretely, rewrite files without the email, expire every snapshot that pointed at old files, remove orphans, then decode remaining Parquet and verify the identifier is gone.

The reason for that specificity is a failure I have seen: DELETE plus metadata expiry left 3 data files on disk with the email for 40 days. strings on the prefix found 0 hits because Parquet encoding hid the address. A forgotten EMR job that decoded the files still read the email.

DELETE plus expire metadata only.

ActionEmail in remaining filesDays
DELETE + expire metadata3 files40
rewrite, expire, remove00

I would not consider it settled without evidence: decode remaining Parquet after the job and require 0 rows with the email. strings on the prefix cannot prove erasure.

Metadata expiry is not a shredder.

Curated: · Written: · Reviewed:

QA-79You hash emails in curated but landing keeps plaintext so "we can debug." Who is the processor of record?(show answer)

I would settle tokenise PII at ingest by replaying the failed path against the contract, not against a green Airflow square.

If landing has plaintext, you are processing PII there. Debugging is an access path. Tokenise or encrypt at ingest with a keyed vault; debug via a break-glass that is audited.

Concretely, tokenise on the first write, store tokens in landing, and require break-glass for reverse lookup.

The reason for that specificity is a failure I have seen: Landing plaintext was queried 140 times by "debug." 140 were unaudited. A contractor dumped 2.4 million emails.

Plaintext landing for debug.

LandingUnaudited debug queriesEmails dumped
plaintext1402,400,000
tokens plus break-glassaudited0 contractor dump

I would not consider it settled without evidence: query landing as an analyst and require tokens only.

Debug is not a second warehouse of PII.

Curated: · Written: · Reviewed:

QA-80You fixed the schema and replayed the DLQ with the same event ids into an append sink. What duplicated?(show answer)

The judgement in DLQ replay after schema fix is which operator wrote the files and whether that operator is idempotent.

DLQ replay is at-least-once. Without MERGE on event id, successful late parses append beside the original if the original also landed, or duplicate if the batch was retried. Replay needs the same idempotent sink as the live path.

Concretely, replay DLQ through MERGE on event_id, and drop DLQ records already in the sink.

The reason for that specificity is a failure I have seen: Replay appended 180,000 events that had also succeeded on a retry. COUNT(*) +180,000. COUNT(DISTINCT event_id) unchanged.

DLQ replay into append.

SinkExtra rowsDistinct event_id
append180,000unchanged
MERGE event_id0unchanged

I would not consider it settled without evidence: replay DLQ after a partial success and require 1 row per event_id.

Replay is a second write; merge it.

Curated: · Written: · Reviewed:

QA-8120 warehouse loads share a pool of 4. You set the DAG concurrency to 16 anyway. What actually runs?(show answer)

Where candidates lose the interview on Airflow pool concurrency is calling the DAG success the load.

The pool is the real parallelism against the warehouse. DAG concurrency above the pool just queues. Ignoring the pool and raising worker slots instead stampedes the warehouse.

Concretely, size the pool to warehouse slots, set DAG concurrency ≤ pool for that queue, and watch queued tasks, not running DAGs.

The reason for that specificity is a failure I have seen: Someone raised worker slots. 16 Spark apps submitted. The warehouse hit 100 percent, all 16 took 4× longer, and the SLA pack missed by 2 hours.

16 apps versus pool of 4.

ThrottleRunning Spark appsPack
worker slots 1616miss 2 hours
pool 44on time

I would not consider it settled without evidence: submit 16 and require the pool to admit 4, with 12 queued.

The pool is the throttle; worker slots are not.

Curated: · Written: · Reviewed:

QA-82You added backoff. The sink is still append. A timeout that actually succeeded then retries. What happens?(show answer)

I would answer exponential backoff versus idempotency by separating delivery, transform location, and the row set the sink is allowed to keep.

Backoff reduces storms. It does not make a timeout safe. A request that succeeded on the server and timed out on the client will retry and duplicate unless the sink is idempotent.

Concretely, keep backoff, and MERGE on a request id the source echoes.

The reason for that specificity is a failure I have seen: Timeouts retried 2,400 "failed" loads that had written. Extra rows 2,400 × 50,000. AVG unchanged, SUM not.

Successful write, client timeout, retry append.

SinkExtra batchesSUMAVG
append2,400upunchanged
MERGE request_id0heldheld

I would not consider it settled without evidence: succeed on the server, timeout the client, retry, and require 1 row per request id.

Backoff without a key still doubles.

Curated: · Written: · Reviewed:

QA-83The API returns pages of 500. You stop when a page is short. The last page was exactly 500. What did you miss?(show answer)

The engineering content of REST pagination extract is the merge key and the quality gate that would refuse the write, not the scheduler.

A full last page is not EOF by itself. After every full page, request the next page; an empty page or a null next_cursor is the stop. Stopping because the current page is short also ends a short last page. Do not stop on an unrelated page-count cap such as 40.

Concretely, follow next_cursor until null or until an empty page after a full page, and test a fixture whose total is a multiple of page size.

The reason for that specificity is a failure I have seen: Exactly 20,000 rows, 40 pages of 500. The job treated 40 pages as complete and never requested page 41. Page 41 would have been empty, which is how a full last page is confirmed. When the set grew to 20,010, the same 40-page cap missed 10 rows.

20,000 rows, page size 500.

Stop rule20,000 loaded20,010 loaded
stop at 40 pages20,00020,000, miss 10
next page until empty/null20,000 (page 41 empty)20,010

I would not consider it settled without evidence: fixture 20,000 and 20,010 and require both complete via cursor.

Page size multiples are not EOF.

Curated: · Written: · Reviewed:

QA-84Row groups are 2 MB because streaming flushed often. Predicates still read almost every group. Why?(show answer)

Before calling Parquet row-group size done I would write down the partition, consumer, or late window nobody checked.

Predicate pushdown skips row groups using min/max stats. Tiny groups do not inherently weaken that skipping: min/max can still exclude a 2 MB group. The downside is metadata and open overhead—thousands of footers and file handles. Huge groups mean you read 1 GB to get 10 rows. Target hundreds of MB for batch tables.

Concretely, compact to ~128 MB row groups, sort on the filter key when possible, and measure groups opened and footer time.

The reason for that specificity is a failure I have seen: 2 MB groups, 9,400 opens for a selective query, 11 minutes of footer and open overhead. After compact to 128 MB sorted on customer_id, 14 groups, 40 seconds. Stats skipping still worked on the tiny groups; opens dominated.

Selective query, two row-group sizes.

Row groupGroups openedWall time
2 MB9,40011 minutes
128 MB sorted1440 seconds

I would not consider it settled without evidence: EXPLAIN or metrics for row groups read before and after compact.

Row-group size is a read SLO.

Curated: · Written: · Reviewed:

QA-85You wrap the partition column in CAST. The scan reads 18 TB. What did CAST do?(show answer)

The first thing I would establish about predicate pushdown is which write actually landed, not which task box turned green.

Functions on the column typically disable partition pruning and parquet min/max pushdown because the engine can no longer match the predicate to directory names or row-group stats. Filter on the stored type and cast the literal instead.

Concretely, write dt = DATE '2026-03-02' not CAST(dt AS string) = '2026-03-02' when dt is a date, and confirm partitionFilters in EXPLAIN before merging the job.

The reason for that specificity is a failure I have seen: CAST(dt AS string) scanned 18 TB, 69 minutes. dt = DATE '2026-03-02' scanned 48 GB, 4 minutes.

CAST on dt.

PredicateBytes scannedWall time
CAST(dt AS string) = '2026-03-02'18 TB69 minutes
dt = DATE '2026-03-02'48 GB4 minutes

I would not consider it settled without evidence: EXPLAIN both predicates and require partitionFilters only on the native form.

Cast the literal; leave the column.

Curated: · Written: · Reviewed:

QA-86You add hour to an Iceberg spec. Old files stay day-partitioned. A job overwrites hour=0 only. What did it replace?(show answer)

I would start partition evolution in the job from the destination checksum and the failed partition, not from the DAG run.

After partition evolution, old files keep the old spec. A one-hour overwrite does not universally delete the whole old day: Iceberg overwrite-by-filter can replace only matching files, while a Hive-style day overwrite can remove mixed-spec files for that day. The job must use the table's overwrite semantics for mixed specs, not assume either outcome.

Concretely, append or MERGE across mixed specs, overwrite only with a predicate the current spec actually matches, and test a point query that spans old and new files.

The reason for that specificity is a failure I have seen: hour=0 overwrite was expected to delete the whole old day. It did not: mixed-spec day files remained, and the hour duplicated 17 million rows against them. In another catalog mode a day-scoped overwrite would have removed 23 hours. Mixed Iceberg specs do not have one overwrite outcome.

hour=0 overwrite after day→hour evolution.

WriteOther hoursUniversal?
1-hour overwritekept or deleted by spec/modeno whole-day delete rule
MERGE that hourkept0 unexpected deletes

I would not consider it settled without evidence: overwrite one hour on a mixed-spec table and require other hours to remain; do not assume a whole-day delete.

New spec plus partial overwrite is how days vanish.

Curated: · Written: · Reviewed:

QA-87You Z-order customer_id weekly. Queries still open most files. What did the weekly job miss?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Clustering only applies to files the job rewrites. New streaming files after the job are unclustered. If most bytes are those new files, weekly Z-order of history does not help today's point lookups.

Concretely, cluster on a cadence that covers the files queries actually read, or cluster on write, and measure files touched for a 1-customer query on the latest day.

The reason for that specificity is a failure I have seen: Weekly Z-order rewrote last week. Today had 8,200 unclustered files. A 1-customer query opened 8,200 files, 16 minutes. After clustering today's files, 6 files, 12 seconds.

Weekly Z-order versus today's files.

Files clustered1-customer query filesTime
last week only8,200 today16 minutes
including today612 seconds

I would not consider it settled without evidence: run a 1-customer query on dt=today after the weekly job and require files opened under a budget.

Cluster the files you query, not last week's.

Curated: · Written: · Reviewed:

QA-88A Flink checkpoint succeeds. Kafka offsets in the checkpoint are behind the sink's last write. After restore, what happens?(show answer)

My answer to Flink checkpoint versus Kafka begins with where extract, load, and transform actually ran, and on which clock.

Restore replays from checkpointed offsets. If the sink already wrote later records without those offsets in the same checkpoint, you get duplicates unless the sink merges. Checkpoint and sink commit must align.

Concretely, use two-phase sink commits with checkpoints, or MERGE on event id after restore.

The reason for that specificity is a failure I have seen: Restore replayed 12 minutes, 540,000 events already in the sink. Append duplicated them. COUNT(DISTINCT event_id) held; COUNT(*) did not.

Restore 12 minutes behind the sink.

SinkExtra rowsDistinct event_id
append540,000unchanged
MERGE0unchanged

I would not consider it settled without evidence: restore from checkpoint and require merge semantics or aligned commits.

A checkpoint is a rewind; the sink must tolerate it.

Curated: · Written: · Reviewed:

QA-89trigger(processingTime='1 second') on a sink that cannot commit that often. What fails first?(show answer)

I would treat structured streaming trigger as a claim about every partition and every late window, not about the one that was demonstrated.

A Structured Streaming query runs one micro-batch at a time. It does not overlap commits. A trigger faster than commit time does not start a second batch; the next trigger waits, so lag grows and you still write tiny files. The trigger is a commit budget: it should sit at or above p99 sink commit time, not at a round number that looks real-time on a slide.

Concretely, set trigger ≥ p99 commit time, compact on the same DAG, and measure batch duration and lag, not concurrent commits.

The reason for that specificity is a failure I have seen: 1-second trigger, 8-second commits. Batches ran serially, lag grew by about 7 seconds plus backlog, and 50,000 tiny files/day accumulated. There were not 900 overlapping Iceberg commits per hour. After a 1-minute trigger, lag cleared, files 1,400/day.

1-second trigger, 8-second commits.

TriggerConcurrent batchesLagFiles/day
1 secondnone; serial~7s plus backlog50,000
1 minutenone; serialcaught up1,400

I would not consider it settled without evidence: compare lag and file count at 1s versus 1 minute on the same sink.

The trigger is a commit budget.

Curated: · Written: · Reviewed:

QA-90You switched to continuous trigger for "real time." The sink is Iceberg MERGE. Why did latency get worse?(show answer)

The useful question for micro-batch versus continuous is what a second consumer would read from the same sink after a retry.

Spark continuous processing generally does not support an Iceberg MERGE sink. Continuous is for low-latency operators that can emit per record. MERGE needs a micro-batch commit. Forcing continuous plus MERGE is a mismatch: use micro-batch sized to the sink.

Concretely, use micro-batch sized to the Iceberg MERGE, and do not pick continuous for a table-format merge sink.

The reason for that specificity is a failure I have seen: Continuous trigger was set on a foreachBatch Iceberg MERGE. Spark continuous processing does not run that sink: the query rejected the combination and never committed MERGE snapshots. A 30-second micro-batch MERGE reached 40-second p99 visibility.

Iceberg MERGE under two triggers.

TriggerIceberg MERGEp99 visibility
continuousnot a supported sinkno MERGE commits
30-second micro-batch2/min40 seconds

I would not consider it settled without evidence: attempt continuous plus Iceberg MERGE and require the job to use micro-batch; measure p99 sink visibility under micro-batch.

Continuous is not faster MERGE.

Curated: · Written: · Reviewed:

QA-91A nested array has p99 length 12 and a tail of 40,000. You explode then join. What dies?(show answer)

I would settle explode cardinality explosion by replaying the failed path against the contract, not against a green Airflow square.

Explode multiplies rows by array length. A tail of 40,000 turns one parent into 40,000 join keys and can dominate shuffle. Cap, split, or pre-aggregate the tail.

Concretely, cap explode length with a documented remainder, page on tail arrays, and never join the uncapped explode of a user-controlled array.

The reason for that specificity is a failure I have seen: One record with 40,000 items exploded to 40,000 rows, skewed a join, and held 80 executors for 5 hours. The other 99.99 percent of records were fine.

One 40,000-item array.

ExplodeExtra rowsJob
uncapped40,000 from 1 parent5 hours skew
cap 100 + remainder path100 + 1 remainderon SLA

I would not consider it settled without evidence: inject a 40,000-item array and require a cap or a split path.

p99 is not the tail that kills you.

Curated: · Written: · Reviewed:

QA-92You filter amount > 0 in WHERE and then ROW_NUMBER to pick latest. Why do you still keep a $0 row as latest?(show answer)

The judgement in window filter versus WHERE is which operator wrote the files and whether that operator is idempotent.

WHERE runs before window functions. A later $0 row can still be the latest if you only filtered in a subquery incorrectly, or a $0 can be dropped before ranking so an older positive looks current. The filter has to apply at the grain you intend.

Concretely, rank on the full stream, then filter, or filter first if $0 must never be current—and test both a later $0 and an earlier $0.

The reason for that specificity is a failure I have seen: WHERE amount > 0 dropped the latest correction of $0 (a void). Rank kept a $40 row. Refunds were understated $620,000 for a week.

Later $0 void.

Order of operationsCurrent amountRefunds
WHERE amount>0 then rank$40$620,000 under
rank then keep including $0$0correct

I would not consider it settled without evidence: send a later $0 void and require it to be current if voids are real.

Filter order is the grain of "current."

Curated: · Written: · Reviewed:

QA-93A fact arrives 3 minutes before the dimension row. The stream inner-joins. Where did the fact go?(show answer)

Where candidates lose the interview on late dimension lookup in a stream is calling the DAG success the load.

An inner join against a dimension that has not arrived yet drops the fact. Streaming lookups need a hold buffer, an inferred member, or an outer join plus a repair MERGE. Dropping with 0 errors is a silent completeness hole the DAG will never page.

Concretely, outer-join or retry the lookup for allowed lateness, mint an inferred member if the fact cannot wait, then repair attributes when the dimension lands.

The reason for that specificity is a failure I have seen: 3-minute dim lag dropped 27,000 first-time orders, $2.1 million, which never retried. The stream task had 0 errors.

Fact before dimension.

JoinFirst-time ordersDollars
inner, no retry0 of 27,000$2.1 million gone
outer plus repair27,000$0 missing

I would not consider it settled without evidence: send a fact 3 minutes early and require it to land after the dim arrives.

Inner join is a drop rule.

Curated: · Written: · Reviewed:

QA-94A broadcast join hangs until spark.sql.broadcastTimeout. The dimension is 900 MB. What is actually timing out?(show answer)

I would answer broadcast timeout by separating delivery, transform location, and the row set the sink is allowed to keep.

Broadcast timeout is how long executors wait to receive the relation. A 900 MB dim on a slow shuffle fetch can miss a 300s timeout even if it would fit. Raising timeout without checking size just waits longer for the same hang.

Concretely, confirm dim size, increase timeout only with evidence, or skip broadcast.

The reason for that specificity is a failure I have seen: Timeout 300s, dim 900 MB, fetch 8 MB/s on a bad node. Four retries are 300s × 4 = 20 minutes; the original attempt plus 4 retries is 25 minutes of SLA miss, not 40. Shuffle join would have finished in 11 minutes.

900 MB broadcast, slow fetch.

StrategyTimeoutFinish
broadcast300s then 4 retries25 minutes miss (5 × 300s)
shuffle joinn/a11 minutes

I would not consider it settled without evidence: measure dim bytes and fetch rate; choose shuffle if timeout would exceed SLA.

Timeout is not a join strategy.

Curated: · Written: · Reviewed:

QA-95Broadcast saved 20 minutes. The dim is 3 GB on 400 executors. What did memory pay?(show answer)

The engineering content of shuffle versus broadcast cost is the merge key and the quality gate that would refuse the write, not the scheduler.

Broadcast copies the dimension to every executor: 3 GB × 400 = 1.2 TB of copies in RAM. Shuffle moves each byte about once plus overhead. Twenty minutes of wall time can still be more expensive than 1.2 TB of executor memory if the cluster OOMs.

Concretely, compute copies × size next to shuffle bytes, and pick the join using both wall time and peak RAM, not wall time alone.

The reason for that specificity is a failure I have seen: 400 executors × 3 GB caused 60 executors to spill and 12 to OOM. The "faster" broadcast never finished. Shuffle finished in 28 minutes.

3 GB dim, 400 executors.

JoinRAM copiesFinished
broadcast1.2 TBOOM 12 executors
shuffle~once28 minutes

I would not consider it settled without evidence: write 3 GB × executor count next to the 20-minute claim.

Minutes are not the only currency.

Curated: · Written: · Reviewed:

QA-96The Spark job SLA is 60 minutes. The data SLA is 09:00 numbers by 09:15. The job started at 08:50 and finished at 09:40, inside 60 minutes. Did you meet the data SLA?(show answer)

Before calling job SLA versus data SLA done I would write down the partition, consumer, or late window nobody checked.

A job duration SLO is not a freshness SLO. Starting late burns the data SLA even when duration is healthy. The number finance bought is when the serving pointer moved, not when Spark's last stage finished inside 60 minutes.

Concretely, alert on serving-pointer time versus the hour, keep duration as a cost and capacity signal, and do not page only on job runtime.

The reason for that specificity is a failure I have seen: Duration 50 minutes, start 08:50, serve 09:40. Data SLA 09:15 missed. Job SLA green. This happened 11 weekdays in a row.

50-minute job, late start.

SLATargetActualColor
job duration60 minutes50green
data serve09:1509:40red, 11 days

I would not consider it settled without evidence: chart start, duration, and serve; require serve ≤ data SLA.

Fast enough after too late is still late.

Curated: · Written: · Reviewed:

QA-97A stage retries after a fetch failure. The first attempt already wrote part files. The sink is a plain path. What is in the partition?(show answer)

The first thing I would establish about Spark stage retry versus duplicate files is which write actually landed, not which task box turned green.

Task retries can write a second copy of the same spill/output if the commit protocol does not discard the first attempt. Path sinks without Spark/Hive/Iceberg commit can keep both. Table formats commit one snapshot.

Concretely, write through a transactional table, or use Spark's unique attempt paths plus a commit that ignores failed attempts.

The reason for that specificity is a failure I have seen: Retry doubled 2.8 million rows in a globbed path. SUM ×2, AVG unchanged.

Stage retry into a glob path.

SinkRowsSUMAVG
glob path×2 (2.8 million extra)×2unchanged
Iceberg commit×1×1unchanged

I would not consider it settled without evidence: fail a task after write and retry into a path sink versus Iceberg, and require Iceberg counts unchanged.

Retries are second attempts; the sink must keep one.

Curated: · Written: · Reviewed:

QA-98remove_orphan_files is scheduled 30 days out. Jobs fail hourly and leave files. What is the storage bill, and what is the GDPR bill?(show answer)

I would start orphan files after failed commit from the destination checksum and the failed partition, not from the DAG run.

Orphans cost storage immediately. If they contain PII, they are also copies the erasure job will miss until they are deleted. A 30-day delay is a 30-day privacy clock.

Concretely, remove orphans on a delay just longer than the longest retry, and include orphans in erasure drills.

The reason for that specificity is a failure I have seen: Hourly failures left 14 TB/month of orphans. Erasure of 1 subject missed 40 orphan files for 30 days. Legal counted 40 extra copies.

Hourly failed writes, 30-day orphan TTL.

TTLExtra storageErasure misses
30 days14 TB/month40 files, 30 days
6 hours (> retry)boundedinside SLA

I would not consider it settled without evidence: fail a job, run erasure, and require orphan removal inside the erasure SLA.

Orphans are tables without a catalog.

Curated: · Written: · Reviewed:

QA-99You expire snapshots older than 7 days to save metadata. A bad load was 8 days ago. What did you give up, and what did you not erase?(show answer)

This is an area where an orchestrator success and a correct table are different observations.

Snapshot expiry drops time travel for those ids. It does not by itself delete all data files still referenced by remaining snapshots, and it is not GDPR erasure. It is an operational retention choice that kills rollback past the horizon.

Concretely, set expiry to the rollback you still owe, and run separate erasure rewrites for subjects.

The reason for that specificity is a failure I have seen: Expiry at 7 days. Bad load at day 8 was un-rollbackable. The files were still live in snapshot 9's lineage. On-call could not undo. Subjects in those files were also not erased by expiry.

7-day expiry, 8-day-old bad load.

NeedAfter expiry
rollback to day 8impossible
GDPR via expiry onlyidentifiers may remain in live files

I would not consider it settled without evidence: attempt rollback past expiry and require a documented "cannot"; attempt erasure via expiry only and require it to fail the privacy test.

Expiry is a rollback horizon, not a shredder.

Curated: · Written: · Reviewed:

QA-100You validated dt=2026-03-02 and shipped. dt=2026-03-01 still has the old schema. Why is that not done?(show answer)

My answer to proving one partition begins with where extract, load, and transform actually ran, and on which clock.

A pipeline change is done when every partition and consumer the contract names has been proven. One good day is a sample. Streaming late windows and the other consumers are separate writers.

Concretely, list partitions and consumers in the change, test a late window and a second sink, and refuse "02 looked fine."

The reason for that specificity is a failure I have seen: 03-02 passed. 03-01 still had null new columns, 6.4 million rows. A late 03-02 window wrote the old schema into a second sink. Two consumers disagreed for 5 days.

Validated one day.

PathRows on old schemaDays disagreed
dt=2026-03-02 batch0—
dt=2026-03-016,400,0005
late window, second sinkold schema5

I would not consider it settled without evidence: run the new job on 03-01, a late 03-02 window, and the second sink before calling the change done.

One green partition is a sample.

Curated: · Written: · Reviewed: