Learn Labs
11. Batch Processing

11. Batch Processing

Chapter 11 of Designing Data-Intensive Applications — 11 sections.

"A system cannot be successful if it is too strongly influenced by a single person. Once the initial design is complete and fairly robust, the real test begins as people with many different viewpoints undertake their own experiments." — Donald Knuth

The two families of data processing:

Online systemsOffline systems (batch)
ShapeYou ask for something, and the system tries to give you an answer AS QUICKLY AS POSSIBLEA job takes READ-ONLY input data and produces output GENERATED FROM SCRATCH EVERY TIME IT RUNS
Primary metricResponse timeTHROUGHPUT — how much data per unit of time
DurationmillisecondsMinutes, hours, or even DAYS. Often scheduled periodically
FaultsRequire fault tolerance for high availabilitySome abort and restart the WHOLE job; others tolerate node crashes

A batch job typically DOES NOT MUTATE DATA the way a read/write transaction would. The output is DERIVED from the input. IF YOU DON'T LIKE THE OUTPUT, YOU CAN DELETE IT, ADJUST THE JOB'S LOGIC, AND RUN THE JOB AGAIN.

The four benefits of immutable inputs and no side effects

  1. Human fault tolerance. “If you introduce a bug and the output is wrong or corrupted, you can simply roll back to a previous version of the code and rerun the job, and the output will be correct again. Or even simpler, keep the old output in a different directory and switch back to it.” Object stores and open table formats support this as time travel.
  2. Faster feature development, because mistakes don't mean irreversible damage. “This principle of minimizing irreversibility is beneficial for Agile software development.” — Chapter 2's evolvability, made operational.
  3. The same files can feed many jobs — including monitoring jobs that calculate metrics and evaluate whether output has the expected characteristics, for example by comparing it to the previous run's output and measuring discrepancies.
  4. Efficient use of computing resources. “It's possible to batch-process data via online systems such as OLTP databases and application servers, but doing so can be much more expensive in terms of the resources required.”

The contrast that makes the first point sharp: most databases with read/write transactions do not have this property — if you deploy buggy code that writes bad data, rolling back the code will do nothing to fix that data.

Two honest costs:

  • With most frameworks, output can be processed by other jobs ONLY AFTER THE WHOLE JOB FINISHES
  • ANY CHANGE TO THE INPUT DATA — EVEN A SINGLE BYTE — REQUIRES THE JOB TO REPROCESS THE ENTIRE INPUT DATASET

The historical arc — and where MapReduce stands now

2004 — Google publishes MapReduceimplemented in Hadoop, CouchDB, MongoDBnow — Spark, Flink, warehouse query enginesMapReduce is largely obsolete and no longer used at Googlethe shift — focus moved to usabilitydataflow APIs, query languages, DataFrame APIsOrchestration: Oozie / Azkaban (Hadoop-centric) → Airflow, Dagster, Prefect.Storage: HDFS / GlusterFS / CephFS → object storage (S3).

MapReduce was “a fairly low-level programming model, less sophisticated than the parallel query execution engines found in data warehouses” — but when it was new it was “a step forward in the scale of processing achievable on commodity hardware.”

Today's engines still rely on sharding and parallel execution, but with far more sophisticated caching and execution strategies. “As these systems matured, operational concerns have been largely solved, so focus has shifted toward usability.” And “scalable cloud data warehouses like BigQuery and Snowflake are blurring the line between data warehouses and batch processing.”

Figure 11.0.2The historical arc — and where MapReduce stands now

On this page