DE Wikiguides / batch

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)

Resources

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)

Resources