Contents

Architecture & System Design › Events & Integration · also in Batch & Distributed Processing

Distributed Batch Processing

Engines like Spark that split big jobs across many machines.

Also known as: distributed batch processing, batch compute cluster, batch jobs distributed

Distributed batch processing crunches huge datasets by splitting work across many machines: partition the input, compute in parallel, shuffle intermediate results, aggregate — the MapReduce-shaped pattern behind Hadoop, Spark and modern dataframes. Throughput comes from parallelism; fault tolerance from re-execution of failed partitions.

input splits → map (parallel) → shuffle/sort → reduce (parallel) → output
failed partition → recompute from lineage (not whole-job restart)

Design centres on data movement (shuffle is the expensive step — minimise and colocate), skew (uneven partitions stall on stragglers), and the batch/stream boundary (bounded historical backfills vs unbounded live streams often share engines now).

The classic mistakes:

  • Shuffling everything. Global repartitioning when local aggregation would do multiplies network traffic hundredfold. Combine locally before shuffling (map-side aggregation).
  • Skew blindness. One giant key’s partition runs for hours while hundreds idle. Salt hot keys, split skewed joins (broadcast the small side), monitor partition sizes.
  • Tiny-file fragmentation. Millions of small files murder scheduling and metadata; coalesce inputs (and outputs) to sane chunk sizes.
  • No idempotence. Retried partitions double-writing unguarded outputs corrupt datasets. Write idempotently (deterministic paths, atomic commits, overwrite semantics).
  • Treating batch as streaming. Hourly jobs papering over needs for minutes-fresh data accumulate scheduling debt; use stream processing where freshness demands it.
  • Resource static allocation. Fixed giant clusters for spiky jobs waste enormously; ephemeral job-scoped clusters (or serverless batch) match spend to work.
  • Untested at scale. Logic verified on samples breaks on skew, nulls and encoding edge cases at full volume. Test on representative slices, validate outputs statistically.

When to use it: bounded, high-volume transformations where minutes-to-hours latency suffices — ETL, backfills, training data, reports. Streaming serves freshness; batch serves completeness and economy.