11.7 Decision cheat sheet
Batch when data freshness isn't important and the input is bounded.
Batch or stream? Batch when data freshness isn't important and the input is bounded. Stream when you need second-level latency or the input is unbounded (Ch 12). Note that "the US banking network runs almost entirely on batch jobs" — freshness matters far less often than people assume.
Single machine or distributed?
Try the Unix pipeline / DuckDB / Polars first. GNU sort spills to disk and parallelizes across cores; the bottleneck is disk read rate. Go distributed only when data exceeds one machine's disk or the job exceeds your time budget. (Ch 1: more nodes are not always faster.)
Hash aggregation or sort? Hash when the number of distinct keys fits in memory — note the working set depends on distinct keys, not record count. Sort when it doesn't; sorting degrades gracefully to disk with sequential I/O.
HDFS or object store? Object store by default now — decoupled scaling, no NameNode, better durability economics. HDFS only if data locality genuinely dominates your bandwidth budget. Either way, use a table format (Iceberg/Delta) so you don't depend on atomic rename.
MapReduce, dataflow engine, or warehouse SQL? Never raw MapReduce for new work. Warehouse SQL for relational analytics your analysts will touch. Spark/Flink for large jobs where warehouse cost bites, for non-relational/multimodal data, for iterative graph and ML work, and for anything awkward in SQL.
How do I get batch output into production?
- ✗ Write directly to the production database from tasks.
- ✔ Push to a stream (Kafka) → downstream systems ingest at their own rate. Add a completion notification so consumers keep the data invisible until the job finishes — read-committed semantics across the boundary.
- ✔ Or build the database file inside the job and bulk-load it, swapping versions atomically — fast and clean, but hard to update incrementally.
- ✔ A hybrid (Venice, for example) when you need both bootstrap and incremental loads.
Spot instances or on-demand? Spot for batch — it's exactly what batch is good at (not time-sensitive, restartable at task granularity, uses otherwise-idle capacity). But budget for preemptions being more frequent than hardware faults, and make sure intermediate-data loss doesn't cascade into whole-stage recomputation.