# Batch vs Stream Processing: MapReduce, Spark, Flink, and Kafka Streams

## Blog Details

- **Author**: Naveen R.
- **Date**: October 7, 2026
- **Tags**: stream processing, batch processing, distributed systems, data engineering, apache flink
- **Read Time**: 20 mins

## Introduction

Any system design interview about analytics, metrics, fraud, or billing eventually reaches the same fork: "Do you compute this in batch or in a stream?" Weak answers name a tool and stop. Strong answers start from two questions: how fresh must the result be, and how correct must it be when data arrives late, twice, or out of order?

This post covers batch (MapReduce, Spark), streaming (Flink, Kafka Streams), the semantics that make streaming hard (event time, watermarks, exactly-once state), Lambda vs Kappa, and micro-batch vs true streaming.

We use one running example throughout: an **ad-click pipeline**. Ad servers emit a click event for every click: `click_id`, `campaign_id`, `advertiser_id`, `user_id`, `event_time`, and some fraud signals. The business needs two outputs:

1. **Per-minute click counts per campaign**, for live dashboards and budget pacing (stop serving a campaign when it exhausts its daily budget).
2. **Daily aggregates per advertiser**, for invoicing. These must be exact, because they become money.

Sizing: assume an average of 20,000 clicks per second, peaks several times higher, and about 500 bytes per event. That is 20,000 × 86,400 ≈ 1.7 billion events per day, roughly 860 GB of raw data per day before compression: large enough that skew, shuffles, and state size all matter.

![An overview of the ad-click pipeline: click events land in a Kafka topic, which is archived to a Parquet data lake for a daily Spark batch job that writes billing aggregates to a warehouse, and is also consumed by a Flink job that writes live per-minute counts to a serving store, with both feeding dashboards and billing.](https://d5osvdbc8um23.cloudfront.net/static-asset/blog_images/batch-vs-stream-processing-compared/01-high-level-architecture.png)

## Bounded vs Unbounded Data

The cleaner framing is **bounded vs unbounded data**, the vocabulary of Google's Dataflow model paper (Akidau et al., 2015).

- **Bounded data** has a known end. "All clicks on 2026-10-04" is bounded once no more clicks for that day will arrive, so you can read it all and produce a final answer.
- **Unbounded data** never ends. Any answer is provisional, and some new data belongs to periods you already reported on.

A **batch engine** assumes it sees the whole input, so it can plan globally and retry tasks from scratch. A **stream engine** keeps long-lived operators running, holds state across events, and must decide *when* an answer is complete enough to emit.

Two consequences matter in interviews:

1. **Batch is a special case of streaming.** A bounded dataset is a stream that ends. Flink runs bounded jobs in batch execution mode on the same runtime, and Spark runs the same DataFrame query over a table or a stream.
2. **Streaming is not inherently less correct.** "Streams are approximate, batch is exact" dates from early stream processors without durable state or event time. Modern engines can be as correct as batch; what they cannot be is correct *and* final at the same instant if data can arrive arbitrarily late. Latency vs completeness vs cost is the real trade-off.

In our pipeline, the dashboard is unbounded and approximate-then-corrected is fine; the invoice is bounded, so we wait a few hours after midnight and compute once.

## Batch Processing: MapReduce and Spark

### When to Reach for Batch

Batch fits when results can be hours old, when the computation needs the full dataset (global sorts, large joins, deduplication across a day, model training), and when cost beats latency. Batch jobs are also the simplest to rerun: fix the bug and recompute the partition. That **cheap, safe reprocessing** is why billing and ML feature backfills tend to live in batch even where streaming is mature.

### MapReduce: Map, Shuffle, Reduce

MapReduce (the 2004 Google paper, popularized by Hadoop) reduced distributed batch processing to two user functions and one framework step:

1. **Map**: each mapper reads one input split (typically one HDFS block, 128 MB by default in Hadoop 2 and later) and emits key-value pairs: `(campaign_id, 1)`.
2. **Shuffle and sort**: the framework partitions mapper output by hash of the key, writes it to local disk, and each reducer pulls and merge-sorts its partition from every mapper. This is the expensive, network-heavy step.
3. **Reduce**: each reducer sees all values for a key together and emits `(campaign_id, total)`.

A **combiner** pre-aggregates on the mapper, so each mapper ships a few partial counts per campaign (Hadoop runs the combiner per spill, zero or more times) instead of one record per click.

```python
# Conceptual MapReduce for daily clicks per campaign
def map_fn(click):
    if click["is_valid"]:
        yield (click["campaign_id"], 1)

def combine_fn(campaign_id, partial_counts):   # runs on the mapper
    yield (campaign_id, sum(partial_counts))

def reduce_fn(campaign_id, counts):            # runs after shuffle
    yield (campaign_id, sum(counts))
```

Fault tolerance is simple: tasks are deterministic functions of their split, so a failed task is rerun elsewhere. The weakness is that every job writes its output to the distributed file system, so multi-step pipelines become chains of jobs that pay full disk and replication cost between steps.

### Spark: DAGs, Stages, and Lazy Evaluation

Spark keeps MapReduce's partitioned, deterministic tasks but generalizes them to a **directed acyclic graph** of operations, keeping intermediate data in memory or on local disk instead of the distributed file system. Key concepts:

- **Partitions**: a DataFrame or RDD is split into partitions, and each partition is processed by one **task**. When reading files, Spark SQL packs input into partitions of up to `spark.sql.files.maxPartitionBytes` (128 MB by default). Too few partitions idle the cores; too many add scheduling overhead.
- **Lazy evaluation**: transformations (`filter`, `select`, `groupBy`, `join`) only build a logical plan. Nothing runs until an **action** (`count`, `collect`, `write`). Laziness lets the Catalyst optimizer see the whole query, push filters into the Parquet reader, and prune columns before any data moves.
- **Narrow vs wide dependencies**: narrow transformations (`filter`, `select`) compute each output partition from one input partition and are pipelined in one task. Wide ones (`groupBy`, most joins, `distinct`) need data from many partitions, which requires a **shuffle**.
- **Stages split at shuffles**: the DAG scheduler cuts the graph at every shuffle boundary. Each **stage** is a run of narrow transformations executed as one set of tasks; it writes shuffle files that the next stage reads.

![A Spark job for daily campaign clicks: the driver's DAG scheduler schedules stage 1, which scans one day of Parquet clicks, filters, and computes partial counts, then writes shuffle files partitioned by hash of campaign ID; stage 2 reads those shuffle files, joins with a broadcast campaigns table, computes final counts, and writes to the warehouse.](https://d5osvdbc8um23.cloudfront.net/static-asset/blog_images/batch-vs-stream-processing-compared/02-spark-job-stages.png)

Here is the daily advertiser rollup in PySpark:

```python
from pyspark.sql import functions as F

spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

clicks = spark.read.parquet("s3://ads-lake/clicks/")            # partitioned by dt
campaigns = spark.read.table("dim.campaigns")                   # small dimension

daily = (
    clicks
    .where((F.col("dt") == "2026-10-04") & F.col("is_valid"))   # partition pruning
    .dropDuplicates(["click_id"])                               # wide: shuffle by click_id
    .groupBy("dt", "campaign_id")                               # wide: shuffle by campaign_id
    .agg(F.count("*").alias("clicks"),
         F.approx_count_distinct("user_id").alias("approx_users"))
    .join(F.broadcast(campaigns), "campaign_id")                # broadcast: no shuffle of clicks
    .groupBy("dt", "advertiser_id")
    .agg(F.sum("clicks").alias("billable_clicks"))
)

(daily.write.mode("overwrite")                                  # replaces only dt=2026-10-04
      .partitionBy("dt")
      .parquet("s3://ads-warehouse/advertiser_daily/"))
```

Nothing executes until `write`. Spark then plans a scan-and-filter stage, a shuffle for deduplication, a shuffle for the campaign aggregate (with map-side partial aggregation, Spark's combiner), a broadcast join that ships the small campaigns table to every executor instead of shuffling the per-campaign aggregate rows, and a final small shuffle.

Two details make the job safe to rerun: deduplicating on `click_id` neutralizes duplicate deliveries upstream, and dynamic partition overwrite replaces exactly one day's output. Idempotent, partition-scoped output is the foundation of every backfill strategy.

### Scaling Batch Jobs

Batch scales by adding executors and partitions.

- `spark.sql.shuffle.partitions` defaults to **200**, far too few if you shuffle 860 GB of raw rows; aim for shuffle partitions in the low hundreds of megabytes. **Adaptive Query Execution** (AQE), enabled by default since Spark 3.2, coalesces small shuffle partitions at runtime and can switch a sort-merge join to a broadcast join when one side turns out small.
- **Columnar formats** (Parquet, ORC) plus partition pruning mean the job reads only the columns and dates it needs. Heavy spill to disk usually signals too few partitions or skew.

### Failure Modes and Anti-Patterns

1. **Data skew**: one key (a viral campaign, or a `null` campaign ID from a bug) lands in one shuffle partition, so one task runs for an hour while the rest finish in a minute. Fixes: handle the hot key separately, **salt** the key (append a random suffix 0 to N, aggregate, then aggregate again without the salt), or rely on AQE's skew join handling (`spark.sql.adaptive.skewJoin.enabled`), which splits oversized partitions in sort-merge joins.
2. **Small files**: many partitions each writing tiny files slows downstream reads. Coalesce before writing, or use a table format with compaction (Delta Lake, Iceberg, Hudi).
3. **Non-deterministic tasks**: `rand()` or the current time inside a transformation makes a retried task differ from the original. Seed randomness and pass "now" in as a parameter.
4. **Processing-time partitioning**: writing clicks to `dt=<arrival date>` instead of `dt=<event date>` means late events silently land in the wrong day's partition. Partition by event time and decide explicitly how long to wait before a day is final.

## Stream Processing: Flink and Kafka Streams

### When to Reach for Streaming

Streaming fits when a result's value decays quickly: budget pacing, fraud detection, alerting, live dashboards, fresh recommendation features. A campaign that overspends because its count was an hour stale is real money lost. It is the wrong default when nobody acts on fresher data or the computation needs a global view. A stream job is a 24/7 service with on-call, upgrades, and state migrations, not a cron entry.

### Apache Flink: Operators, Keyed State, Checkpoints

A Flink job is a dataflow graph of **operators** (sources, `map`, `keyBy`, windows, sinks), each running as parallel **subtasks** on TaskManagers, coordinated by a JobManager. Records flow continuously; there are no batch boundaries.

**Keyed state** is what makes Flink powerful. After `keyBy(campaign_id)`, every record for a campaign goes to the same subtask, and that subtask can keep per-key state (`ValueState`, `ListState`, `MapState`, window contents) that Flink manages for you. State is sharded into **key groups**, the unit Flink redistributes when you change parallelism. The maximum parallelism setting fixes the number of key groups, so set it generously up front; changing it later breaks restoring from existing state.

**State backends** decide where that state lives:

- **HashMapStateBackend** keeps state as objects on the JVM heap. Fast, but limited by heap size and garbage collection.
- **EmbeddedRocksDBStateBackend** keeps state in an embedded RocksDB instance on local disk, with hot data in memory. State can be much larger than memory, at the cost of serialization on every access. It supports **incremental checkpoints**, which upload only new SST files rather than the full state each time.

(These names date from Flink 1.13; earlier versions used `MemoryStateBackend`, `FsStateBackend`, and `RocksDBStateBackend`, which mixed where state lives with where checkpoints go.)

**Checkpoints** make that state fault tolerant. Flink's algorithm is asynchronous barrier snapshotting, a variant of the **Chandy-Lamport** distributed snapshot algorithm:

1. The JobManager tells each source to inject a **checkpoint barrier** with checkpoint ID *n* into its output streams, and the source records its current position (for Kafka, the partition offsets).
2. Barriers flow downstream with the data. When an operator has received barrier *n* on all of its inputs, it snapshots its state (asynchronously, for RocksDB) and forwards the barrier.
3. When every task has acknowledged its snapshot for barrier *n* to the JobManager, checkpoint *n* is complete: a consistent cut where the state reflects exactly the records before the recorded source offsets.

On failure, Flink restarts the job, restores every operator's state from the last completed checkpoint, and rewinds sources to the stored offsets. Records after the checkpoint are reprocessed, but because the state was rolled back too, each record affects the state exactly once.

With **aligned** checkpoints, an operator pauses inputs that already delivered the barrier until the slower ones catch up, which can stall under backpressure. **Unaligned checkpoints** (Flink 1.11 and later) let barriers overtake buffered records and store in-flight data in the checkpoint.

**Savepoints** use the same mechanism but are triggered manually and owned by the user. You take one to upgrade code, change parallelism, or migrate clusters, then start the new version from it. Stable operator IDs (`.uid("per-minute-window")`) let state map to the right operator after code changes.

### Kafka Streams: A Library, Not a Cluster

Kafka Streams is a **Java library** embedded in your application. There is no processing cluster: you run more instances of your service, and Kafka's consumer group protocol spreads the work.

- **Tasks map to input partitions.** A 24-partition topic yields 24 tasks, so parallelism tops out at the partition count.
- **State stores** hold local state, RocksDB by default, on each instance's disk.
- **Changelog topics** make that state durable. Every store update is also written to a compacted Kafka topic; if an instance dies, another rebuilds the store by replaying it. Because large restores are slow, **standby replicas** (`num.standby.replicas`) keep warm copies elsewhere.
- **Repartition topics**: changing the key before aggregating routes data through an internal repartition topic so all records for a key reach one task. This is Kafka Streams' shuffle, costing an extra write and read through the brokers.

```java
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "clicks-per-minute");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);

StreamsBuilder builder = new StreamsBuilder();
builder.stream("clicks", Consumed.with(Serdes.String(), clickSerde))   // key = campaign_id
    .filter((campaignId, click) -> click.isValid())
    .groupByKey()
    .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
    .count(Materialized.as("clicks-per-minute-store"))                // RocksDB + changelog
    .toStream()
    .map((windowedKey, count) -> KeyValue.pair(
        windowedKey.key() + "@" + windowedKey.window().start(), count))
    .to("campaign-clicks-per-minute", Produced.with(Serdes.String(), Serdes.Long()));

new KafkaStreams(builder.build(), props).start();
```

Because clicks are already keyed by `campaign_id`, `groupByKey` needs no repartition.

The trade-off is mostly operational. Kafka Streams fits teams that already run Kafka and want stream processing to be "just a service"; other sources and sinks need Kafka Connect. Flink fits when you need many connectors, very large state, rich event-time handling, or SQL over streams, and can afford to run (or buy) a Flink cluster.

### Scaling Stream Jobs

Stream jobs scale by parallelism, bounded at the source by partition count. Provision above peak: a job that handles exactly the peak rate can never drain the backlog after an outage.

**State size** is the other axis. Per-minute counts per campaign are tiny: even 1 million active campaigns × a handful of open windows × tens of bytes is well under a gigabyte. Deduplicating clicks by `click_id` over 24 hours is very different: 1.7 billion IDs at roughly 50 bytes each in RocksDB is on the order of 85 GB of state before overhead. Feasible with RocksDB and incremental checkpoints, but it changes recovery time and disk sizing. Estimate state size before choosing the design.

### Failure Modes and Anti-Patterns

1. **Unbounded state**: keyed state with no TTL or window end, such as "remember every user ever seen", grows until disk or checkpoints fail. Use windows, state TTL, or bounded dedup horizons.
2. **Hot keys**: the stream equivalent of skew. One viral campaign pins one subtask at full CPU and backpressures the whole job. Pre-aggregate with a salted key, then combine.
3. **Changing code without a savepoint plan**: removing an operator or changing its state type without stable UIDs and state compatibility can make the job unable to restore. Treat state schemas like database schemas.
4. **Side effects in operators**: calling an external API or sending an email from inside a `map` means replay after failure repeats the side effect. Side effects belong in sinks with idempotency or transactions.

## Event Time, Watermarks, and Windows

### Event Time vs Processing Time

**Event time** is when the click happened. **Processing time** is when the engine handles it. They differ because of network delay, offline mobile clients, retries, and backlog after an outage.

Window by processing time and a backlog after an outage dumps an hour of clicks into the current minute, corrupting history. Windowing by **event time** assigns each click to the minute it happened, so replaying the same input assigns events to the same windows; outputs match exactly only if late events are not dropped, because watermarks can advance differently during a fast replay. That determinism is what makes reprocessing (and Kappa) possible.

### Window Types

- **Tumbling windows**: fixed size, non-overlapping. "Clicks per campaign per minute" uses 1-minute tumbling windows; each click belongs to exactly one window.
- **Sliding (hopping) windows**: fixed size, overlapping, advancing by a slide interval. "Clicks in the last 10 minutes, updated every minute" is a 10-minute window sliding by 1 minute; each click belongs to 10 windows, so state and output are roughly 10 times larger.
- **Session windows**: per key, a window that stays open while events keep arriving and closes after a gap of inactivity. "A user's browsing session ends after 30 minutes without a click." Sessions have no fixed size and can merge when a late event bridges two sessions.

### Watermarks

With event time, the engine needs to know when a window is complete. It cannot know for certain, so it estimates with a **watermark**: a timestamp *W* asserting "I do not expect more events with event time earlier than *W*." When the watermark passes the end of a window, the window fires.

The common strategy is **bounded out-of-orderness**: maximum event time seen minus a fixed delay. With a 30-second bound, the 12:00 to 12:01 window fires once an event stamped 12:01:30 or later arrives. An operator's watermark is the **minimum** across its input partitions, so one idle partition holds back the job; Flink's `withIdleness` marks quiet partitions idle.

![A Flink pipeline for per-minute click counts: Kafka clicks with event time pass through a watermark generator set to maximum timestamp minus 30 seconds, are keyed by campaign ID into a one-minute tumbling window backed by RocksDB keyed state that is snapshotted to S3 by checkpoint barriers, and fired windows go to a transactional Kafka sink while clicks that arrive too late go to a side output.](https://d5osvdbc8um23.cloudfront.net/static-asset/blog_images/batch-vs-stream-processing-compared/03-flink-windows-and-watermarks.png)

### Allowed Lateness and Late Data

The watermark is a heuristic, so some events arrive after it. You have three choices, and good designs use all three:

1. **Choose the watermark delay** as a latency vs completeness dial. A longer delay means more complete first results but later dashboards.
2. **Allowed lateness**: keep window state for an extra period after the watermark passes, and re-fire the window with an updated count when late events arrive. Downstream must handle updates (upserts keyed by campaign and minute), not just appends.
3. **Side outputs** for events later than allowed lateness, so they are counted somewhere (a correction path, the batch job, or an audit table) instead of being silently dropped.

```java
WatermarkStrategy<Click> watermarks = WatermarkStrategy
    .<Click>forBoundedOutOfOrderness(Duration.ofSeconds(30))
    .withTimestampAssigner((click, recordTs) -> click.getEventTimeMillis())
    .withIdleness(Duration.ofMinutes(1));

OutputTag<Click> lateClicks = new OutputTag<Click>("late-clicks") {};

SingleOutputStreamOperator<CampaignMinuteCount> perMinute = env
    .fromSource(kafkaClickSource, watermarks, "clicks")
    .filter(Click::isValid)
    .keyBy(Click::getCampaignId)
    .window(TumblingEventTimeWindows.of(Duration.ofMinutes(1)))  // Time.minutes(1) before Flink 1.19
    .allowedLateness(Duration.ofMinutes(5))
    .sideOutputLateData(lateClicks)
    .aggregate(new CountAggregate(), new EmitCampaignMinute())
    .uid("per-minute-window");

perMinute.sinkTo(perMinuteKafkaSink);
perMinute.getSideOutput(lateClicks).sinkTo(lateClicksSink);
```

Kafka Streams expresses the same idea differently. It has no separate watermark; each task tracks **stream time**, the maximum record timestamp seen on its partitions, and a window accepts records until stream time passes the window end plus the **grace period** (`ofSizeAndGrace` above). Spark Structured Streaming uses `withWatermark("event_time", "30 seconds")`, which plays the role of both the watermark delay and the allowed lateness.

For the ad pipeline, the dashboard uses a 30-second watermark and 5 minutes of allowed lateness. The invoice sidesteps all of this: the batch job runs hours after midnight UTC, and anything later is an explicit correction.

## Exactly-Once Processing and State

### What "Exactly-Once" Actually Means

No system delivers each message exactly once over a network. Engines provide **exactly-once state semantics**: each record's effect on state and committed output is as if it were processed once, even if it is physically reprocessed after failures.

End-to-end exactly-once needs three parts, and missing any one breaks it:

1. **A replayable source**: the engine must be able to rewind and reread from a known position. Kafka (by offset), Kinesis (by sequence number), and files are replayable. A plain HTTP endpoint or a queue that deletes messages on acknowledgment is not, unless you put a log in front of it.
2. **Checkpointed state, consistent with source positions**: Flink checkpoints, Kafka Streams changelogs with transactional commits, or Spark's checkpointed offsets and state store. After failure, state and source position are restored together.
3. **A transactional or idempotent sink**: output must either commit atomically with the checkpoint (transactional) or be safe to write twice (idempotent).

### Transactional vs Idempotent Sinks

**Transactional sinks** use two-phase commit tied to checkpoints. Flink's `KafkaSink` with `DeliveryGuarantee.EXACTLY_ONCE` writes each checkpoint interval's output in a Kafka transaction, pre-commits at the barrier, and commits when the checkpoint completes; `read_committed` consumers never see uncommitted output. Consequences: a 60-second checkpoint interval adds up to 60 seconds of visible latency, and the producer's `transaction.timeout.ms` must stay below the broker's `transaction.max.timeout.ms` (15 minutes by default).

```java
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
env.setStateBackend(new EmbeddedRocksDBStateBackend(true));          // incremental
env.getCheckpointConfig().setCheckpointStorage("s3://ads-flink/checkpoints");

KafkaSink<CampaignMinuteCount> perMinuteKafkaSink = KafkaSink.<CampaignMinuteCount>builder()
    .setBootstrapServers(brokers)
    .setRecordSerializer(perMinuteSerializer)
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("ads-per-minute")
    .setProperty("transaction.timeout.ms", "600000")                 // below broker max
    .build();
```

**Kafka transactions** (Kafka 0.11 and later) let a producer write to multiple partitions and commit consumer offsets atomically: exactly the consume-process-produce loop, which is why Kafka Streams offers exactly-once with one config line. `exactly_once_v2` (Kafka 3.0 and later; `exactly_once_beta` in 2.6 to 2.8) uses one transactional producer per thread instead of per task, scaling better than the deprecated original `exactly_once`.

**Idempotent sinks** are often simpler: upsert `(campaign_id, minute) -> count` into a key-value store or `MERGE` into a table, and a replay produces the same final row with no transaction coordination. This does not work for append-only sinks or side effects like emails unless you deduplicate on a stable ID.

### Duplicates Upstream Still Need Handling

Exactly-once inside the engine does not fix duplicates already in the input. If the ad server's application code re-sends a click after a timeout (a new send the idempotent producer cannot deduplicate) and Kafka stores it twice, the engine faithfully counts both. Billing-grade counts deduplicate on `click_id`: keyed state with a TTL in a stream, or `dropDuplicates` over the day in batch. That is one more reason the invoice comes from batch: day-long dedup is a simple shuffle in Spark but a large, long-lived state store in a stream.

## Lambda vs Kappa Architecture

### Lambda: Two Paths, Two Codebases

The **Lambda architecture**, named by Nathan Marz, runs two pipelines over the same input:

- A **batch layer** that periodically recomputes complete, accurate views from the immutable raw dataset (the raw clicks in the lake).
- A **speed layer** that processes recent data in a stream to cover the gap since the last batch run, accepting approximation.
- A **serving layer** that merges the two at query time: batch results for older periods, speed results for the most recent hours.

Lambda emerged when stream processors could not be trusted for correctness, so the batch layer overwrote the speed layer's mistakes every night. Its well-known cost is **two codebases implementing the same logic**, often in different frameworks. Every business rule must be implemented twice and kept in sync, and when they disagree, the dashboard number changes at midnight. Unified APIs (Spark DataFrames, Flink SQL, Apache Beam) reduce this, but the two paths still run, fail, and deploy separately.

### Kappa: Replay the Log

The **Kappa architecture**, proposed by Jay Kreps in "Questioning the Lambda Architecture" (2014), drops the batch layer. Everything is a stream job reading from a durable log. To reprocess, whether to fix a bug, change logic, or backfill a new metric, you:

1. Start a **new version** of the job reading the log from the beginning (or from the earliest offset you need), writing to a new output table.
2. Let it catch up to the head of the log.
3. Switch readers to the new table and stop the old job.

![A comparison of Lambda and Kappa architectures: in Lambda, the Kafka clicks log is archived to a data lake for a Spark batch layer while a Flink speed layer consumes it directly, and a serving layer merges batch and speed results; in Kappa, a live job version 1 writes serving tables while a version 2 job replays the log from offset zero and its output is swapped in once it has caught up.](https://d5osvdbc8um23.cloudfront.net/static-asset/blog_images/batch-vs-stream-processing-compared/04-lambda-vs-kappa.png)

Kappa works because event-time logic makes stream results reproducible, so a replay gives the same answer a batch job would, provided late data is retained rather than dropped (watermarks advance differently during a fast replay). Its requirements are real, though:

- **The log must retain enough history.** Kafka's default retention is 7 days (`log.retention.hours=168`). Replaying a year means retaining a year, either with long retention, tiered storage (KIP-405, which moves older segments to object storage and became production-ready in recent Kafka releases), or by reading history from the lake with the same stream code in bounded mode.
- **Replay throughput.** 90 days at 860 GB per day is about 77 TB; at 1 GB per second that is roughly 21 hours. Know the number before promising "we'll just replay".
- **Event-time logic and double capacity.** Replays run far faster than real time, so wall-clock timers and lookups against changed external data break; and old and new versions run side by side.

### Backfills and Reprocessing in Practice

Whichever architecture you choose, backfills are where designs are tested. The patterns that hold up:

- **Keep raw, immutable input** (the lake or the log) as the source of truth. Derived tables are disposable.
- **Make outputs idempotent and partitioned by event time**, so rerunning a day or an hour replaces exactly that slice.
- **Version outputs** (`advertiser_daily_v2`) and switch readers atomically, rather than overwriting tables that people are querying.
- **Share code** between paths: a Flink job in batch mode over the lake, or one Spark DataFrame function used by both batch and streaming.

Many teams land on a hybrid: a Kappa-style streaming core for fresh results, plus batch jobs over the lake for billing-grade reconciliation and large backfills. That is technically Lambda, and that is fine, as long as both paths share logic and batch is clearly the system of record for money.

## Micro-Batch vs True Streaming

### Spark Structured Streaming: Micro-Batch by Default

Structured Streaming treats a stream as an unbounded table. By default it uses **micro-batch** execution: each trigger, it records the new source offsets in a write-ahead log in the checkpoint directory, processes that slice as a small batch job, updates the state store, and commits. The Spark documentation describes micro-batch latencies as low as about 100 milliseconds; in practice, with shuffles and a stateful aggregation, seconds are a more typical expectation.

```python
from pyspark.sql import functions as F

clicks = (spark.readStream.format("kafka")
    .option("kafka.bootstrap.servers", brokers)
    .option("subscribe", "clicks")
    .load()
    .select(F.from_json(F.col("value").cast("string"), click_schema).alias("c"))
    .select("c.*"))

per_minute = (clicks
    .where("is_valid")
    .withWatermark("event_time", "30 seconds")
    .groupBy(F.window("event_time", "1 minute"), "campaign_id")
    .count())

def upsert(batch_df, batch_id):
    # Idempotent MERGE into a Delta or Iceberg table; replaying a batch_id is harmless
    batch_df.createOrReplaceTempView("updates")
    batch_df.sparkSession.sql("""
        MERGE INTO serving.campaign_minute t
        USING updates u
        ON t.campaign_id = u.campaign_id AND t.window_start = u.window.start
        WHEN MATCHED THEN UPDATE SET t.clicks = u.count
        WHEN NOT MATCHED THEN INSERT (campaign_id, window_start, clicks)
                              VALUES (u.campaign_id, u.window.start, u.count)
    """)

(per_minute.writeStream
    .outputMode("update")
    .option("checkpointLocation", "s3://ads-spark/checkpoints/per_minute")
    .trigger(processingTime="10 seconds")
    .foreachBatch(upsert)
    .start())
```

Exactly-once comes from the same three parts: replayable Kafka, offsets and state checkpointed per batch, and an idempotent `MERGE` that tolerates a rerun of the same `batch_id`.

Spark also offers **continuous processing** (experimental since Spark 2.3) with millisecond-scale latency per the documentation, but only for map-like operations (no aggregations) and with at-least-once guarantees. It cannot run our per-minute counts; in interviews, treat Spark streaming as micro-batch unless you specifically cite a newer low-latency mode and its limits.

### True Streaming: Flink and Kafka Streams

Flink and Kafka Streams process records one at a time through long-running operators, with no per-batch scheduling, so latency is bounded mainly by buffering and window semantics. But Kafka Streams under exactly-once commits on an interval (`commit.interval.ms`, 100 ms by default in that mode), and Flink's transactional sinks expose output per checkpoint: record-at-a-time processing is not record-at-a-time visibility.

### Does the Difference Matter?

Less often than people argue. With 1-minute windows, a 30-second watermark, and final-result emission, results are 30 to 90 seconds old regardless of engine. The difference matters when:

- You need **sub-second reaction** to individual events (fraud blocking during an authorization, real-time bidding signals). True streaming wins.
- You need **fine-grained timers and per-key logic** (complex event processing, custom session logic, per-key timeouts). Flink's `KeyedProcessFunction` with event-time timers is more mature and expressive than Spark's `applyInPandasWithState` or `mapGroupsWithState`, though Spark 4.0's `transformWithState` adds timers and state TTL.
- Your team **already runs Spark** for batch and wants one engine, one set of skills, and shared code. Micro-batch wins on organizational cost.

## Head-to-Head Comparison

| Dimension | MapReduce | Spark (batch) | Spark Structured Streaming | Flink | Kafka Streams |
|---|---|---|---|---|---|
| Processing model | Batch, disk between jobs | Batch DAG, in-memory pipelining | Micro-batch (continuous mode limited) | Record-at-a-time; batch mode for bounded input | Record-at-a-time |
| Typical latency | Minutes to hours | Minutes to hours | Seconds | Sub-second to seconds | Sub-second to seconds |
| Deployment | Hadoop cluster (YARN) | Cluster (YARN, Kubernetes, managed) | Same as Spark | Dedicated cluster or managed service | Library inside your app |
| State management | None across jobs | None across jobs (cache within a job) | State store, checkpointed per batch | Keyed state; heap or RocksDB; checkpoints and savepoints | RocksDB state stores backed by changelog topics |
| Event time and windows | Manual | Manual (group by time columns) | Watermarks, tumbling, sliding, session | Full: watermarks, allowed lateness, side outputs, timers | Stream time, grace period, tumbling, hopping, session |
| Exactly-once | Deterministic rerun of tasks | Deterministic rerun plus idempotent output | Checkpointed offsets plus idempotent sink | Checkpoints plus transactional or idempotent sink | Kafka transactions (`exactly_once_v2`) |
| Sources and sinks | Files | Files, tables, JDBC, many connectors | Kafka, files, Kinesis via connectors | Many connectors | Kafka only (Connect for the rest) |
| Reprocessing | Rerun the job | Rerun per partition, cheap | Restart with new checkpoint from earlier offsets | Replay log or run bounded mode over lake | Reset offsets and rebuild state |
| Scaling unit | Map and reduce tasks | Partitions and executors | Partitions and executors | Operator parallelism, key groups | Input partitions, app instances |
| Operational cost | High, largely legacy | Moderate; scheduled clusters | Moderate; long-running cluster | High; 24/7 service, state ops | Low to moderate; it is your service |
| Best fit | Legacy pipelines | ETL, billing, backfills, ML features | Teams on Spark wanting fresh-enough results | Low-latency, large-state, event-time-heavy jobs | Kafka-native microservices doing stream logic |

Two rows deserve emphasis. **Reprocessing** is where batch quietly wins: rerunning one day in Spark is a routine operation, while reprocessing in a stream engine requires retained history and extra capacity. **Operational cost** is where Kafka Streams quietly wins: there is no new cluster to run, which is often the deciding factor for a small team.

## When to Pick Which

Start from freshness and correctness, not tools. The flowchart is the order in which I ask the questions.

![A decision flowchart: if results are not needed within minutes, use a scheduled Spark batch job; otherwise, if the pipeline reads from and writes to Kafka within a JVM service, use Kafka Streams inside that service; otherwise, if it needs large state, event-time logic, and low latency, use a Flink cluster; otherwise use Spark Structured Streaming.](https://d5osvdbc8um23.cloudfront.net/static-asset/blog_images/batch-vs-stream-processing-compared/05-decision-flowchart.png)

The first question is deliberately about batch: if nobody acts on a result within minutes, a scheduled Spark job is cheaper and easier to correct.

### Scenario 1: Ad-Click Billing and Pacing

**Requirements**: advertisers are invoiced on daily valid clicks; budgets must stop serving within about two minutes of exhaustion; about 20,000 clicks per second on average.

**Choice**: both paths, with clear ownership. Flink computes per-minute counts per campaign with a 30-second watermark and writes upserts to a key-value serving store that the ad server checks for pacing. A nightly Spark job over the lake deduplicates by `click_id`, applies the full fraud filtering logic, and writes the daily advertiser aggregates that feed invoicing. The batch output is the system of record for money; the stream output is the system of record for pacing. Share the click-validity logic as a library both jobs import.

**If pushed on Lambda's two-codebase cost**: the shared logic is small (validity rules), lives in one library, and a daily job reconciles stream totals against batch totals.

### Scenario 2: E-commerce Clickstream for Recommendations

**Requirements**: product-view and add-to-cart events power session-based recommendations; features must reflect the current session within seconds; models retrain weekly over months of history.

**Choice**: Flink with session windows (30-minute gap) keyed by `user_id`, writing features to a low-latency feature store. The weekly training set is built by Spark over the lake, the natural place for months of history and wide catalog joins.

### Scenario 3: Order Events to Downstream Microservices

**Requirements**: an order service publishes `OrderPlaced` events to Kafka; an inventory service must maintain per-SKU reserved counts and emit `StockLow` events; everything is Kafka in, Kafka out, written in Java.

**Choice**: Kafka Streams inside the inventory service, with a per-SKU state store and `exactly_once_v2`. No new cluster, the same deployment pipeline as other services, and exactly-once from input topic to output topic. Flink would work but adds a platform the team does not need.

## In the Interview

### How to Justify a Choice

Use the same four-step structure every time:

1. **State the freshness requirement per output**: "Pacing needs counts within about two minutes; invoices need exact daily totals by 6 AM."
2. **State the correctness requirement**: "A pacing count that is briefly low is acceptable; an invoice that double-counts retried clicks is not."
3. **Name the processing style, then the engine**: "Pacing is an unbounded, event-time windowed aggregation, so streaming, specifically Flink with RocksDB state. Invoicing is a bounded daily computation with global dedup, so a Spark batch job."
4. **Name the cost you are accepting**: "We run two paths, so we share validity logic as a library and reconcile daily."

The fourth step is the one most candidates skip, and it is the one interviewers listen for.

### Common Traps

- **"Streaming is real-time, so it's better."** If nobody acts on fresher data, a 24/7 stateful service is pure cost. Say what decision the fresher number enables.
- **"Kafka gives exactly-once."** Kafka transactions give atomic writes within Kafka. End-to-end exactly-once still needs a replayable source, checkpointed state, and a transactional or idempotent sink, and it does not remove duplicates already in the input.
- **Windowing by processing time.** Any replay or backlog corrupts the results. Always say "event time" and name the watermark strategy.
- **Ignoring late data.** Say what happens to a click that arrives 10 minutes late: allowed lateness with upserts, a side output, or batch correction.
- **"We'll just replay the log" without numbers.** State retention and replay time: "90 days is about 77 TB; at 1 GB per second that's most of a day, with double capacity."

### Likely Follow-Up Questions

- **"What happens if the Flink job crashes mid-window?"** It restores state from the last completed checkpoint, rewinds Kafka offsets, and reprocesses; uncommitted transactional output is aborted, so consumers see each result once.
- **"How do you deploy a new version of the stream job without losing state?"** Take a savepoint, stop the job, start the new version from the savepoint. Stable operator UIDs and compatible state schemas make that work.
- **"One campaign is getting 50 percent of all clicks. What breaks?"** One subtask or one shuffle partition. Salt the key into N sub-keys, pre-aggregate, then combine; in Spark batch, AQE skew handling also helps for joins.
- **"How would you backfill a new metric for 30 days?"** Run it in batch over the lake by event date into a new versioned table, validate, then switch readers. In Kappa, start a new job version from the earliest needed offset if retention covers it.
- **"Why not Spark Structured Streaming for everything?"** You can, if seconds of latency are fine and the team already runs Spark. Choose Flink when you need sub-second latency, fine-grained timers, or very large keyed state.

Related reading: [Designing a Log Aggregation System](/blog/system-design-log-aggregation), [Designing an Ad Serving and Real-Time Bidding System](/blog/system-design-ad-serving-rtb), [Designing a Recommendation Engine](/blog/system-design-recommendation-engine), [Message Queues vs Logs vs Pub/Sub](/blog/message-queues-vs-logs-vs-pubsub), and, for a lighter introduction, [The Real Talk on Batch vs Stream Processing](/blog/batch-vs-stream-processing).
