Batch Processing
Batch processing is the execution of data jobs on a scheduled, finite dataset. Despite the rise of streaming, batch remains the dominant paradigm for data warehousing, reporting, and large-scale transformations where real-time is not required.
Core Concepts
Job Scheduling
Batch jobs run on a cadence: hourly, daily, weekly. The scheduler manages dependencies, retries, and resource allocation. Common schedulers include:
- Apache Airflow - DAG-based scheduling with rich dependency management
- Prefect - Python-native, with automatic retries and notifications
- Dagster - Asset-based, with built-in data quality checks
- Cron + bash - Simple scheduled scripts (fine for prototypes)
Partitioning
Partitioning splits data into manageable chunks, typically by date. A daily batch job processes one partition per run, reducing the dataset from terabytes to gigabytes.
# Date-partitioned batch job structure
/data
/landing/events/2025/01/15/events_000.parquet
/landing/events/2025/01/15/events_001.parquet
/staging/events/dt=2025-01-15/
/prod/events/dt=2025-01-15/
# Spark reads only relevant partition
df = spark.read.parquet(f"/data/landing/events/{year}/{month}/{day}/")MapReduce Paradigm
MapReduce is the foundational batch processing model, splitting work into two phases:
- Map - Apply a function to each record, emitting key-value pairs
- Reduce - Aggregate values by key
# Conceptual MapReduce: word count def mapper(line): for word in line.split(): emit(word, 1) def reducer(word, counts): emit(word, sum(counts))Apache Spark for Batch
Spark is the most widely used batch processing engine, offering in-memory computation that is 10-100x faster than Hadoop MapReduce for iterative algorithms.
Spark Batch Example
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum, avg, count
spark = SparkSession.builder \
.appName("daily_aggregation") \
.config("spark.sql.shuffle.partitions", "200") \
.getOrCreate()
# Read daily partition
orders = spark.read.parquet(f"/data/orders/dt={execution_date}/")
# Transform
daily_metrics = (
orders
.groupBy("product_id", "category")
.agg(
count("order_id").alias("order_count"),
sum("amount").alias("total_revenue"),
avg("amount").alias("avg_order_value"),
)
.filter(col("order_count") > 10)
)
# Write result
daily_metrics.write.mode("overwrite") \
.parquet(f"/data/metrics/dt={execution_date}/")Optimization Techniques
- Partition pruning - Only read relevant partitions based on filter predicates
- Bucketing - Pre-shard data by join keys to avoid expensive shuffles
- Predicate pushdown - Push filters into the data source (Parquet, Delta) to skip irrelevant rows
- Broadcast joins - Distribute small tables to all executors instead of shuffling
- Caching - Keep intermediate results in memory when reused across transformations
Monitoring Batch Jobs
- Data volume metrics (rows in/out, bytes processed)
- Duration and resource utilization (CPU, memory, shuffle I/O)
- Failure rate and retry counts
- Data quality constraints (freshness, completeness, uniqueness)
- Cost per run (especially on cloud: Databricks DBUs, Snowflake credits)