Hiring BarSupport

Flink vs Spark Streaming vs Kafka Streams

Comparison19 min read3 diagrams

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 whenPick Spark Structured Streaming whenPick Kafka Streams when
Latency budgets are milliseconds to low seconds per eventLatency 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 featuresThe same logic runs in batch and streaming, and the team is fluent in SparkThe owning team wants to deploy it like any other microservice, with no cluster
You need rich event-time semantics: watermarks, late data, timersYou already run Spark on a managed platform with a lakehouse sinkState per instance stays in tens of GB and restore from changelog is tolerable
Many sources and sinks beyond Kafka; a central team offers streaming SQLThroughput is high but per-event latency is not the productScale 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."

Diagram: The Verdict

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#

DimensionFlinkSpark Structured StreamingKafka Streams
Deployment modelCluster: 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 modelRecord at a time, pipelinedMicro-batches by default; a real-time mode for stateless Scala queries arrived in Spark 4.1Record at a time
Sources / sinksKafka, Kinesis, Pulsar, files, CDC, databases, lakehouse tables, many connectorsKafka, files, lakehouse tables, many connectorsKafka only (other systems via Kafka Connect)
StateKeyed state in RocksDB or heap; disaggregated state backend in Flink 2.xState store per partition (RocksDB provider available), checkpointed per batchRocksDB stores per task, backed by Kafka changelog topics
Fault toleranceDistributed snapshots (checkpoint barriers) to object storage; savepoints for upgradesOffsets and state checkpointed to object storage each micro-batchChangelog topics + Kafka offsets; standby replicas speed recovery
Exactly-onceEnd to end with transactional or idempotent sinks (two-phase commit)End to end with replayable sources and idempotent or transactional sinksKafka-to-Kafka via Kafka transactions
Event timeFirst-class: watermarks, allowed lateness, timers, side outputs for late dataWatermarks and event-time windows; fewer low-level controlsEvent-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
ThroughputMillions of events/s per job on large clustersMillions of events/s; strong for wide transformsBounded by partitions × per-instance throughput; ~10K–100K+ events/s per instance
Scaling modelChange parallelism; rescale from savepoint; key groups redistribute stateAdd executors; partitions of input and shuffleAdd instances up to the input partition count; rebalance moves tasks
Operational burdenHigh: a platform to run, upgrade, monitorMedium: often managed; streaming jobs on a batch platformLow per app; spread across every team using it
Managed optionsSeveral cloud and vendor-managed Flink servicesDatabricks, EMR, Dataproc, other Spark platformsRuns wherever your app runs; Kafka can be managed
Cost shapeAlways-on cluster sized for peak state and throughputAlways-on cluster; batch intervals trade cost for latencyYour 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:

FigureValueCondition
Flink end-to-end latency~10–100msWithout transactional sinks; checkpoint interval adds visibility delay with them
Spark micro-batch latency~0.5–10sDepends on trigger and batch planning overhead
Spark real-time modeSingle-digit ms p99Spark 4.1: stateless, single-stage Scala queries, Kafka sources
Kafka Streams parallelism ceilingOne task per input partitionExtra instances sit idle
Changelog restore~50–100 MB/s per task50 GB of state ≈ 10–20 minutes without standbys
Flink checkpoint interval10s–5min typicalShorter = fresher transactional output, more overhead
Flink 2.0Released March 2025Disaggregated 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.

FlinkSparkKafka Streams
Unit of deploymentA job submitted to a clusterAn application on a Spark clusterYour service's container
Who is paged when it stallsThe streaming platform team, then the job ownerThe data platform team, then the job ownerThe service team, like any other outage in their service
How code shipsSavepoint, stop, start new version from savepointStop query, deploy, restart from checkpoint (compatible changes only)Rolling deploy of the service
Resource isolationSlots per TaskManager; jobs share or get dedicated clustersExecutors per applicationProcess 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.

Diagram: 2. Per-Event vs Micro-Batch

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 jobFlinkSparkKafka Streams
Under ~10 GBEasyEasyEasy
10–100 GBEasyWorkable; watch batch timeNeeds standby replicas; restore time matters
100 GB–10 TBDesigned for it: incremental checkpoints, rescaling by key groupsHard; long batches and recoveriesPainful: 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 needFlinkSparkKafka Streams
Per-source watermarks with idle-source handlingYesGlobal watermark across sourcesStream time per partition
Keep windows open for late eventsAllowed lateness per windowWatermark delayGrace period per window
Route too-late events elsewhereSide outputsDropped (or handled by a separate batch correction)Dropped, with a metric
Custom per-key timers (sessions, timeouts)Process functions with timersArbitrary stateful processing with timeoutsPunctuators on processors

Where Each One Breaks#

SystemFailure modeSymptomDetectionMitigationOwner
FlinkCheckpoints time outCheckpoint duration climbs past interval; eventually fails; recovery replays far backlastCheckpointDuration, failed checkpoints, alignment timeUnaligned checkpoints under backpressure, incremental checkpoints, fix the slow operatorJob owner + platform
FlinkBackpressure from a slow sinkSource lag grows; busy operators at 100%Backpressure ratio per operator, consumer lagAsync I/O, batching sink writes, scale the sinkJob owner
FlinkState growth without TTLRocksDB grows forever; checkpoints and restores slow downState size per operatorState TTL, key expiry, compaction tuningJob owner
FlinkUpgrade breaks state compatibilityNew job cannot restore the savepointFailed restore in stagingStable operator UIDs, schema evolution rules, savepoint compatibility testsJob owner
SparkBatch duration exceeds trigger intervalMicro-batches queue; latency grows without boundBatch duration vs trigger, input rows/s vs processed rows/sMore executors, smaller state, better partitioningData team
SparkState store bloatEach batch slower as state grows; executor memory pressureState rows and memory per operatorRocksDB state store, watermarks to drop old stateData team
SparkIncompatible query changeRestart from checkpoint fails after changing the queryFailed restartPlan compatible changes; new checkpoint + backfill for breaking onesData team
Kafka StreamsRebalance stormsInstances repeatedly rebalance during deploys; processing pauses minutes each timeRebalance count, task restore timeStatic membership, cooperative rebalancing, standby replicas, slow rolling deploysService team
Kafka StreamsLong state restoreAfter a crash, a task replays a 50 GB changelog; partition stalls ~15 minRestore progress metrics, consumer lag per partitionStandby replicas, smaller state per task, more partitionsService team
Kafka StreamsPartition ceilingAdding instances past the partition count adds idle instancesIdle tasks, lagRepartition topics (a migration); plan partitions up frontService team
AllPoison recordOne malformed event crashes the job; restart loop; lag growsRestart count, exception rateDead-letter output for bad records; schema validation at ingestJob 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#

FlinkSpark Structured StreamingKafka Streams
Who runs itA streaming platform team (commonly 3–6 engineers for a shared platform) or a managed serviceA data platform team, often on a managed Spark platformEach service team
What the bill scales withTaskManager cores and memory for peak throughput and state; checkpoint storageExecutor hours (always-on for streaming); trigger interval trades latency for costService instances; Kafka storage and traffic for changelog and repartition topics
Hidden costPlatform headcount, upgrade cycles, connector maintenanceAlways-on clusters on a platform priced for batchChangelog topics can double Kafka storage and traffic for stateful apps; duplicated expertise across teams
Cheapest whenMany jobs share a platform; large stateLogic shared with batch; latency relaxedFew, 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.

ScaleSensible choiceRough infrastructurePeople
Small (1–3 Kafka-to-Kafka jobs, small state)Kafka Streams inside the owning servicesThe services' own instances plus changelog topicsOwning teams
Medium (10–30 jobs, several sinks, some large state)A managed Flink service, or Spark if the team is already thereLow to mid thousands of dollars a month1–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 servicesTens of thousands a monthA 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 doesFlinkSparkKafka Streams
Most common pageCheckpoint failures, backpressure, consumer lagBatch duration over trigger, lagRebalances, restore time, lag
First lookJob graph backpressure and checkpoint historyStreaming query progress, batch durationsRebalance and restore metrics per instance
Usual fixScale the slow operator or sink; restore from savepointAdd executors; trim state with watermarksStatic membership, standbys, more partitions
Upgrade ritualSavepoint, deploy, restore, verifyStop, deploy, restart from checkpointRolling 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#

MigrationDifficultyWhat's hard to undo
Kafka Streams → FlinkModerate: logic rewrites; state rebuilt from Kafka history or bootstrappedState held only in changelog topics must be replayed or exported
Spark Streaming → FlinkModerate: DataFrame logic maps to Flink SQL or Table API; state cannot be transferredCheckpoints are engine-specific; rebuild state by replay
Flink → Kafka StreamsHard if you rely on non-Kafka sources, sinks or large stateConnectors and event-time features
Any engine upgrade with stateModerateState schema compatibility; operator IDs in Flink; query compatibility in Spark
Adding partitions to input topicsHard for Kafka Streams and keyed stateKey-to-partition mapping changes; state must be reshuffled
Diagram: Switching Later

Moving a Kafka Streams app to Flink, in order:

  1. Run the Flink job in shadow, reading the same input topics and writing to a parallel output topic.
  2. Bootstrap state by replaying retained input (or the Kafka Streams changelog) before comparing outputs.
  3. Compare outputs per key for a full business cycle (a week covers weekday and weekend patterns).
  4. Switch consumers to the new output topic; keep the old app running read-only for rollback.
  5. Retire the old app and its changelog and repartition topics only after the soak period.

The one-way doors:

  1. 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).
  2. Input partition counts for keyed state. Pick for 3–5 years of growth; changing them reshuffles every key.
  3. 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 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."

  1. Loading the index…