Learn Labs
11. Batch Processing

11.3 Batch Processing Models

3.1 MapReduce

The four steps — and how they map exactly onto the Unix pipeline:

read input filesa newline separates one log line from the nextmapperextract a key and value from each recordsort all pairs by keyimplicit — you do not write itreduceriterate over the sorted pairsinput format parserfiles in HDFS/S3; Parquet or Avroawk '{print $7}'URL as key, empty valuesortuniq -cMAPREDUCETHE UNIX EQUIVALENT

The mapper may emit any number of pairs, including none, and keeps no state between records — so each record is handled independently and many mappers can run in parallel. The sort is not something you write: “the output from the mapper is always sorted before it is given to the reducer.”

Sorting has made same-key records adjacent, “so it is easy to combine those values without having to keep a lot of state in memory.” Reducers for different keys also run in parallel. The second sort -r -n in the Unix example becomes a second MapReduce job, using the first job's output as its input.

“The role of the mapper is to prepare the data by putting it into a form suitable for sorting; the role of the reducer is to process the data that has been sorted.”

Figure 11.3.1The four steps — and how they map exactly onto the Unix pipeline

Why the functional-programming heritage matters:

Lisp introduced map and reduce (or fold) as higher-order functions on lists.

The functional principle of AVOIDING MUTABLE STATE IS WHAT ENABLES PARALLEL EXECUTION. As every call depends ONLY on the data the framework EXPLICITLY PASSES, the framework is free to RUN INDEPENDENT CALLS IN PARALLEL ON DIFFERENT NODES — AND IF A TASK FAILS, TO CALL THE MAPPER OR REDUCER AGAIN WITH THE SAME INPUT ON ANOTHER NODE.

Two damning limitations:

  • Implementing a complex job with raw MapReduce APIs is QUITE LABORIOUS — ANY JOIN ALGORITHMS WOULD NEED TO BE IMPLEMENTED FROM SCRATCH
  • MapReduce is QUITE SLOW compared to modern batch processors. One reason: ITS FILE-BASED I/O PREVENTS JOB PIPELINING — processing output in a downstream job before the upstream job is complete

They handle AN ENTIRE WORKFLOW AS ONE JOB, rather than breaking it into independent subjobs. Since they EXPLICITLY MODEL THE FLOW OF DATA through several processing stages, they are known as DATAFLOW ENGINES.

Like MapReduce they support a low-level record-at-a-time API, but they also offer HIGHER-LEVEL OPERATORS such as JOIN and GROUP BY. They parallelize by SHARDING inputs and COPY THE OUTPUT OF ONE TASK OVER THE NETWORK to become another's input. UNLIKE MAPREDUCE, OPERATORS NEED NOT TAKE THE STRICT ROLES OF ALTERNATING MAP AND REDUCE — they can be ASSEMBLED IN MORE FLEXIBLE WAYS.

The six concrete advantages over MapReduce:

  1. “Expensive work such as sorting needs to be performed only where it is required, rather than always happening by default between every map and reduce stage.”
  2. “When several operators in a row don't change the sharding of the dataset (such as map or filter), they can be combined into a single task, reducing data copying overheads.” — operator fusion.
  3. “Because all joins and data dependencies are explicitly declared, the scheduler has an overview of what data is required where, so it can make locality optimizations — placing a consumer on the same machine as its producer so the data can be exchanged through a shared memory buffer rather than the network.”
  4. “Intermediate state can usually be kept in memory or on local disk, which requires less I/O than writing it to a DFS or object store, where it must be replicated to several machines and written to disk on each replica.” MapReduce already does this for mapper output; dataflow engines generalize it to all intermediate state.
  5. “Operators can start executing as soon as their input is ready; no need to wait for the entire preceding stage to finish.” — pipelining, the fix for MapReduce's flaw.
  6. “Existing processes can be reused to run new operators, reducing startup overheads compared to MapReduce, which launches a new JVM for each task.”

3.3 Shuffling — the foundational algorithm

⚠️ SHUFFLE IS NOT RANDOM. "When you shuffle a deck of cards, you end up with a RANDOM order. In contrast, the shuffle we're talking about PRODUCES A SORTED ORDER, WITH NO RANDOMNESS."

A distributed sorting algorithm where BOTH THE INPUT AND THE OUTPUT ARE SHARDED. Batch processors must sort datasets PETABYTES in size.

input shard 1input shard 2input shard 3mapper 1mapper 2mapper 3reducer 1reducer 2reducer 3output r1output r2output r3One map task per input shard — the count is fixed by the input. The number of reduce tasks is configured by the job’s author, and may differ.Each arrow across the middle is a file: “m1,r2” is mapper 1’s data destined for reducer 2.
  • Each mapper creates a separate output file on its local disk for every reducer.
  • A hash of the key typically determines which reducer file a pair goes to.
  • While writing, the mapper also sorts the pairs within each file — using the LSM technique of Chapter 4: sorted in-memory structure → sorted segment files → progressively merged.
  • After each mapper finishes, reducers connect to it and copy the appropriate file to their local disk.
  • Once a reducer has its share from all mappers, it merges them preserving sort order, mergesort-style. Same-key pairs are now consecutive even though they came from different mappers.
  • The reducer is called once per key, with an iterator over all its values.
  • Reducer output is written sequentially to one file per reduce task; these become the shards of the job's output.
Figure 11.3.33.3 Shuffling — the foundational algorithm

Modern dataflow engines and cloud warehouses are more sophisticated: BigQuery has optimized its shuffle to KEEP DATA IN MEMORY and to write to EXTERNAL SORTING SERVICES — speeding up shuffling and REPLICATING SHUFFLED DATA TO PROVIDE RESILIENCE.

3.4 Joins and grouping — the sort-merge join

The scenario: activity events (clickstream) on the left, a user database on the right. In star-schema terms, the event log is the FACT TABLE and the user database is one of the DIMENSIONS.

mapper A — activity eventsemits (user_id, page_view_url)mapper B — user databaseemits (user_id, date_of_birth)the shufflesame key → same reducerreducer for user_id = 42user record first, then events by timestampThe shuffle brings all pairs with the same key to the same reducer, no matter which shard they were on originally.

Inside that reducer:

  1. A secondary sort arranges records so the reducer always sees the user database record first, followed by the activity events in timestamp order.
  2. The first value is the date of birth → store it in a local variable.
  3. Iterate over the activity events, outputting each URL with the date of birth.

“A reducer processes all the records for a particular user ID in one go, so it needs to keep only one user record in memory at any one time, and it never needs to make any requests over the network.” This is a sort-merge join — mapper output sorted by key, reducers merging the sorted lists from both sides.

The next job does the GROUP BY: shuffle by URL, then reducers iterate over all page views for one URL, keeping a counter per age group.

Figure 11.3.43.4 Joins and grouping — the sort-merge join

3.5 Query languages

With the problem of physically operating batch processes at scale CONSIDERED MORE OR LESS SOLVED, attention has turned to IMPROVING THE PROGRAMMING MODEL.

MapReduce, dataflow engines, and cloud warehouses have all embraced SQL AS THE LINGUA FRANCA. It's a natural fit: legacy warehouses used SQL, analytics and ETL tools already support it, and ALL DEVELOPERS AND ANALYSTS KNOW IT.

Two payoffs:

  1. Human: less code, and INTERACTIVE USE — an efficient and natural way for business analysts, product managers, sales and finance teams to explore data. SQL support has made distributed batch systems suitable for EXPLORATORY QUERIES.
  2. Machine: the translation from query → syntax tree → physical operators ALLOWS THE ENGINE TO OPTIMIZE. Hive, Trino, Spark, and Flink have COST-BASED OPTIMIZERS that analyze the properties of join inputs and AUTOMATICALLY DECIDE WHICH ALGORITHM IS MOST SUITABLE. Optimizers MIGHT EVEN CHANGE THE ORDER OF JOINS so the amount of INTERMEDIATE STATE IS MINIMIZED.

Niche languages: Apache Pig (relational operators, pipelines specified step by step rather than as one big SQL query), Morel (a modern language influenced by Pig), JSON query languages (jq, JMESPath, JSONPath), and graph languages (Apache TinkerPop's Gremlin).

3.6 Batch processing and cloud warehouses converge

ThenNow
Data warehouses: specialized hardware appliances, SQL over relational data.Batch frameworks support SQL, and achieve good performance on relational queries using columnar formats (Parquet) and optimized execution — compilation and vectorization, Chapter 4.
Batch frameworks: greater scalability and flexibility, logic in a general-purpose language, arbitrary data formats.Warehouses became scalable by moving to the cloud and implementing the same scheduling, fault tolerance, and shuffling techniques. Many use distributed filesystems.
—Warehouses also adopted alternative processing models: BigQuery has a DataFrames library, Snowflake's Snowpark integrates with Pandas, and Airflow/Prefect/Dagster integrate with warehouses.

But they haven't fully merged. What SQL/warehouses still struggle with:

  • ITERATIVE GRAPH ALGORITHMS such as PageRank, complex ML tasks
  • AI DATA PROCESSING — nonrelational and MULTIMODAL data such as images, video, audio
  • ROW-BY-ROW COMPUTATION is less efficient with column-oriented storage
  • Cloud warehouses TEND TO BE MORE EXPENSIVE. It can be MORE COST-EFFICIENT to run large jobs in Spark or Flink

The decision often comes down to COST, CONVENIENCE, EASE OF IMPLEMENTATION, AND AVAILABILITY. Most large enterprises have MANY data processing systems, giving them flexibility. SMALLER COMPANIES OFTEN GET BY WITH JUST ONE.

3.7 DataFrames in a distributed setting

Why they exist here: "Data scientists wanted to interact with the LARGE DATASETS found in batch environments USING THE DATAFRAME APIs THEY WERE USED TO, since SQL AND MAPREDUCE ARE NOT WELL SUITED TO THEIR NEEDS."

⚠️ Two traps:

① "LOCAL DATAFRAMES ARE USUALLY INDEXED AND ORDERED, WHILE DISTRIBUTED DATAFRAMES ARE GENERALLY NOT. THIS CAN LEAD TO PERFORMANCE SURPRISES WHEN MIGRATING TO BATCH FRAMEWORKS."

② "PANDAS EXECUTES OPERATIONS IMMEDIATELY when DataFrame methods are called; SPARK FIRST TRANSLATES ALL THE API CALLS INTO A QUERY PLAN AND RUNS QUERY OPTIMIZATION before executing." (Eager vs lazy.)

Hybrid execution: Daft supports BOTH client- and server-side computation — smaller in-memory operations on the client, larger datasets on a server. Columnar formats such as APACHE ARROW offer a UNIFIED DATA MODEL that both execution engines can share.


On this page