Skip to content

Batch vs Stream Processing#

Problem statement (interviewer prompt)

Design the analytics pipeline for a streaming product. You need real-time dashboards (sub-minute freshness) AND historical backfills over 5 years. Pick between Lambda and Kappa, choose Flink vs Spark vs Kafka Streams, and discuss watermarks + exactly-once.

MapReduce data flow: input split into map tasks, shuffled, then reduced, the canonical batch processing model
Source: Wikimedia Commons. CC BY-SA.
flowchart LR
  E[Events]
  B[Batch: Spark / Airflow<br/>periodic jobs]
  S[Stream: Flink / Kafka Streams<br/>continuous]
  DW[(Data warehouse)]
  RT[(Real-time view)]
  E --> B --> DW
  E --> S --> RT

  classDef p fill:#dbeafe,stroke:#1e40af,stroke-width:1px,color:#0f172a;
  classDef s fill:#fef3c7,stroke:#92400e,stroke-width:1px,color:#0f172a;
  classDef r fill:#fee2e2,stroke:#991b1b,stroke-width:1px,color:#0f172a;
  class E p;
  class B,S s;
  class DW,RT r;

    classDef client fill:#dbeafe,stroke:#1e40af,stroke-width:1px,color:#0f172a;
    classDef edge fill:#cffafe,stroke:#0e7490,stroke-width:1px,color:#0f172a;
    classDef service fill:#fef3c7,stroke:#92400e,stroke-width:1px,color:#0f172a;
    classDef datastore fill:#fee2e2,stroke:#991b1b,stroke-width:1px,color:#0f172a;
    classDef cache fill:#fed7aa,stroke:#9a3412,stroke-width:1px,color:#0f172a;
    classDef queue fill:#ede9fe,stroke:#5b21b6,stroke-width:1px,color:#0f172a;
    classDef compute fill:#d1fae5,stroke:#065f46,stroke-width:1px,color:#0f172a;
    classDef storage fill:#e5e7eb,stroke:#374151,stroke-width:1px,color:#0f172a;
    classDef external fill:#fce7f3,stroke:#9d174d,stroke-width:1px,color:#0f172a;
    classDef obs fill:#f3e8ff,stroke:#6b21a8,stroke-width:1px,color:#0f172a;
    class E service;
    class DW,RT datastore;
    class S queue;
    class B compute;

Batch processes bounded chunks of data periodically (cheap, easy). Stream processes events one-at-a-time, continuously (fresh, harder). Modern systems use both - historically as Lambda architecture, increasingly as Kappa.

flowchart TB
  subgraph Sources
    APP[App events]
    DB[DB CDC]
    LOGS[Logs]
  end

  subgraph Bus[Streaming bus]
    K[Kafka / Kinesis / Pub-Sub]
  end

  subgraph Lake[Storage]
    DL[(Data lake - S3 / Parquet)]
    DWH[(Warehouse - Snowflake / BigQuery / Redshift)]
  end

  subgraph Batch[Batch processing]
    BJ[Airflow / Spark / dbt]
    SCHED[Cron-style schedule]
  end

  subgraph Stream[Stream processing]
    FL[Flink / Kafka Streams / Spark Structured Streaming]
    STATE[(State store - RocksDB)]
    WIN[Tumbling / hopping / session windows]
    WM[Watermarks for late events]
  end

  Sources --> K
  K --> Stream
  K --> DL
  DL --> Batch
  Batch --> DWH
  Stream --> DWH
  Stream --> RT[(Real-time serving store)]

  classDef compute fill:#d1fae5,stroke:#065f46,stroke-width:1px,color:#0f172a;
  class BJ,FL,WIN,WM,SCHED,STATE compute;

    classDef client fill:#dbeafe,stroke:#1e40af,stroke-width:1px,color:#0f172a;
    classDef edge fill:#cffafe,stroke:#0e7490,stroke-width:1px,color:#0f172a;
    classDef service fill:#fef3c7,stroke:#92400e,stroke-width:1px,color:#0f172a;
    classDef datastore fill:#fee2e2,stroke:#991b1b,stroke-width:1px,color:#0f172a;
    classDef cache fill:#fed7aa,stroke:#9a3412,stroke-width:1px,color:#0f172a;
    classDef queue fill:#ede9fe,stroke:#5b21b6,stroke-width:1px,color:#0f172a;
    classDef compute fill:#d1fae5,stroke:#065f46,stroke-width:1px,color:#0f172a;
    classDef storage fill:#e5e7eb,stroke:#374151,stroke-width:1px,color:#0f172a;
    classDef external fill:#fce7f3,stroke:#9d174d,stroke-width:1px,color:#0f172a;
    classDef obs fill:#f3e8ff,stroke:#6b21a8,stroke-width:1px,color:#0f172a;
    class APP,WIN,WM service;
    class DB,DWH,STATE,RT datastore;
    class K,FL queue;
    class BJ,SCHED compute;
    class DL storage;
    class LOGS obs;

When each wins#

Need Pick
Daily / hourly reports Batch
Heavy backfills, joins on TB Batch
Real-time dashboards Stream
Real-time ML features Stream
Anomaly / fraud detection Stream
Strict reproducibility Batch (deterministic re-run)
Latency-sensitive personalization Stream

Lambda architecture#

flowchart LR
  Ev[Events] --> Bf[Batch: full re-compute]
  Ev --> Sf[Stream: incremental]
  Bf --> Serve[Serving layer]
  Sf --> Serve
  Serve --> User

    classDef client fill:#dbeafe,stroke:#1e40af,stroke-width:1px,color:#0f172a;
    classDef edge fill:#cffafe,stroke:#0e7490,stroke-width:1px,color:#0f172a;
    classDef service fill:#fef3c7,stroke:#92400e,stroke-width:1px,color:#0f172a;
    classDef datastore fill:#fee2e2,stroke:#991b1b,stroke-width:1px,color:#0f172a;
    classDef cache fill:#fed7aa,stroke:#9a3412,stroke-width:1px,color:#0f172a;
    classDef queue fill:#ede9fe,stroke:#5b21b6,stroke-width:1px,color:#0f172a;
    classDef compute fill:#d1fae5,stroke:#065f46,stroke-width:1px,color:#0f172a;
    classDef storage fill:#e5e7eb,stroke:#374151,stroke-width:1px,color:#0f172a;
    classDef external fill:#fce7f3,stroke:#9d174d,stroke-width:1px,color:#0f172a;
    classDef obs fill:#f3e8ff,stroke:#6b21a8,stroke-width:1px,color:#0f172a;
    class Ev,Bf,Serve service;
    class Sf queue;

Two code paths, two storage paths - heavy maintenance.

Kappa architecture#

flowchart LR
  Ev[Events forever-log] --> Sf[Stream: one path]
  Sf --> Serve[Serving layer]
  Serve --> User
  Ev -.replay.-> Sf

    classDef client fill:#dbeafe,stroke:#1e40af,stroke-width:1px,color:#0f172a;
    classDef edge fill:#cffafe,stroke:#0e7490,stroke-width:1px,color:#0f172a;
    classDef service fill:#fef3c7,stroke:#92400e,stroke-width:1px,color:#0f172a;
    classDef datastore fill:#fee2e2,stroke:#991b1b,stroke-width:1px,color:#0f172a;
    classDef cache fill:#fed7aa,stroke:#9a3412,stroke-width:1px,color:#0f172a;
    classDef queue fill:#ede9fe,stroke:#5b21b6,stroke-width:1px,color:#0f172a;
    classDef compute fill:#d1fae5,stroke:#065f46,stroke-width:1px,color:#0f172a;
    classDef storage fill:#e5e7eb,stroke:#374151,stroke-width:1px,color:#0f172a;
    classDef external fill:#fce7f3,stroke:#9d174d,stroke-width:1px,color:#0f172a;
    classDef obs fill:#f3e8ff,stroke:#6b21a8,stroke-width:1px,color:#0f172a;
    class Ev,Serve service;
    class Sf queue;

One code path. Reprocess history by rewinding the stream.

Stream-processing fundamentals#

Time#

  • Event time: when the event happened (clock on the device).
  • Ingestion time: when it entered the system.
  • Processing time: when the operator processed it.
  • Watermark: "we've seen all events with event-time ≤ T" - triggers window close.

Windows#

  • Tumbling - fixed, non-overlapping (every 1 min).
  • Hopping / sliding - fixed size, overlap (5 min window, slide 1 min).
  • Session - gap-defined (close after 30 min idle).

State#

  • Stream operators maintain state (counts, joins, sessions).
  • Backed by an embedded store (Flink + RocksDB).
  • Checkpointed to durable storage for recovery (every 1-30 s).

Exactly-once#

  • Source: replayable, dedupable offsets (Kafka).
  • Operator: deterministic + idempotent + transactional sinks (Flink + Kafka tx).
  • Sink: upsert by key.

OLAP store choices#

Store Best for Trade-off
Snowflake / BigQuery / Redshift analyst SQL, federated per-query cost
ClickHouse / Druid / Pinot sub-second OLAP at scale ops complexity
Iceberg / Delta / Hudi on S3 data-lake table format governance + ETL friction

ETL vs ELT#

  • ETL (transform before load) - older; precomputed shape.
  • ELT (transform inside warehouse with dbt) - modern; raw data preserved.

Common pitfalls#

  • Cardinality blow-up in stream aggregations (group by user_id × second).
  • Watermark stuck because one slow source - partition the bus, run shadow watermarks.
  • Reprocessing in Lambda - code drift between batch and stream branches.
  • Backfills: stream platforms struggle with multi-year historical data; bring batch back for that.

Glossary & fundamentals#

Concepts referenced in this design. Each row links to its canonical page; the tag column shows whether it is a high-level (HLD) or low-level (LLD) concept.

Tag Concept What it is Page
HLD Pub/Sub & message brokers topics, consumer groups, delivery semantics pub-sub-pattern
HLD LSM vs B-Tree engines WAL, memtable, SSTables, compaction storage-engines-lsm-btree
HLD Change Data Capture WAL/binlog tailing, outbox publishing change-data-capture
HLD Batch & stream processing Lambda vs Kappa, watermarks, windows batch-stream-processing

Quick reference#

Choosing freshness#

  • ≤ 100 ms - in-process (Kafka Streams, Flink low-latency mode).
  • ≤ 5 s - mainstream stream processing.
  • ≤ 5 min - micro-batches (Spark Structured Streaming with 30 s triggers).
  • ≥ 15 min - batch jobs.

Schema management#

  • Use Avro / Protobuf with a Schema Registry.
  • Forward + backward compat rules: add fields with defaults, never remove.
  • Compact topics for "latest value per key" semantics.

Cost levers#

  • Compact / dedupe early to shrink downstream.
  • Tier-out hot vs cold storage in the warehouse.
  • Time-travel features (Snowflake / Iceberg) - fast but pricey.

Operational checklist#

  • Consumer lag monitor (Burrow / Kafka exporter).
  • Watermark monitor for stream jobs.
  • Backfill SOP: pause downstream, replay, validate, resume.
  • Schema-evolution playbook.

Refs#

  • "Streaming Systems" - Akidau, Chernyak, Lax (Google).
  • Apache Flink docs, Kafka Streams DSL.
  • Jay Kreps: "Questioning the Lambda Architecture" (origin of Kappa).
  • "Designing Data-Intensive Applications" - ch.10-11.

FAQ#

What is the difference between batch and stream processing?#

Batch processing runs jobs over bounded datasets on a schedule, optimized for throughput. Stream processing handles unbounded event streams continuously, optimized for low latency and freshness.

What is Lambda architecture?#

Lambda combines a batch layer for accurate historical results with a speed layer for near-real-time approximations. A serving layer merges both views. It is reliable but maintains two codebases.

What is Kappa architecture?#

Kappa removes the batch layer and uses a single stream-processing pipeline for both real-time and backfill, replaying the event log from offset zero when reprocessing is needed.

When should I pick streaming over batch?#

Pick streaming when business value drops with latency, like fraud detection, dashboards, or alerting. Stick with batch when freshness above minutes does not change the decision.

What is a watermark in stream processing?#

A watermark is a heuristic timestamp that says no events older than this should arrive. It lets the engine close event-time windows and emit results despite out-of-order events.

  • Event Sourcing and CQRS: event sourcing produces the event streams that stream processing pipelines consume
  • Pub/Sub Pattern: pub/sub messaging systems are the backbone of real-time stream processing architectures

Further reading#

Curated, high-credibility sources for going deeper on this topic.