This complete guide covers all 13 Data Engineering tutorials on TheCodeForge, organised by topic.
Data engineering failures have a signature: the job does not crash, it gets slow, and then it gets slow enough to matter. A Spark stage where 199 tasks finish in seconds and one runs for forty minutes. A Kafka consumer group whose lag climbs while every broker sits at 10% CPU. A table that reads fine on Monday and takes six minutes on Friday because it now contains 400,000 files of 3 KB each.
All three are distribution problems. Distributed systems parallelise by partitioning, and every one of these symptoms is partitioning that has stopped matching the data. That is why the fixes in this track are rarely about tuning a memory flag and usually about changing how the work is divided — a salted join key, a rebalanced consumer, a compaction pass.
Spark's Job aborted: task failed 4 times is a wrapper, not a cause. The real exception is inside the failed task's log, and raising spark.executor.memory because the word 'memory' appeared somewhere is how teams end up paying for a cluster four times too large that still fails.
The Spark UI answers the question directly. Open the slowest stage and compare the max task duration and shuffle-read size against the median. If the max is orders of magnitude larger, you have skew, and no amount of executor memory fixes skew — one key is simply too big for one task.
| What the stage view shows | Diagnosis | Correct lever |
|---|---|---|
| Max task ≫ median task, one long tail | Key skew — a few values dominate the join or group | Salt the hot keys, broadcast the small side, or enable adaptive skew join handling |
| All tasks slow, high GC time | Genuinely undersized executors or too much cached data | More memory per executor, or stop caching what is read once |
| Thousands of tiny tasks, low CPU | Small-files problem — one task per file | Compact on write; coalesce output; use a table format that supports file compaction |
| Spill to disk dominates shuffle | Partition count too low for the data volume | Raise shuffle partitions, or let adaptive execution coalesce them for you |
If lag is growing and the brokers are bored, the bottleneck is downstream. The usual shape is a consumer whose per-record work is slower than the produce rate, and the failure mode compounds: processing takes longer than max.poll.interval.ms, the coordinator decides the member is dead, the group rebalances, the partially-processed batch is reprocessed by someone else, and the lag climbs further while nobody makes progress.
That loop is also what produces CommitFailedException: by the time the consumer tries to commit, its partitions have been reassigned and the commit is rejected. The exception is a consequence of the timeout, so fixing the commit call fixes nothing.
| Setting | What it controls | Tune it when |
|---|---|---|
max.poll.records | Records handed to one poll() | Per-record work is slow — shrink the batch so the loop returns inside the interval |
max.poll.interval.ms | How long the coordinator waits before assuming death | Batches are legitimately long-running and cannot be split |
session.timeout.ms / heartbeat | Liveness, sent on a background thread | Network is genuinely flaky — not for slow processing, which heartbeats no longer mask |
| Partition count / consumer count | Achievable parallelism | Every consumer is saturated — a group cannot exceed one consumer per partition |
Kafka's transactional producer plus read-committed consumers give you exactly-once within Kafka. What breaks it is almost always the boundary: a consumer that reads a record, writes to Postgres, and then commits its offset has two systems and no shared transaction, so a crash between the write and the commit duplicates the record on restart.
There are only two honest answers. Either make the write idempotent — an upsert on a natural key, or a dedupe table keyed by the record's identity — or store the offset in the same transactional store as the data so one commit covers both. Enabling enable.idempotence and hoping is the third option, and it does not work, because producer idempotence protects against retry duplication inside Kafka, not against your consumer's side effects.
Airflow tasks stuck in queued and dbt's model depends on a node that was not found look unrelated, and are the same kind of problem: the tool's model of the world no longer matches reality. Airflow queued-forever means the scheduler handed the task to an executor slot that does not exist or cannot accept it — pool exhausted, concurrency limit reached, worker not running, or a queue name with no worker listening.
dbt's missing node means the ref() graph references a model that is not in the current selection or not in the project at all — commonly a renamed file, a model excluded by the selector you passed, or a package not installed. In both cases the fix is to inspect the tool's own view of state rather than the code.
max.poll.records instead.df.printSchema() at the failing point rather than trusting the schema you expect.Every tutorial starts with a plain-English analogy — then real code, then interview questions.
Browse Data Engineering Tutorials →