People compare these three because each one can consume a Kafka topic, keep state, aggregate over windows and write results somewhere with exactly-once semantics. They differ in where the processing lives and what it costs to run. Flink is a dedicated streaming cluster built for per-event processing with large keyed state, event time and fine-grained checkpoints. Spark Structured Streaming is the batch engine extended to streams, processing micro-batches with the same DataFrame API your batch jobs already use. Kafka Streams is a library inside your service: no cluster, just more instances of your app, with state backed by Kafka topics. The question that decides it is: is stream processing a platform your organization runs, or a feature inside one service? A feature inside a service, Kafka in and Kafka out, points to Kafka Streams. A platform with large state, many sources and sinks, and strict latency points to Flink. A data team already living in Spark, with seconds of latency acceptable, points to Spark.
The Verdict#
Default: Flink for a shared stream processing platform; Kafka Streams for a single service that transforms Kafka topics; Spark Structured Streaming when the team already runs Spark and seconds of latency are fine.
| Pick Flink when | Pick Spark Structured Streaming when | Pick Kafka Streams when |
|---|---|---|
| Latency budgets are milliseconds to low seconds per event | Latency of ~1–10 seconds is acceptable (micro-batch) | Input and output are both Kafka and the logic belongs to one service |
| State is large (hundreds of GB to TBs) and keyed: sessions, joins, fraud features | The same logic runs in batch and streaming, and the team is fluent in Spark | The owning team wants to deploy it like any other microservice, with no cluster |
| You need rich event-time semantics: watermarks, late data, timers | You already run Spark on a managed platform with a lakehouse sink | State per instance stays in tens of GB and restore from changelog is tolerable |
| Many sources and sinks beyond Kafka; a central team offers streaming SQL | Throughput is high but per-event latency is not the product | Scale fits within the input topic's partition count |
🎯 Staff Move: "This enrichment job reads one topic and writes another, and the payments team owns it end to end, so I'd write it with Kafka Streams and ship it as part of their service. The fraud feature pipeline is different: 400 GB of keyed state, event-time windows and three sinks — that belongs on the Flink platform."
When the non-default wins:
- Kafka Streams with large state can still win when the team owns the domain end to end and invests in standby replicas, static membership and generous partition counts — Wise's ~50 TB of state across hundreds of apps shows it can be run at scale with standardization.
- Spark over Flink for a platform wins when the organization's data engineering already runs on Spark with a lakehouse, latency SLOs are in seconds, and one engine for batch and streaming saves a second platform team.
- No stream processor at all wins more often than people admit: a consumer that updates a database row per event, with an idempotency key, covers many "streaming" requirements without windows or managed state.
At a Glance#
| Dimension | Flink | Spark Structured Streaming | Kafka Streams |
|---|---|---|---|
| Deployment model | Cluster: JobManager + TaskManagers (on Kubernetes, YARN or managed) | Spark cluster: driver + executors (on Kubernetes, YARN or managed) | Library embedded in your JVM application; no separate cluster |
| Processing model | Record at a time, pipelined | Micro-batches by default; a real-time mode for stateless Scala queries arrived in Spark 4.1 | Record at a time |
| Sources / sinks | Kafka, Kinesis, Pulsar, files, CDC, databases, lakehouse tables, many connectors | Kafka, files, lakehouse tables, many connectors | Kafka only (other systems via Kafka Connect) |
| State | Keyed state in RocksDB or heap; disaggregated state backend in Flink 2.x | State store per partition (RocksDB provider available), checkpointed per batch | RocksDB stores per task, backed by Kafka changelog topics |
| Fault tolerance | Distributed snapshots (checkpoint barriers) to object storage; savepoints for upgrades | Offsets and state checkpointed to object storage each micro-batch | Changelog topics + Kafka offsets; standby replicas speed recovery |
| Exactly-once | End to end with transactional or idempotent sinks (two-phase commit) | End to end with replayable sources and idempotent or transactional sinks | Kafka-to-Kafka via Kafka transactions |
| Event time | First-class: watermarks, allowed lateness, timers, side outputs for late data | Watermarks and event-time windows; fewer low-level controls | Event-time windows with grace periods; stream-time advanced by records |
| Latency | ~10–100ms typical end to end | ~0.5–10s typical with micro-batches | ~10–100ms typical |
| Throughput | Millions of events/s per job on large clusters | Millions of events/s; strong for wide transforms | Bounded by partitions × per-instance throughput; ~10K–100K+ events/s per instance |
| Scaling model | Change parallelism; rescale from savepoint; key groups redistribute state | Add executors; partitions of input and shuffle | Add instances up to the input partition count; rebalance moves tasks |
| Operational burden | High: a platform to run, upgrade, monitor | Medium: often managed; streaming jobs on a batch platform | Low per app; spread across every team using it |
| Managed options | Several cloud and vendor-managed Flink services | Databricks, EMR, Dataproc, other Spark platforms | Runs wherever your app runs; Kafka can be managed |
| Cost shape | Always-on cluster sized for peak state and throughput | Always-on cluster; batch intervals trade cost for latency | Your app's instances plus Kafka changelog storage and traffic |
Latency and throughput figures are typical ranges for well-tuned jobs; checkpoint intervals, state size and sink behavior move them substantially.
Numbers to bring:
| Figure | Value | Condition |
|---|---|---|
| Flink end-to-end latency | ~10–100ms | Without transactional sinks; checkpoint interval adds visibility delay with them |
| Spark micro-batch latency | ~0.5–10s | Depends on trigger and batch planning overhead |
| Spark real-time mode | Single-digit ms p99 | Spark 4.1: stateless, single-stage Scala queries, Kafka sources |
| Kafka Streams parallelism ceiling | One task per input partition | Extra instances sit idle |
| Changelog restore | ~50–100 MB/s per task | 50 GB of state ≈ 10–20 minutes without standbys |
| Flink checkpoint interval | 10s–5min typical | Shorter = fresher transactional output, more overhead |
| Flink 2.0 | Released March 2025 | Disaggregated state backend, DataSet and Scala APIs removed |
How They Actually Differ#
1. Cluster vs Library: Who Owns the Runtime#
This is the deepest difference, and it is organizational.
| Flink | Spark | Kafka Streams | |
|---|---|---|---|
| Unit of deployment | A job submitted to a cluster | An application on a Spark cluster | Your service's container |
| Who is paged when it stalls | The streaming platform team, then the job owner | The data platform team, then the job owner | The service team, like any other outage in their service |
| How code ships | Savepoint, stop, start new version from savepoint | Stop query, deploy, restart from checkpoint (compatible changes only) | Rolling deploy of the service |
| Resource isolation | Slots per TaskManager; jobs share or get dedicated clusters | Executors per application | Process boundaries of the service |
Kafka Streams wins when a team should own its processing end to end: no ticket to a platform team, the same CI/CD, the same on-call. It loses when 30 teams each reinvent windowing, state recovery and monitoring, and every one of them hits the same rebalance incidents. Flink wins when a platform team offers streaming as a service — often through SQL — so product teams write queries, not operators.
🎯 Staff Insight: "No cluster" does not mean "no operations". With Kafka Streams the operations are distributed across every team that uses it. That is the right trade for 3 services and the wrong one for 30.
2. Per-Event vs Micro-Batch#
Flink and Kafka Streams process each record as it arrives, pipelined through operators. Spark Structured Streaming, by default, collects a micro-batch (for example every 1 second), plans it, runs it as a small batch job, commits offsets and state, and repeats. Planning and commit overhead put a floor under latency — commonly hundreds of milliseconds to seconds — but micro-batches amortize per-record overhead and make throughput and exactly-once bookkeeping simple.
As of Spark 4.1 (released December 2025), a real-time mode runs long-lived tasks that process records continuously, with documented p99 latencies in the single-digit milliseconds for supported queries. In 4.1 it covers stateless, single-stage Scala queries with Kafka sources and Kafka or foreach sinks. That narrows the latency gap for simple transforms; for stateful, event-time-heavy pipelines, Flink remains the reference design.
3. State: Where It Lives and How It Recovers#
Large state is where these systems separate under stress.
- Flink keeps keyed state in RocksDB on local disk (or, in Flink 2.x, a disaggregated backend that keeps the primary copy in object storage) and takes periodic distributed snapshots: barriers flow through the dataflow, each operator snapshots its state, and a complete checkpoint goes to object storage. Incremental checkpoints upload only changed files, so TB-scale state checkpoints in seconds to minutes. Recovery restores from the last checkpoint and replays the source from the matching offsets.
- Spark checkpoints offsets and state store files to object storage each micro-batch. Recovery reloads state for the partitions an executor owns; large state makes batches slower and recovery longer.
- Kafka Streams writes every state update to a compacted changelog topic. When an instance dies, another rebuilds that task's RocksDB store by replaying the changelog — which, for 50 GB of state at ~50–100 MB/s, is roughly 10–20 minutes of that partition not processing unless a standby replica is already warm.
| State size per job | Flink | Spark | Kafka Streams |
|---|---|---|---|
| Under ~10 GB | Easy | Easy | Easy |
| 10–100 GB | Easy | Workable; watch batch time | Needs standby replicas; restore time matters |
| 100 GB–10 TB | Designed for it: incremental checkpoints, rescaling by key groups | Hard; long batches and recoveries | Painful: changelog storage and restore time dominate |
4. Exactly-Once Is a Property of the Whole Pipeline#
All three offer exactly-once processing, and all three mean "the state and outputs reflect each input once, as long as the sink cooperates". Flink uses two-phase commit with transactional sinks (Kafka transactions, some databases) or idempotent writes, which means output becomes visible only when a checkpoint completes — a checkpoint interval of 60s means downstream readers with read-committed isolation see results up to a minute late. Kafka Streams uses Kafka transactions to atomically commit output records, changelog writes and input offsets — clean, but only inside Kafka. Spark relies on replayable sources plus idempotent or transactional sinks per micro-batch.
The Staff line: exactly-once ends at the first non-transactional side effect — an email, an HTTP call, a cache write. Make those idempotent with a key, or move them to a consumer that is.
5. Event Time and Late Data#
Flink is the most expressive: per-source watermarks, idle-source handling, allowed lateness on windows, timers per key for custom session logic, and side outputs that route late events somewhere instead of dropping them. Spark supports watermarks and event-time windows with less low-level control. Kafka Streams advances "stream time" from the records it sees per partition and closes windows after a grace period, which works well until a partition goes quiet and its windows never close. For pipelines where late or out-of-order data changes money — billing, ads, fraud — this is often the deciding factor.
Late data, side by side:
| Late-data need | Flink | Spark | Kafka Streams |
|---|---|---|---|
| Per-source watermarks with idle-source handling | Yes | Global watermark across sources | Stream time per partition |
| Keep windows open for late events | Allowed lateness per window | Watermark delay | Grace period per window |
| Route too-late events elsewhere | Side outputs | Dropped (or handled by a separate batch correction) | Dropped, with a metric |
| Custom per-key timers (sessions, timeouts) | Process functions with timers | Arbitrary stateful processing with timeouts | Punctuators on processors |
Where Each One Breaks#
| System | Failure mode | Symptom | Detection | Mitigation | Owner |
|---|---|---|---|---|---|
| Flink | Checkpoints time out | Checkpoint duration climbs past interval; eventually fails; recovery replays far back | lastCheckpointDuration, failed checkpoints, alignment time | Unaligned checkpoints under backpressure, incremental checkpoints, fix the slow operator | Job owner + platform |
| Flink | Backpressure from a slow sink | Source lag grows; busy operators at 100% | Backpressure ratio per operator, consumer lag | Async I/O, batching sink writes, scale the sink | Job owner |
| Flink | State growth without TTL | RocksDB grows forever; checkpoints and restores slow down | State size per operator | State TTL, key expiry, compaction tuning | Job owner |
| Flink | Upgrade breaks state compatibility | New job cannot restore the savepoint | Failed restore in staging | Stable operator UIDs, schema evolution rules, savepoint compatibility tests | Job owner |
| Spark | Batch duration exceeds trigger interval | Micro-batches queue; latency grows without bound | Batch duration vs trigger, input rows/s vs processed rows/s | More executors, smaller state, better partitioning | Data team |
| Spark | State store bloat | Each batch slower as state grows; executor memory pressure | State rows and memory per operator | RocksDB state store, watermarks to drop old state | Data team |
| Spark | Incompatible query change | Restart from checkpoint fails after changing the query | Failed restart | Plan compatible changes; new checkpoint + backfill for breaking ones | Data team |
| Kafka Streams | Rebalance storms | Instances repeatedly rebalance during deploys; processing pauses minutes each time | Rebalance count, task restore time | Static membership, cooperative rebalancing, standby replicas, slow rolling deploys | Service team |
| Kafka Streams | Long state restore | After a crash, a task replays a 50 GB changelog; partition stalls ~15 min | Restore progress metrics, consumer lag per partition | Standby replicas, smaller state per task, more partitions | Service team |
| Kafka Streams | Partition ceiling | Adding instances past the partition count adds idle instances | Idle tasks, lag | Repartition topics (a migration); plan partitions up front | Service team |
| All | Poison record | One malformed event crashes the job; restart loop; lag grows | Restart count, exception rate | Dead-letter output for bad records; schema validation at ingest | Job owner |
The production surprise for each:
- Flink: the job is fine; the checkpoint is not. Most Flink incidents start as checkpoints slowing down under backpressure until recovery has to replay an hour.
- Spark: latency is stable at 2 seconds until state grows, then each batch takes longer than the trigger and lag compounds.
- Kafka Streams: deploys are outages. A rolling restart of 20 instances without static membership and standbys can mean 20 rebalances and repeated state restores.
Incident Sketch: The Checkpoint Spiral#
t=0 Downstream database slows; the Flink sink's writes take 5x longer
t=+2min Backpressure reaches sources; aligned checkpoint barriers wait behind full buffers
t=+10min Checkpoint duration 9 minutes against a 1-minute interval; then timeouts
t=+25min A TaskManager dies; job restores from a checkpoint 25 minutes old
t=+26min Replays 25 minutes of Kafka at full speed into the still-slow sink
t=+60min Lag still growing; transactional output to consumers stalled the whole time
Detection: checkpoint duration vs interval, failed checkpoint count, backpressure per operator, consumer lag. Mitigation: unaligned checkpoints under backpressure, async and batched sink writes, a sink with its own capacity alert. Owner: job owner for the sink, platform for checkpoint defaults.
Incident Sketch: The Deploy That Paused Processing for 40 Minutes#
t=0 Rolling restart of 20 Kafka Streams instances, 30 seconds apart
t=+30s Each instance leaving and joining triggers a rebalance; tasks move
t=+5min Moved stateful tasks restore RocksDB from changelog topics
t=+12min Some tasks moved twice; restores restart
t=+40min All tasks restored; lag peaks at 35 minutes of events
Prevention: static group membership so a restart within the session timeout does not rebalance, standby replicas so moved tasks are already warm, and slower rolling deploys gated on restore completion. Owner: the service team, which is exactly why the library model needs a shared template.
Cost and Operations#
| Flink | Spark Structured Streaming | Kafka Streams | |
|---|---|---|---|
| Who runs it | A streaming platform team (commonly 3–6 engineers for a shared platform) or a managed service | A data platform team, often on a managed Spark platform | Each service team |
| What the bill scales with | TaskManager cores and memory for peak throughput and state; checkpoint storage | Executor hours (always-on for streaming); trigger interval trades latency for cost | Service instances; Kafka storage and traffic for changelog and repartition topics |
| Hidden cost | Platform headcount, upgrade cycles, connector maintenance | Always-on clusters on a platform priced for batch | Changelog topics can double Kafka storage and traffic for stateful apps; duplicated expertise across teams |
| Cheapest when | Many jobs share a platform; large state | Logic shared with batch; latency relaxed | Few, small, Kafka-to-Kafka apps |
Rough comparison for one pipeline doing 50K events/s with 100 GB of state: Flink needs on the order of 8–16 cores plus RocksDB disks and checkpoint storage; Spark similar cores but more memory headroom for batch processing; Kafka Streams similar app cores plus changelog topics holding a compacted copy of the 100 GB (×3 replication in Kafka) and the network to write every state change twice. The infrastructure is comparable. The difference is who carries the pager and how many teams each learn the same lessons.
| Scale | Sensible choice | Rough infrastructure | People |
|---|---|---|---|
| Small (1–3 Kafka-to-Kafka jobs, small state) | Kafka Streams inside the owning services | The services' own instances plus changelog topics | Owning teams |
| Medium (10–30 jobs, several sinks, some large state) | A managed Flink service, or Spark if the team is already there | Low to mid thousands of dollars a month | 1–3 engineers owning the platform |
| Large (100+ jobs, TB-scale state, SQL users) | Self-run or managed Flink platform with SQL; Kafka Streams allowed for team-owned services | Tens of thousands a month | A streaming platform team of 3–6 |
Assumptions: cloud list prices, always-on jobs, state on SSD with checkpoints in object storage; order of magnitude only.
What the pager looks like for each:
| What the on-call actually does | Flink | Spark | Kafka Streams |
|---|---|---|---|
| Most common page | Checkpoint failures, backpressure, consumer lag | Batch duration over trigger, lag | Rebalances, restore time, lag |
| First look | Job graph backpressure and checkpoint history | Streaming query progress, batch durations | Rebalance and restore metrics per instance |
| Usual fix | Scale the slow operator or sink; restore from savepoint | Add executors; trim state with watermarks | Static membership, standbys, more partitions |
| Upgrade ritual | Savepoint, deploy, restore, verify | Stop, deploy, restart from checkpoint | Rolling deploy gated on restores |
🧭 Principal Insight: Standardize on one cluster engine for the platform and allow Kafka Streams for team-owned Kafka-to-Kafka services. Two platform engines double the expertise you need; banning the library pushes small jobs onto a platform queue they do not need.
Switching Later#
| Migration | Difficulty | What's hard to undo |
|---|---|---|
| Kafka Streams → Flink | Moderate: logic rewrites; state rebuilt from Kafka history or bootstrapped | State held only in changelog topics must be replayed or exported |
| Spark Streaming → Flink | Moderate: DataFrame logic maps to Flink SQL or Table API; state cannot be transferred | Checkpoints are engine-specific; rebuild state by replay |
| Flink → Kafka Streams | Hard if you rely on non-Kafka sources, sinks or large state | Connectors and event-time features |
| Any engine upgrade with state | Moderate | State schema compatibility; operator IDs in Flink; query compatibility in Spark |
| Adding partitions to input topics | Hard for Kafka Streams and keyed state | Key-to-partition mapping changes; state must be reshuffled |
Moving a Kafka Streams app to Flink, in order:
- Run the Flink job in shadow, reading the same input topics and writing to a parallel output topic.
- Bootstrap state by replaying retained input (or the Kafka Streams changelog) before comparing outputs.
- Compare outputs per key for a full business cycle (a week covers weekday and weekend patterns).
- Switch consumers to the new output topic; keep the old app running read-only for rollback.
- Retire the old app and its changelog and repartition topics only after the soak period.
The one-way doors:
- State that exists only in the stream processor. If a session store or feature table lives only in Flink state or a Kafka Streams changelog, migrating engines means replaying history. Keep a replayable source (retained topics or a lakehouse table).
- Input partition counts for keyed state. Pick for 3–5 years of growth; changing them reshuffles every key.
- Exactly-once wiring into downstream consumers. Once consumers depend on transactional read-committed semantics, switching engines changes visibility timing they rely on.
How Real Companies Chose#
Uber — Flink SQL as a Streaming Platform (AthenaX)#
Uber built AthenaX so that both engineers and non-engineers could write streaming analytics in SQL, compiled into Flink dataflows. Uber reported more than a trillion messages a day through Kafka, over 220 AthenaX applications in production, processing of several million messages per second with only eight YARN containers for some workloads, and that more than 70% of streaming applications could be expressed in SQL (Uber Engineering).
Staff insight: This is the platform model — a central team runs Flink, product teams write SQL. The "70% expressible in SQL" figure is the argument for building a platform rather than letting every team write operators.
Wise — Hundreds of Kafka Streams Applications#
Wise runs roughly 300 stateful stream processing applications on Kafka Streams among about 400 microservices, holding around 50 TB of state on Kubernetes, to aggregate, join and process the event streams behind real-time international money transfers. Teams use a shared DSL and container images for governance, and interactive queries to read application state (Confluent community).
Staff insight: The library model scales organizationally when a platform team standardizes how it is used — shared images, conventions, governance — even without a shared cluster.
Zalando — Kafka Streams Instead of a Spark Cluster#
Zalando's fashion insight team built a real-time ranking of fashion websites on Kafka Streams. They explicitly ruled out Hadoop for lack of experience in a small team, and found that running Spark full time as separate infrastructure added cost they did not need, while a library let them stay close to the data and deploy containerized services (Confluent Blog).
Staff insight: For a small team with a Kafka-to-Kafka problem, "no new cluster" was the deciding feature. Team size and existing infrastructure are legitimate inputs to an engine choice.
Follow-Ups to Expect#
| After You Say... | They Will Ask... | What They're Testing |
|---|---|---|
| "Flink with exactly-once" | "When does a downstream consumer see the output?" | Transactional sinks commit on checkpoint; latency equals checkpoint interval |
| "Flink" | "The job has 2 TB of state. How do you deploy a new version?" | Savepoints, operator UIDs, state schema evolution |
| "Kafka Streams" | "An instance dies holding 50 GB of state. What happens?" | Changelog restore time; standby replicas |
| "Kafka Streams" | "You need 4× throughput. Can you add instances?" | Partition ceiling; repartitioning is a migration |
| "Spark Structured Streaming" | "Can you get end-to-end latency under 100ms?" | Micro-batch floor; real-time mode limits; when to switch engines |
| "Event-time windows" | "Events arrive 2 hours late. What happens to them?" | Watermarks, allowed lateness, side outputs, correction pipelines |
| "Exactly-once" | "The job also sends an email per event." | Side effects outside the transaction; idempotency keys |
| "One platform for all teams" | "Who is paged when a team's job is stuck?" | Ownership split between platform and job owners |
| "Flink SQL for product teams" | "Who reviews a query that will hold 500 GB of state?" | Platform guardrails: state TTL, cost review, quotas |
| "Standby replicas" | "What do they cost?" | Extra instances and changelog consumption vs restore time |
| "Spark for batch and streaming" | "How do you backfill a month without hurting the live job?" | Separate backfill job, same code, idempotent sink |
What to Say in the Interview#
"The first question is whether stream processing is a platform or a feature. A Kafka-to-Kafka transform owned by one team is a Kafka Streams library in their service; large keyed state with event-time windows and several sinks belongs on Flink."
"Exactly-once only covers the engine and transactional sinks. Anything else — emails, HTTP calls, cache writes — gets an idempotency key."
"With Flink, downstream readers see transactional output when a checkpoint completes, so a 30-second checkpoint interval is a 30-second visibility delay, and I'd size it against the freshness requirement."
"If the team already lives in Spark and two seconds of latency is fine, Structured Streaming lets them share one codebase between batch backfills and the live pipeline — that's worth more than milliseconds they don't need."
Related Guides#
- Flink — checkpoints, state backends, watermarks and operations in depth
- Kafka — partitions, transactions and the log every engine reads
- Design a Stream Processing Engine — the internals these systems share
- Batch and Stream Pipelines — lambda vs kappa, backfills and reconciliation
- Design Ad Click Aggregation — windowed aggregation with late data and exactly-once
- Backpressure — what happens when a sink cannot keep up
- Design a Message Broker — the log underneath every pipeline
- Idempotency — making side effects safe under replay