Key Differences Between Batch and Streaming Processing: Architectural Patterns and Implementation Guide
Batch processing handles finite datasets on scheduled intervals with high latency tolerance, while streaming processing consumes unbounded data continuously with sub-second latency requirements and complex state management.
Understanding the key differences between batch and streaming processing is fundamental to building efficient data pipelines. The DataExpert-io/data-engineer-handbook repository demonstrates both paradigms through practical implementations using Apache Spark and Apache Flink. This article explores the architectural distinctions, operational trade-offs, and hands-on code examples found in the handbook's intermediate bootcamp materials.
Key Differences in Execution Model and Data Boundaries
The fundamental distinction lies in how each paradigm handles data volume and execution lifecycle.
Batch Processing: Finite Datasets
Batch processing operates on a finite set of records collected over a specific time window, such as a day or hour. According to intermediate-bootcamp/materials/3-spark-fundamentals/README.md, batch jobs follow a start-terminate pattern: the job initializes, reads the entire dataset snapshot, performs transformations, writes results, and then releases resources. This model relies on schedulers like Airflow, Dagster, or cron to trigger execution at predetermined intervals.
Streaming Processing: Unbounded Streams
Streaming processing handles unbounded, continuously arriving data streams that never terminate. As implemented in intermediate-bootcamp/materials/4-apache-flink-training/README.md, streaming jobs run indefinitely, ingesting events as they arrive from sources like Kafka or Kinesis. The pipeline maintains persistent execution contexts to handle real-time ingestion, requiring continuous resource allocation rather than ephemeral job clusters.
Latency and Timing Characteristics
Latency requirements represent a critical differentiator between these paradigms.
Batch processing accepts high latency ranging from minutes to hours because results are only needed after the collection window closes. The system optimizes for throughput over speed, processing terabytes of data in a single pass without immediate delivery constraints.
Streaming processing demands low latency measured in seconds to sub-seconds. Real-time dashboards, fraud detection systems, and event-driven microservices require immediate insights. The Flink training material in intermediate-bootcamp/materials/4-apache-flink-training/homework/homework.md demonstrates how frameworks handle out-of-order events using watermarks and event-time processing to maintain correctness while minimizing latency.
State Management and Fault Tolerance
These key differences between batch and streaming processing manifest most clearly in how each system handles state and recovery.
Batch Recovery Mechanisms
Batch jobs are typically stateless or use intermediate temporary storage that clears after completion. Fault tolerance relies on idempotent writes and simple re-execution: if a job fails, operators restart the entire pipeline from the beginning. The Spark fundamentals section in intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md illustrates this pattern through PySpark jobs that read static Parquet files and overwrite destination directories atomically.
Streaming State and Checkpoints
Streaming requires persistent state management that survives across individual events. State backends like RocksDB or HDFS store keyed windows and aggregates continuously. According to the Flink implementation, systems use checkpointing and replay mechanisms—such as Kafka offsets or Flink snapshots—to recover only the lost portion of the stream without duplicating already-processed records. This enables exactly-once processing semantics while maintaining back-pressure to prevent system overload.
Scalability and Resource Patterns
Both paradigms scale horizontally but handle growth differently.
Batch processing scales by partitioning the input dataset and running parallel tasks for a single job execution. Resources are provisioned for the expected runtime and released upon completion, making cost optimization straightforward.
Streaming processing scales both horizontally through parallel operators and vertically to handle ever-growing event rates. The system must maintain ordering guarantees per key while dynamically adjusting to traffic spikes, requiring sophisticated load balancing as shown in the Spark Structured Streaming examples.
Practical Implementation Examples
The Data Engineer Handbook provides concrete implementations demonstrating these architectural differences.
Batch Processing with PySpark
The following example from intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md demonstrates a classic batch job that processes daily sales data:
# File: intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md
# Reads a static Parquet dataset, performs a daily aggregation, and writes the result.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as _sum, to_date
spark = SparkSession.builder.appName("DailySalesBatch").getOrCreate()
# Input: a static dataset stored on S3 (or local for the demo)
sales_df = spark.read.parquet("s3://my-bucket/sales/")
# Daily aggregation (batch‑style)
daily_sales = (
sales_df
.withColumn("date", to_date(col("event_timestamp")))
.groupBy("date")
.agg(_sum("amount").alias("total_amount"))
)
daily_sales.write.mode("overwrite").parquet("s3://my-bucket/aggregated/daily_sales/")
spark.stop()
This job exhibits classic batch characteristics: it reads a static snapshot, performs a complete aggregation, writes the output, and terminates.
Stateful Stream Processing with Apache Flink
The Flink training material in intermediate-bootcamp/materials/4-apache-flink-training/homework/homework.md implements continuous sessionization:
# File: intermediate-bootcamp/materials/4-apache-flink-training/homework/homework.md
# A simple Flink streaming job that reads from Kafka, sessionizes by IP, and writes to a sink.
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes, EnvironmentSettings
from pyflink.table.descriptors import Kafka, Json, Schema
env = StreamExecutionEnvironment.get_execution_environment()
settings = EnvironmentSettings.new_instance().in_streaming_mode().use_blink_planner().build()
t_env = StreamTableEnvironment.create(env, environment_settings=settings)
# Define source – Kafka topic "events"
t_env.connect(
Kafka()
.version("universal")
.topic("events")
.start_from_latest()
.property("bootstrap.servers", "kafka:9092")
).with_format(
Json().fail_on_missing_field(True)
).with_schema(
Schema()
.field("ip", DataTypes.STRING())
.field("timestamp", DataTypes.TIMESTAMP(3))
.field("payload", DataTypes.STRING())
).create_temporary_table("KafkaSource")
# Sessionize by IP address (streaming window)
t_env.sql_query("""
SELECT
ip,
SESSION_START(ts) as session_start,
SESSION_END(ts) as session_end,
COUNT(*) as event_count
FROM KafkaSource
WINDOW SESSION(10 MINUTES) AS w
GROUP BY ip, w
""").execute_insert("ProcessedEventsSink")
This pipeline runs continuously, maintaining session state across events and using Kafka offsets for fault tolerance.
Micro-Batch Processing with Spark Structured Streaming
The handbook also covers hybrid approaches in intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md:
# File: intermediate-bootcamp/materials/3-spark-fundamentals/homework/homework.md
# Reads a continuous stream from Kafka, aggregates per 5‑second window, and writes to console.
from pyspark.sql import SparkSession
from pyspark.sql.functions import window, col, sum as _sum
spark = SparkSession.builder.appName("StreamingSales").getOrCreate()
kafka_df = (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafka:9092")
.option("subscribe", "sales")
.load()
)
sales_df = kafka_df.selectExpr("CAST(value AS STRING) as json")
sales_parsed = spark.read.json(sales_df.rdd.map(lambda r: r.json))
agg_df = (
sales_parsed
.withColumn("event_time", col("event_timestamp").cast("timestamp"))
.groupBy(window("event_time", "5 seconds"))
.agg(_sum("amount").alias("total_amount"))
)
query = agg_df.writeStream \
.outputMode("complete") \
.format("console") \
.option("truncate", "false") \
.start()
query.awaitTermination()
This example uses Spark Structured Streaming to process Kafka topics in micro-batches, bridging the gap between pure batch and continuous streaming.
Summary
- Batch processing handles finite datasets with scheduled execution, high latency tolerance, and simple fault tolerance through job re-execution.
- Streaming processing manages unbounded data with continuous execution, sub-second latency requirements, and complex state management using checkpointing.
- State persistence distinguishes the two paradigms: batch uses temporary or no state, while streaming requires durable state stores like RocksDB.
- Recovery mechanisms differ significantly: batch restarts entire jobs, while streaming uses offset replay and snapshots for granular recovery.
- The DataExpert-io/data-engineer-handbook provides production-ready examples of both paradigms in
intermediate-bootcamp/materials/3-spark-fundamentals/andintermediate-bootcamp/materials/4-apache-flink-training/.
Frequently Asked Questions
When should I choose batch processing over streaming?
Choose batch processing when your business requirements tolerate latency of minutes to hours and you need comprehensive analytics on complete datasets. Daily ETL jobs, data warehouse loads, and large-scale historical analysis are ideal batch use cases because they simplify development and resource management compared to continuous pipelines.
Can batch systems handle real-time data requirements?
Traditional batch systems cannot meet true real-time requirements, but micro-batch architectures like Spark Structured Streaming provide a hybrid approach. As shown in the handbook's Spark fundamentals section, micro-batching processes small increments of data every few seconds, offering near-real-time latency while maintaining simpler fault tolerance semantics than pure streaming.
What makes streaming processing more complex than batch?
Streaming introduces complexity around state management, event ordering, and fault tolerance. Unlike batch jobs that process static snapshots, streaming systems must handle out-of-order events using watermarks, maintain persistent aggregates across infinite data streams, and implement checkpointing mechanisms to recover from failures without data loss or duplication.
How does the Data Engineer Handbook demonstrate these concepts?
The repository provides hands-on implementations in intermediate-bootcamp/materials/3-spark-fundamentals/README.md for batch and micro-batch processing using PySpark, and intermediate-bootcamp/materials/4-apache-flink-training/README.md for stateful stream processing with Apache Flink. The accompanying homework files contain runnable code examples for daily aggregation jobs, Kafka sessionization pipelines, and windowed stream analytics.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →