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 systems | Offline systems (batch) | |
|---|---|---|
| Shape | You ask for something, and the system tries to give you an answer AS QUICKLY AS POSSIBLE | A job takes READ-ONLY input data and produces output GENERATED FROM SCRATCH EVERY TIME IT RUNS |
| Primary metric | Response time | THROUGHPUT — how much data per unit of time |
| Duration | milliseconds | Minutes, hours, or even DAYS. Often scheduled periodically |
| Faults | Require fault tolerance for high availability | Some 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
- 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.
- 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.
- 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.
- 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
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.”
Sections
- 11.11Batch Processing with Unix Tools
- 11.23Batch Processing in Distributed Systems
- 11.35Batch Processing Models
- 11.42Batch Use Cases
- 11.5Deep divesTechnology deep dives
- 11.6Failure catalogProduction failure catalog for this chapter
- 11.7Decision sheet1Decision cheat sheetBatch when data freshness isn't important and the input is bounded.
- 11.8Worked examplesWorked examples
- 11.9Self-testSelf-test
- 11.10TerminologyTerminology introduced here
- 11.11Forward linksForward links