Home› Data Engineering› Complete Guide
Complete Guide

Complete Data Engineering Tutorial

This complete guide covers all 13 Data Engineering tutorials on TheCodeForge, organised by topic.

Learning Roadmap
Beginner → Build a strong foundation
Intermediate → Deepen your understanding with practical topics
Advanced → Master advanced concepts and real-world applications
13
Topics
4
Beginner
6
Intermediate
3
Advanced
Jump to section
Kafka (4)Airflow (1)dbt (1)Spark (5)Databricks (2)

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.

Read the skew before you touch a memory setting

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 showsDiagnosisCorrect lever
Max task ≫ median task, one long tailKey skew — a few values dominate the join or groupSalt the hot keys, broadcast the small side, or enable adaptive skew join handling
All tasks slow, high GC timeGenuinely undersized executors or too much cached dataMore memory per executor, or stop caching what is read once
Thousands of tiny tasks, low CPUSmall-files problem — one task per fileCompact on write; coalesce output; use a table format that supports file compaction
Spill to disk dominates shufflePartition count too low for the data volumeRaise shuffle partitions, or let adaptive execution coalesce them for you
python
# Find the skew instead of guessing at it
(df.groupBy("customer_id").count()
   .orderBy(F.desc("count"))
   .show(20, truncate=False))

# Salt the hot side so one key spreads across N tasks
N = 64
left  = df.withColumn("salt", (F.rand() * N).cast("int"))
right = (dim.withColumn("salt", F.explode(F.array(*[F.lit(i) for i in range(N)]))))
joined = left.join(right, ["customer_id", "salt"], "left")

# Or let Spark handle it — on by default in Spark 3.2+, worth asserting
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")

Kafka consumer lag with idle brokers is always the consumer

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.

SettingWhat it controlsTune it when
max.poll.recordsRecords handed to one poll()Per-record work is slow — shrink the batch so the loop returns inside the interval
max.poll.interval.msHow long the coordinator waits before assuming deathBatches are legitimately long-running and cannot be split
session.timeout.ms / heartbeatLiveness, sent on a background threadNetwork is genuinely flaky — not for slow processing, which heartbeats no longer mask
Partition count / consumer countAchievable parallelismEvery consumer is saturated — a group cannot exceed one consumer per partition
In practiceParallelism inside a group is capped by partition count. Adding a ninth consumer to a topic with eight partitions gives you an idle process, not more throughput. If you are already at the cap and still behind, the topic needs more partitions — which also means accepting that ordering is only guaranteed within a partition.

Exactly-once is a property of the whole pipeline

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.

Orchestration and transformation: failures of state, not logic

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.

bash
# Airflow: is anything actually able to run this task?
airflow pools list
airflow celery inspect active            # Celery executor: are workers alive?
airflow tasks states-for-dag-run <dag_id> <run_id>

# Check the three limits that produce silent queueing
#   parallelism, max_active_tasks_per_dag, pool slots, and the task's queue

# dbt: what does the graph think exists?
dbt ls --select +my_model               # upstream closure of the model
dbt deps                                # install packages before compiling
dbt compile --select my_model           # see the resolved SQL and refs

Frequently Asked Questions

What actually causes the small-files problem?
Writing many partitions, each producing at least one file per task, on a schedule. A stream micro-batching every minute with 200 shuffle partitions creates 288,000 files a day. Reads then pay per-file overhead — object store listing, open, metadata — which dominates when files are smaller than a few megabytes. Compact on a schedule, coalesce before writing, and prefer a table format with built-in compaction.
Why does my Spark job fail only in production when the data looks the same?
Because it usually is not the same. Production has the long tail: the one customer with ten million rows, the null key that collects every unmatched record, the day a upstream backfill tripled volume. Sample-based testing hides skew by construction. Run the group-by-count on the join key in production and compare the top value against the median before concluding the data is identical.
Can I just increase max.poll.interval.ms and move on?
Sometimes, and it is worth understanding the trade. A longer interval means genuinely dead consumers go undetected for longer, so a real failure stalls its partitions for that duration. It is the right lever when a batch is legitimately slow and atomic. It is the wrong lever when the real problem is that each record takes 400ms and you are pulling 500 at a time — shrink max.poll.records instead.
What is a Delta concurrent append exception telling me?
That two writers committed to the same table version and the optimistic concurrency check caught the conflict. It is the table format doing its job. Fix it by partitioning writers so they touch disjoint partitions, serialising the writes, or retrying with backoff — Delta's conflict detection is partition-aware, so disjoint writers usually stop colliding once the partitioning matches the write pattern.
Is 'cannot resolve column name' ever not a typo?
Often. Case sensitivity differs between Spark's own resolution and the underlying source; a column may exist in the file but not in the catalog's cached schema; a nested field needs dotted access the error does not hint at; and a column present in most Parquet files but absent in older ones disappears depending on schema merging. Print df.printSchema() at the failing point rather than trusting the schema you expect.
How do I know whether to scale the cluster or fix the code?
Compare max task time to median task time in the slowest stage. If they are close, the cluster is genuinely the limit and scaling helps proportionally. If max is many multiples of median, the work is not divisible as written and a bigger cluster adds idle executors waiting on one task. That single ratio settles the argument faster than any other measurement.

Kafka

Airflow

dbt

Spark

Databricks

Also Explore
Database 139 tutorials → DevOps 304 tutorials → Python 171 tutorials → System Design 145 tutorials → ML / AI 206 tutorials → Observability 10 tutorials →
Start from the beginning

Every tutorial starts with a plain-English analogy — then real code, then interview questions.

Browse Data Engineering Tutorials →