Skip to content
Reliable Data Engineering
Practice problem hard oombroadcastjoinscachingmemory
Practise with timer, notes and rubric

Debug a Nightly Job: Driver OOM, Broadcast Timeouts and a Stale Cache

Difficulty: Hard · Topics: OOM, broadcast joins, caching, memory · Time: 40 minutes

Scenario

A nightly PySpark job builds gold.customer_daily by joining events (3 TB/day) with customers (dimension, ~400M rows, 60 GB in Parquet) and enriching with a scoring step. It ran fine for a year. After last week’s customer migration (an extra 300M rows), three things started happening on different nights:

  1. The job fails within 10 minutes with java.lang.OutOfMemoryError: Java heap space in the driver log, or with Could not execute broadcast in 300 secs.
  2. On nights it gets past the join, stage 7 (the scoring step, a pandas UDF) fails with Container killed by YARN for exceeding memory limits. 21.3 GB of 21 GB physical memory used on several executors, then FetchFailedException loops.
  3. An analyst reports that a notebook dashboard built on the same tables shows yesterday’s totals even after the job succeeds.

Evidence

spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 2 * 1024**3)   # set by someone "to speed up joins"
customers = spark.table("silver.customers").filter("is_active")
joined = events.join(customers, "customer_id")
scored = joined.groupBy("customer_id").applyInPandas(score_fn, schema)   # stage 7

Your task

Explain the root cause of each symptom, fix each one properly (not just “add memory”), and say what you’d put in place so it doesn’t recur.

Hints

Hint 1

Where is a broadcast relation built before it’s shipped? What size does the planner think customers is, and what size is it really, especially once deserialized into a hash table?

Hint 2

Where does pandas UDF memory live, and what does max 9 GB vs median 300 MB tell you about groupBy(customer_id)?

Solution

1. Driver OOM / broadcast timeout: a broadcast from stale statistics.

Fix:

2. Containers killed in stage 7: skew plus Python memory in overhead.

Fix:

3. Stale dashboard: the cache is a snapshot.

Fix:

Prevent recurrence:

What interviewers look for