Learn Labs
11. Batch Processing

11.5 Technology deep dives


5.1 Apache Spark

Problem it solves. Run an entire multi-stage workflow as one optimized job, keeping intermediate data in memory, with fault tolerance that doesn't require writing every stage to a replicated filesystem.

Why wasn't MapReduce enough? §3.1's two limits plus §3.2's six: mandatory sort between every stage, no operator fusion, no locality planning, all intermediate state written to the DFS (replicated, on disk, on every replica), no pipelining, and a new JVM per task.

How it works internally. The RDD/DataFrame is the unit; transformations build a lazy DAG, and only an action triggers execution. Catalyst optimizes the logical plan (predicate pushdown, column pruning, join reordering) and Tungsten generates whole-stage code (Ch 4's query compilation). The DAG is cut into stages at shuffle boundaries; within a stage, narrow transformations are fused into one task (advantage ②). Lineage — the record of how each partition was computed — is what lets Spark recompute lost intermediate data instead of replicating it (§2.3's fault-tolerance table). Adaptive Query Execution re-plans mid-job using actual shuffle statistics.

Deployment. On YARN, Kubernetes, or standalone; executors sized as cores × memory; data in an object store via Parquet/Iceberg; driven by Airflow/Dagster.

Monitoring.

  • Task-duration skew within a stage — the single most diagnostic metric. 199 tasks in 10 s and one in 4 hours means a skewed key.
  • Shuffle read/write bytes and spill (memory and disk) — spill volume is your "this doesn't fit" signal
  • GC time as a fraction of task time (>10% means executor memory is wrong)
  • Stage retry count and FetchFailedException rate (lost shuffle files → whole-stage recomputation)
  • Executor loss reason — distinguishing preemption (spot instances, §2.3) from OOM matters enormously

Scaling. More executors helps until shuffle becomes the bottleneck; then the lever is reducing shuffle volume (broadcast joins, better partitioning, pre-aggregation) rather than adding machines.

What actually breaks.

  • Data skew. One key holds 40% of the rows; one task runs for hours while the cluster idles. Fixes: salting the key, AQE skew join handling, or a broadcast join.
  • OOM on the driver from collect() on a large DataFrame — the classic mistake of pulling a distributed dataset into one JVM.
  • Cascading recomputation when an executor holding shuffle output dies: every downstream task that needed that block fails, and Spark recomputes the whole upstream stage. On spot instances this can dominate runtime. (External shuffle service / shuffle-service-on-object-storage exists precisely for this.)
  • The small-file problem on write — repartition before writing, or you produce 10,000 tiny Parquet files (Ch 4 §9.3).
  • Eager-vs-lazy confusion for people arriving from Pandas (§3.7) — a chain of transformations that "runs instantly" and then takes an hour on the first action.
  • spark.sql.shuffle.partitions = 200 (the default) being wildly wrong for both tiny and huge jobs.

5.2 Apache Airflow (workflow orchestration)

Problem it solves. Schedule and manage the dependency graph between jobs — which the per-job schedulers (YARN, Spark's own) explicitly do not do (§2.3).

Why wasn't cron enough? Cron has no notion of dependencies, no backfill, no retry semantics, no visibility into which of 100 jobs failed, and no way to express "run when all three upstream jobs succeed."

How it works internally. A DAG of tasks defined in Python. The scheduler parses DAG files, creates DagRuns per schedule interval, and marks tasks runnable when upstream dependencies are met; executors (Celery/Kubernetes) run them; state lives in a metadata database. Operators encapsulate integrations — §4.1's "built-in source, sink, and query operators for MySQL, PostgreSQL, Snowflake, Spark, Flink, and dozens of others."

Monitoring. Task duration trend per task (a slowly growing job is a future incident); SLA misses; scheduler loop latency and DAG parse time (heavy top-level code in DAG files silently throttles the whole scheduler); queued-vs-running task counts; retry counts by task — distinguishing transient from systematic failure, exactly the distinction §4.1 says makes debugging easy.

Scaling. More workers for task throughput; but the scheduler and metadata DB are the real ceiling. Keep DAG files light.

What actually breaks.

  • Top-level code in DAG files — an API call or heavy import at module scope runs on every parse, every few seconds, for every DAG.
  • Tasks that aren't idempotent, meeting Airflow's retry behaviour → duplicated side effects. This is why §1's "batch jobs avoid side effects" matters operationally.
  • Backfills running hundreds of DagRuns concurrently and melting a source database.
  • Timezone and execution_date semantics — the single most common source of "the job ran but processed the wrong day's data."
  • Using Airflow as the compute engine (heavy work inside a PythonOperator) rather than as an orchestrator.
  • A silently succeeding job that produced no data — success ≠ correctness; this is why §1's third benefit (monitoring jobs comparing to the previous run) exists.

5.3 HDFS vs S3 as the batch storage layer

Problem each solves. Store petabyte datasets durably and read them at aggregate bandwidths a single machine can't reach.

Why HDFS first? Data locality (§2.2) — run the task on the node holding the block, so you ship the code to the data, not the data to the code. In 2006, when datacenter networks were slow relative to disk, this was decisive.

Why S3 now? §2.2: "modern datacenter networks are very fast, so this is often acceptable" — and decoupling lets CPU and storage scale independently. Plus no NameNode to operate, no rebalancing, and vastly better durability economics.

How they differ operationally — the list that causes real bugs:

HDFSS3
RenameAtomicCopy + delete, nonatomic, O(n) for a "directory"
ListingDirectory listingRecursive prefix scan; paginated; eventually complete
AppendSupportedGenerally not
MetadataNameNode (SPOF, memory-bound on file count)Managed, effectively unbounded
LocalityYesNo
Cost modelCluster you runPer-request fees + storage; batching matters

Monitoring. HDFS: NameNode heap and total file/block count (the small-file problem is a NameNode memory problem here, not just a query-planning one), under-replicated blocks, DataNode volume failures. S3: request rate and 503 SlowDown per prefix, first-byte latency, incomplete multipart uploads, cost per prefix.

What actually breaks.

  • The commit protocol. Spark/Hadoop output committers historically relied on atomic rename to publish results. On S3, rename isn't atomic, so the naive committer is both slow and unsafe — hence the S3A committers and, better, table formats (Iceberg/Delta) whose atomic metadata commit replaces rename entirely. This is the classic HDFS→S3 migration failure.
  • Millions of small files killing the NameNode (HDFS) or query planning and request cost (S3).
  • Eventual-consistency assumptions in old code — "list then process" silently missing files.
  • Hot prefix throttling when a job writes every part file under one date prefix.

5.4 Kubernetes / YARN as the job orchestrator

Problem it solves. Decide where and when each task runs, enforce isolation, and reclaim resources on failure — §2.3's three components.

Why not just SSH and run it? No global view of capacity, no fairness, no isolation, no automatic retry on node loss, no bin-packing.

How it works internally. Resource manager holds cluster state in a consensus store (ZooKeeper for YARN, etcd for Kubernetes — Ch 10 §4's "outsource the consensus"). Scheduler matches pending tasks to nodes using heuristics (§2.3: FIFO, DRF, priority queues, capacity/quota, bin-packing — because the optimal problem is NP-hard). Executors enforce limits with cgroups.

Monitoring. Pending/unschedulable task count (and the reason — "no node fits the request" is a capacity-planning signal, not a bug); queue wait time per priority class; node resource fragmentation; preemption rate (§2.3 — expected and healthy on spot capacity, alarming on on-demand); scheduler decision latency.

What actually breaks.

  • Gang-scheduling deadlock (§2.3): two jobs each hold half the cores they need, neither can proceed, neither releases. Requires a gang/coscheduling plugin with all-or-nothing admission.
  • Starvation of large jobs by a stream of small ones.
  • Preemption thrash — low-priority tasks killed and restarted repeatedly, so they never finish and all the work is wasted.
  • Resource requests set from guesswork — too high wastes the cluster, too low gets tasks OOMKilled or CPU-throttled.
  • The centralized resource manager as a bottleneck — §2.3 warns about it explicitly, and it shows up as scheduler latency growing with cluster size.

On this page