Data Engineering Reference/Production Pipeline Operations

Delivery Semantics & Failure Recovery

At-least-once vs exactly-once in pipelines, dead-letter queues, checkpoint recovery, and designing for partial failures.

4/5Overview: 30m

Delivery semantics (pipeline view)

GuaranteeMeaningTypical mechanism
At-most-onceMay lose dataFire-and-forget
At-least-onceMay duplicateRetry + idempotent sink
Exactly-onceNo dup, no lossTransactions + idempotent + checkpoint

True end-to-end exactly-once is hard — requires cooperation across source, processor, sink.

Link to Distributed Systems for theory — Kafka transactions, two-phase commit, idempotency keys.

Pipeline failure modes

FailureSymptomMitigation
Task retryDuplicate partitionsIdempotent overwrite/merge
Partial cluster lossIncomplete writeWrite to temp path, atomic rename
Upstream schema breakJob errorContract tests, quarantine
Skew/stragglerSLA missAQE, repartition
Consumer read too earlyWrong dashboardPartition publish flag / _SUCCESS marker

Dead-letter queues

Bad records → DLQ table/topic for inspection instead of failing entire batch:

parse JSON → valid → silver → invalid → dlq.raw_payload + error_reason

Checkpoint recovery (streaming)

Flink/Spark Streaming store offsets + state to durable storage. On failure, restart from last checkpoint.

Savepoints (Flink) — manual checkpoint for code upgrades with state compatibility.

_SUCCESS and publish patterns

Hive-style _SUCCESS empty file signals partition complete. Consumers filter WHERE _partition_ready.

Delta/Iceberg: transaction commit is atomic — readers see consistent snapshot.

Interview answer template

"Effective exactly-once: Kafka transactional producer + Flink checkpoint + Delta idempotent merge on event_id. Task retries can't duplicate because merge is keyed. Partial writes use replaceWhere on single dt only after job succeeds."

Further Reading

Hands-On Tasks (Optional)

Pipeline design drills and whiteboard exercises — DAG sketches, partition plans, backfill strategies. Assumes Databases and SQL fundamentals are in place.

  • Recover from a partial write

    Spark job wrote 800/1000 partitions before cluster loss. Describe detection, idempotent retry strategy, and how consumers avoid reading incomplete data.

    15m