Learn Labs
11. Batch Processing

11.2 Batch Processing in Distributed Systems

Batch join strategy

strategy
8 GB
worker RAM
data across the network200 GB
job time20 min
Safe

No shuffle of the large side. The 2 GB table is copied to all 100 workers and held in a hash table, so the 500 GB side is read once, locally. Skew does not matter here — every worker has the whole small side.

A distributed join is decided by whether the small side fits in a worker's memory. If it does, ship it everywhere and skip the shuffle; if it doesn't, both sides get sorted into the same partitions.

The organizing analogy — the distributed operating system

Single machineDistributed framework
Storage via the filesystem interfaceDistributed filesystem or object store
A scheduler allocating CPUA job orchestrator — scheduler + resource manager + task executors
Programs connected by pipesTasks sending data via the filesystem or other communication channels

“You can think of these frameworks as distributed operating systems.”

2.1 Distributed filesystems

The local filesystem stack, layer by layer:

VFS — virtual filesystema common API, whatever the filesystem underneathfilesystem layerfiles into blocks; inodes, directories, metadata (ext4, XFS)page cacherecently accessed blocks kept in memoryblock device driversraw block read/writethe DFS protocolNameNode (Hadoop)3FS metadata servicedata nodes' own OS page cachesplus extra tiers — JuiceFS client-side + local diskLOCAL FILESYSTEM STACKDFS ANALOGUEEvery local layer has a distributed counterpart — which is why one DFS can replace another behind the same protocol.
Figure 11.2.2The local filesystem stack, layer by layer

Block sizes — and why they're so much bigger:

SystemBlock size
ext44,096 bytes
JuiceFS, many object stores4 MB
HDFS128 MB

Larger blocks mean LESS METADATA to keep track of, which MAKES A BIG DIFFERENCE ON PETABYTE-SIZED DATASETS. Larger blocks also LOWER THE OVERHEAD OF SEEKING TO A BLOCK RELATIVE TO READING IT.

And unlike physical devices, DFSs DON'T need to write partial blocks: a 900 MB file with 128 MB blocks has SEVEN blocks of 128 MB and ONE BLOCK OF 4 MB.

Data nodes: each machine runs a daemon exposing an API to read/write blocks as files on its local filesystem — HDFS calls them DataNodes, GlusterFS calls them glusterfsd.

The protocol as the pluggable interface:

Distributed filesystems must expose a protocol so batch systems can read and write. THIS PROTOCOL ACTS AS A PLUGGABLE INTERFACE; ANY DFS MAY BE USED SO LONG AS IT IMPLEMENTS THE PROTOCOL. For example, AMAZON S3's API HAS BEEN WIDELY ADOPTED by MinIO, Cloudflare R2, Tigris, Backblaze B2, and many others.

POSIX compatibility via FUSE or NFS. (NFS was originally developed to let multiple clients read/write on a SINGLE SERVER; more recently Amazon EFS and Archil provide NFS-compatible implementations that are FAR MORE SCALABLE — clients still connect to one endpoint, but underneath these systems talk to distributed metadata services and data nodes.)

DFS vs NAS/SAN — the shared-nothing point:

Distributed filesystems are based on the SHARED-NOTHING principle, in contrast to the SHARED-DISK approach of NAS and SAN. Shared-disk storage uses a CENTRALIZED STORAGE APPLIANCE, often with CUSTOM HARDWARE and special network infrastructure such as FIBRE CHANNEL. The shared-nothing approach requires NO SPECIAL HARDWARE, only computers connected by a conventional datacenter network.

Many DFSs are built on COMMODITY HARDWARE — less expensive but with HIGHER FAILURE RATES. To tolerate machine and disk failures, file blocks are REPLICATED on multiple machines. THIS ALSO ALLOWS SCHEDULERS TO MORE EVENLY DISTRIBUTE WORKLOADS, since they can execute a task on ANY node holding a replica of the task's input data.

Replication: either several copies (Ch 6) or erasure coding (Reed–Solomon), which allows lost data to be recovered with LOWER STORAGE OVERHEAD than full replication. (Similar to RAID; the difference is that here file access and replication are done over a conventional datacenter network without special hardware.)

2.2 Object stores

The URL anatomy: s3://my-photo-bucket/2025/04/01/birthday.png Host = the BUCKET (globally unique name); the rest = the object's KEY (unique within its bucket).

The differences that bite:

Distributed filesystemObject store
MutabilityFiles are mutable; fopen/fseek file handlesOBJECTS ARE IMMUTABLE ONCE WRITTEN. To update, FULLY REWRITE with a put. (Azure Blob and S3 Express One Zone support appends; most others don't.) NO FILE HANDLE APIs
DirectoriesRealDO NOT EXIST. The path structure is SIMPLY A CONVENTION — the slashes are PART OF THE KEY
Listingls of one levelA prefix list behaves like a RECURSIVE ls -R — all objects starting with the prefix, including subpaths
Empty directoriesPossibleNOT POSSIBLE. Delete everything under .../2025/04/01 and 01 disappears from the listing of .../2025/04. Common practice: create a ZERO-BYTE OBJECT to represent an empty directory
Hard/symbolic links, file lockingOften supportedTypically NOT supported
RenamesAtomicNONATOMIC — copy to the new key, then delete the old. TO RENAME A "DIRECTORY" YOU MUST INDIVIDUALLY RENAME EVERY OBJECT WITHIN IT
Data localityHDFS allows tasks to RUN ON THE MACHINE STORING A COPY of the file — reading without sending it over the network, saving bandwidth IF THE TASK'S CODE IS SMALLER THAN THE FILEStorage and computation are SEPARATE. Might use more bandwidth, BUT MODERN DATACENTER NETWORKS ARE VERY FAST, so this is often acceptable — and it lets CPU/memory SCALE INDEPENDENTLY OF STORAGE

Size/latency positioning:

The key-value stores of Ch 4 are optimized for SMALL values (kilobytes) and FREQUENT, LOW-LATENCY reads/writes. Distributed filesystems and object stores are optimized for LARGE objects (megabytes to gigabytes) and LESS FREQUENT, LARGER reads. Recently, though, object stores have begun adding support for frequent, smaller I/O — S3 EXPRESS ONE ZONE now offers SINGLE-MILLISECOND LATENCY and a pricing model more similar to key-value stores.

⚠️ The line is blurry and dangerous: FUSE drivers let you treat S3 as a filesystem; JuiceFS and Ceph offer both APIs. However, their APIs, PERFORMANCE, AND CONSISTENCY GUARANTEES ARE VERY DIFFERENT. CARE MUST BE TAKEN to make sure they behave as expected, EVEN IF THEY SEEM TO IMPLEMENT THE REQUISITE APIs.

2.3 Distributed job orchestration

A job-start request carries: number of tasks · memory/CPU/disk per task · a job identifier · access credentials · job parameters (input and output data) · required hardware details such as GPUs or disk types · the location of the job's executable code.

Three components you'll find in nearly every orchestrator:

ComponentRoleYARNKubernetes
Task executorsDaemon on each node: runs tasks, sends HEARTBEATS to signal liveness, tracks task status and resource allocation. Retrieves the job's executable code, starts the task, monitors until it finishes or fails. Also works with the OS for SECURITY AND PERFORMANCE ISOLATION — both use Linux CGROUPS — preventing tasks from accessing data without permission or degrading other tasksNodeManagerkubelet
Resource managerMetadata about each node: available hardware, task statuses, network location, node status. Provides a GLOBAL VIEW of cluster state. ⚠️ Its CENTRALIZED nature CAN LEAD TO BOTH SCALABILITY AND AVAILABILITY BOTTLENECKSState in ZooKeeperState in etcd
SchedulerReceives start/stop/status requests; uses the request plus resource-manager state to decide WHICH TASKS RUN ON WHICH NODESResourceManagerkube-scheduler
(Application-specific sub-schedulers)For requirements the central scheduler can't know — e.g. autoscaling read replicas at a query threshold. They work together with the central schedulerApplicationMastersoperators
Resource allocation — the genuinely hard part

The five-node, 160-core example with two jobs each wanting 100 cores:

Option A
run 80 tasks for each job; start the remaining 20 each as tasks complete.
Option B: gang scheduling

run All of one job’s tasks, then start the second’s only when 100 cores are free.

  • ✗ “If the scheduler Reserves cores until all 100 are available at the same time, Nodes will sit idle. Cluster utilization Drops, and a deadlock might occur if other jobs also attempt to reserve cores.”
  • ✗ “If it simply Waits for 100 cores, Other jobs might grab them in the meantime. The cluster might not have 100 cores available for a very long time, which leads to Starvation.”
  • ✗ Preemption — kill some of job 1’s tasks to make room. “Decreases cluster efficiency as well, since the killed tasks Will need to be restarted later.”
Option C
the second request arrives much later ⇒ the scheduler has Incomplete information. Allocate all 100 to job 1, or Hold some back in anticipation of a future job That might or might not ever come?

Now imagine hundreds or even MILLIONS of such requests. Finding an optimal solution seems intractable. IN FACT, THE PROBLEM IS NP-HARD — prohibitively slow to solve optimally for all but the smallest examples.

In practice, schedulers therefore use HEURISTICS to make NONOPTIMAL BUT REASONABLE decisions: FIFO · dominant resource fairness (DRF) · priority queues · capacity/quota-based scheduling · bin-packing algorithms.

Scheduling workflows

A WORKFLOW (or DAG) of jobs: the output of one job becomes the input to one or more others.

⚠️ Terminology collision with Ch 5: "In 'Durable Execution and Workflows' we saw workflow engines offering durable execution of a sequence of steps, typically performing RPCs. In BATCH processing, 'workflow' has a DIFFERENT MEANING: a sequence of BATCH PROCESSES, each taking input data and producing output data, but NORMALLY NOT MAKING RPCs TO EXTERNAL SERVICES. Durable execution engines typically process LESS DATA PER REQUEST, though the line is somewhat fuzzy."

Three reasons a workflow is needed:

  1. The output feeds several jobs MAINTAINED BY DIFFERENT TEAMS → write it where all can read it, and schedule consumers on data update or their own schedule
  2. Transfer data between processing TOOLS — a Spark job writes HDFS, a Python script triggers a Trino SQL query, which outputs to S3
  3. Multiple internal stages — if one stage needs data sharded by one key and the next by a different key, the first stage can output data sharded the way the second requires

Coupling choice — pipe vs file:

Unix pipe modelFile / object store model — more typical
A small in-memory buffer; if it fills, the producer waits — a form of backpressure.The job writes output to a DFS or object store; the next job reads it from there.
Spark and Flink support a similar model: the output of one task is passed directly to another, over the network if they are on different machines.This decouples the jobs, allowing them to run at different times. A workflow scheduler waits until all jobs producing its inputs have completed successfully before running the consumer.

Orchestration-framework schedulers (YARN's ResourceManager, Spark's built-in scheduler) DO NOT MANAGE ENTIRE WORKFLOWS; they schedule PER JOB. To handle dependencies BETWEEN job executions, WORKFLOW SCHEDULERS were developed: AIRFLOW, DAGSTER, PREFECT.

Workflows of 50 TO 100 JOBS are common in many data pipelines, and in a large organization MANY TEAMS MAY BE RUNNING JOBS THAT READ ONE ANOTHER'S OUTPUT ACROSS MANY SYSTEMS. TOOL SUPPORT IS IMPORTANT FOR MANAGING SUCH COMPLEX DATAFLOWS.

Handling faults — and why batch has it easy

Two reasons a task doesn't finish: hardware faults / network interruptions (Ch 2, Ch 9), and deliberate PREEMPTION by the scheduler.

Preemption is particularly useful with MULTIPLE PRIORITY LEVELS: low-priority tasks are CHEAPER and run whenever there's spare capacity, but RISK BEING PREEMPTED AT ANY MOMENT. These are spot instances (EC2), spot virtual machines (Azure), preemptible instances (Google Cloud).

Batch processing is often not time-sensitive, so it's WELL SUITED to spot instances — using spare resources that would otherwise be idle, INCREASING CLUSTER UTILIZATION. However, THOSE TASKS ARE MORE LIKELY TO BE KILLED, BECAUSE PREEMPTIONS OCCUR MORE FREQUENTLY THAN HARDWARE FAULTS.

Since batch jobs REGENERATE THEIR OUTPUT FROM SCRATCH every time, TASK FAILURES ARE EASIER TO HANDLE THAN IN ONLINE SYSTEMS: delete the partial output from the failed execution and reschedule the task on another machine.

It would be WASTEFUL to rerun the ENTIRE job for one task failure. MapReduce and successors therefore KEEP THE EXECUTION OF PARALLEL TASKS INDEPENDENT, so they can RETRY AT THE GRANULARITY OF AN INDIVIDUAL TASK.

Three approaches to intermediate-data fault tolerance:

SystemApproachTrade-off
MapReduceALWAYS writes intermediate data back to the DFS, and WAITS for the writing task to complete successfully before others read itWorks even where preemption is common — BUT MEANS A LOT OF WRITES TO THE DFS, WHICH CAN BE INEFFICIENT
SparkKeeps intermediate data IN MEMORY (spilling to local disk if it won't fit) and writes ONLY THE FINAL RESULT to the DFS. TRACKS HOW THE INTERMEDIATE DATA WAS COMPUTED, allowing RECOMPUTATION if lost (lineage)Much faster; recomputation cost on loss
FlinkPeriodic CHECKPOINTING of a snapshot of tasksDifferent trade-off again

On this page