Stream Processing — Event Time, Windows, State & Exactly-Once
Stream processing is continuous computation over unbounded data: you keep results fresh as events arrive instead of recomputing everything on a schedule. The hard part is not reading from Kafka fast. It is answering four questions when data arrives late, out of order, and forever: what you compute, where in event time, when you emit, and how refinements relate. A production engine must also keep large keyed state consistent across crashes and rescaling, and get results into sinks exactly once.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What is the difference between event time and processing time, and why does it matter?
Answer
Event time is when the thing happened, embedded in the record. Processing time is when the engine sees it. Network delay, offline clients, and backlogs make them diverge, so per-minute results grouped by processing time are wrong whenever the pipeline lags or replays.
L2
Lambda versus Kappa. Which would you choose today?
Answer
Kappa by default: one codebase on an engine with checkpointed state and event-time windows, plus a replayable log or a lakehouse table for reprocessing. Lambda only if history exists solely in the warehouse and the streaming path can tolerate being approximate.
L3
Why is micro-batch latency bounded below?
Answer
Each micro-batch must be planned, scheduled, and committed. Results for data arriving just after a batch starts wait for the next trigger. Latency is roughly the trigger interval plus batch processing time.
L4
When would you choose Kafka Streams over Flink?
Answer
When inputs and outputs are Kafka topics, the logic is per-key enrichment or aggregation, and you want a library deployed like any microservice. Choose Flink for heterogeneous sources and sinks, very large state, complex event-time joins, or a dedicated platform team.
L5
What does exactly-once really mean in streaming?
Answer
Each input affects state and visible output exactly once, even if records are physically reprocessed after failures. It needs a replayable source, consistent snapshots of state and offsets, and an idempotent or transactional sink.
L6
Give a use case where streaming is the wrong tool.
Answer
Monthly invoicing with late corrections from upstream ledgers. Freshness is irrelevant, inputs get restated, and a deterministic batch re-run is simpler to audit.
L7
What metrics tell you a streaming job is healthy?
Answer
Consumer lag in records and time, watermark lag versus wall clock, checkpoint duration, size, and failures, backpressure ratio per operator, late-record and side-output counts, and end-to-end latency at the sink.
Failure modes
Processing-time windows for a business minute
Results change whenever the pipeline lags or a backlog is replayed.
Exactly-once treated as a checkbox
Engine state can be exactly-once while an HTTP sink still double-charges.
Unbounded state
Joins and dedupe without TTL or watermark cleanup grow until the job falls over.
No reprocessing plan
A code change has nowhere to replay from except a full batch rewrite.
Misconceptions
Streaming means reading from Kafka as fast as possible.
The hard part is event time, windows, triggers, state, and the sink. Throughput is the easy sentence.
Exactly-once means each record is processed once.
Crashes replay records. Exactly-once means the effect on state and on visible output happens once.
Lambda is still the default architecture.
One streaming codebase plus a replayable log is the default. Two code paths drift.
Interviewer traps
Picking Flink because it is the serious engine.
Ask for the freshness the business will pay for. A one-hour report does not need a record-at-a-time cluster.
Explaining consumer groups when the question is event time.
Point at the Kafka lesson for partitions and groups, then answer where in event time and when you emit.
Design scenario
Same prompt for every reader.
Requirements
Same alert set as a batch recompute, with detection latency near a second for fraud, and an honest nightly job for finance. Name the engine for each path.
Traffic / scale
A few hundred transactions per second, with mobile clients that upload an hour late.
Latency
Fraud wants about a second. Finance is happy the next morning.
Consistency
A replay of yesterday must produce the same fraud alerts. Finance must be able to restate a day.
Availability
A crashed fraud job restores from a checkpoint and does not double-alert. Finance can rerun.
Failure assumptions
- Events arrive out of order.
- One sink is Kafka and another is an HTTP webhook.
Constraints
- Do not run two copies of the fraud rule.
- Do not put the daily ledger on a 24x7 Flink job.
Prompt
A card fraud rule should flag three or more transactions for the same card inside one minute of event time. Finance also wants a daily reconciliation that absorbs late ledger corrections. The team already runs Spark.
API
What does the fraud alert record contain so a retry does not page twice?
Data
Which timestamp is the window key, and what happens to a click that arrives two hours late?
Architecture
Which path is Flink or Kafka Streams, which path stays batch, and where do you replay from?
When the result has to stay fresh
Prefer
One incremental job on a replayable log
Each record is applied once to keyed state. A replay of the same log with the same code produces the same windows.
- Latency is the trigger interval or a single record, not the next nightly job.
- Late data is a watermark and lateness policy, not a surprise rerun.
- Reprocessing is a new job version from an earlier offset, one codebase.
Alternative
Recompute the whole history on a timer
Simple and auditable, and stale between runs. A naive rerun every ten seconds matches micro-batch latency and rescans history every time.
- Sessions split across batch boundaries.
- A backlog replay with processing-time windows stuffs a day into the current minute.
- Two code paths, batch and speed, drift apart.
Four questions before you name an engine
Akidau's framing. The rest of this cluster is where, when, and how, plus state and the sink.
- 1
What are you computing?
A transform or an aggregate. The fraud rule below is a count per card per minute. - 2
Where in event time?
Windows. Tumbling, hopping, session, or global. Processing time is a different question. - 3
When do you emit?
Watermarks plus triggers. Early, on time, and late panes are choices. - 4
How do refinements relate?
Discard the delta, accumulate and overwrite, or retract the previous pane. - 5
Where does the effect land?
Keyed state must survive crashes. The sink must be idempotent or transactional.
Overview
Interviewers ask how you keep a fraud rule, a dashboard, or a session metric fresh. The weak answer is "we consume Kafka." The strong answer separates the business question from the engine.
Event time is when the thing happened. Processing time is when your operator saw it. They diverge because of network delay, offline phones, and backlogs. A per-minute billing number grouped by processing time changes whenever the job lags or you replay. That is the whole reason this cluster exists.
This page is the decision map. Windowing is where. Watermarks are when. Keyed state is how the memory survives a crash. Exactly-once sinks are how the output survives a crash. Joins and backpressure are how the job falls over in production.
Partitions, consumer groups, and delivery semantics for the log itself stay on Apache Kafka and delivery semantics. Change capture stays on CDC. Domain events that must commit with a row stay on event-driven architecture. Approximate counts stay on sketches. The log landing in a lake stays on object storage.
Lesson map
Stream Processing — Event Time, Windows, State & Exactly-Once
Stream processing is continuous computation over unbounded data: you keep results fresh as events arrive instead of recomputing everything on a schedule. The hard part is not reading from Kafka fast. It is answering four questions when data arrives late, out of order, and forever: what you compute, where in event time, when you emit, and how refinements relate. A production engine must also keep large keyed state consistent across crashes and rescaling, and get results into sinks exactly once.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB log["1. Log splits into batch and speed"] merge["2. Serving merges both views"] job["3. Kappa: one job replays v2"] cut["4. Cut over once v2 catches up"] log -->|1. Log splits into batch and speed| merge merge -->|2. Serving merges both views| job job -->|3. Kappa: one job replays v2| cut
Batch, micro-batch, and record-at-a-time
| Model | How it runs | Typical latency | Strengths | Costs | Examples |
|---|---|---|---|---|---|
| Batch | Scheduled job over a bounded snapshot | Minutes to a day | Simple, cheap, easy to reprocess, great for joins over full history | Stale results; sessions split across runs; late data needs re-runs | Spark, SQL, or dbt on a lake or warehouse |
| Micro-batch | The engine slices the stream into small batches, often every 1-10 seconds, with an incremental state store | About 1 second to minutes | Reuses a batch engine and SQL; high throughput; exactly-once via batch ids | Latency floor is the trigger interval; per-batch scheduling overhead | Spark Structured Streaming, default trigger |
| Record-at-a-time | Long-running operators process each event as it arrives | Milliseconds | Lowest latency; rich event-time semantics; huge keyed state | State backends, checkpoints, and backpressure | Apache Flink, Kafka Streams |
The same fraud rule, three or more transactions for one card inside one minute of event time, produces the same alerts in every model. Latency and redundant work change.
Press Run. Snippets must be self-contained — no network, files, or native modules.
batch every 120s alerts=True latency_s={'A@min0': 89, 'B@min1': 25} records_processed=19
micro-batch 10s alerts=True latency_s={'A@min0': 9, 'B@min1': 5} records_processed=11
batch rescan 10s alerts=True latency_s={'A@min0': 9, 'B@min1': 5} records_processed=181
streaming alerts=True latency_s={'A@min0': 1.2, 'B@min1': 1.2} records_processed=11All four runs emit the same alerts. Streaming notices about a second after the third transaction and touches each record once. Micro-batch at 10 seconds is a few seconds later and also touches each record once, because the state store carries history. A batch job every 10 seconds matches that latency and rescans history, 181 record-touches instead of 11. That rescan is why micro-batch engines keep an incremental state store.
When not to stream
Streaming is a cost you pay forever: 24x7 jobs, state, and on-call. Prefer batch when:
- Freshness need is hours. Daily finance, monthly billing reconciliations, model training sets.
- The logic needs the whole dataset. Global ranking, full re-dedupe, complex multi-way joins over history.
- Correctness beats latency and inputs are restated. Late corrections from upstream systems are easier to absorb with a re-run.
- The team cannot operate it. A broken stream with silent lag is worse than an honest nightly job.
A good senior answer: start with batch or micro-batch, and move a path to record-at-a-time only where the business value of seconds is real. Fraud, alerting, personalization, and ops are the usual yes.
Lambda versus Kappa
Decisions
- ?
1. Lambda or Kappa?
- batch plus speed2a. Batch layer
- batch plus speed2b. Speed layer
- one log2c. One job, replay
- 2
2a. Batch layer
- next3a. Serving merges views
- 3
2b. Speed layer
- next3a. Serving merges views
- 4
2c. One job, replay
- next3b. Cut over after replay
- 5
3a. Serving merges views
- 6
3b. Cut over after replay
| Lambda | Kappa | |
|---|---|---|
| Code paths | Two, batch plus streaming, that must agree | One |
| Reprocessing | The batch layer recomputes continuously | Replay the log with a new job version |
| Pain | Dual logic drift; merging views | Needs long log retention, or tiered storage, and a fast replay |
| Pick when | The streaming engine cannot be made correct, or history lives only in the lake | The engine has exactly-once state and event time |
Jay Kreps' argument is still the interview answer: maintaining the same logic in two distributed systems is exactly as painful as it seems. Modern lakehouse setups often land the log in object storage as Iceberg, Delta, or Hudi tables and replay from there. That is a Kappa variant. The storage contract is object storage, not a second copy of this page.
Choosing an engine
| Dimension | Apache Flink | Spark Structured Streaming | Kafka Streams |
|---|---|---|---|
| Execution | Record-at-a-time dataflow. JobManager plus TaskManagers | Micro-batch on Spark. Continuous mode exists and is limited | A library inside your app. Scale means more instances |
| Latency | Milliseconds | About seconds, the trigger interval | Milliseconds |
| Event time | First class, per source, idleness, alignment | withWatermark, one watermark per query | Stream time per partition, plus grace |
| State | Keyed state, heap or RocksDB, incremental checkpoints, savepoints | State store in the checkpoint location, HDFS or RocksDB | RocksDB backed by compacted changelog topics |
| Exactly-once | Barrier checkpoints plus two-phase-commit sinks | Replayable source plus an idempotent or transactional sink keyed by batch id | exactly_once_v2, Kafka-to-Kafka only |
| Sources and sinks | Kafka, Kinesis, files, CDC, JDBC, Iceberg, and more | Anything Spark reads or writes, strong lake integration | Kafka in, Kafka out. Connect for the rest |
| Ops cost | Highest. Cluster, state, checkpoints | Medium if you already run Spark | Lowest, no cluster. Rebalance and restore still hurt |
| Pick when | Low latency, big state, event-time correctness | The team is on Spark, seconds are fine, lakehouse sinks | Kafka-centric services: enrich, filter, route, moderate aggregations |
Press Run. Snippets must be self-contained — no network, files, or native modules.
daily finance report -> batch (Spark/SQL on a schedule)
enrich + route orders -> Kafka Streams
10s dashboard on lake -> Spark Structured Streaming
fraud: click+payment join -> Apache FlinkRead the branches out loud. Hours of freshness stay on a schedule even if the logic is stateful. Kafka in and Kafka out, without a big join, is a library. Seconds plus an existing Spark platform is Structured Streaming. A 200 ms event-time join is Flink even when the team has Spark, because the latency floor of micro-batch is the trigger.
Anatomy of a streaming job
Flow
- 1
1. Source offsets
- next2. Barrier rides the records
- 2
2. Barrier rides the records
- next3. Event time and watermarks
- 3
3. Event time and watermarks
- next4. keyBy, windows, state
- 4
4. keyBy, windows, state
- next5. Idempotent or txn sink
- 5
5. Idempotent or txn sink
One record, five responsibilities
The checkpoint coordinator rides along. It is not a second pipeline.
- 1
Source offsets
Kafka offsets, Kinesis sequence numbers, or files. If you cannot rewind, nothing downstream is exactly-once. - 2
Timestamps and watermarks
Arrival order becomes event-time progress. Idle partitions are a watermark bug, not a traffic dip. - 3
keyBy
All events for a key sit with that key's state. Skew here is the hottest risk. - 4
Window or stateful operator
Windows and keyed state live here. Unbounded buffers are how disks fill. - 5
Sink
Idempotent upsert or a transaction that commits with the checkpoint. An HTTP POST is neither.
The checkpoint coordinator injects barriers at the sources. Keyed state snapshots land next to the offsets, often in object storage. That cut is the state lesson. Whether the last hop is exactly-once is the sink lesson.
Pitfalls
Take the four cases in the TypeScript chooser. For each one, name the latency, whether state is large, and which sentence in the table you are using. Then change the fraud case to a one-hour SLO and watch it fall back to Spark.
Interview Q&A
What is the difference between event time and processing time, and why does it matter?
Answer
Event time is when the thing happened, embedded in the record. Processing time is when the engine sees it. Network delays, offline mobile clients, and backlogs make them diverge, so per-minute results grouped by processing time are wrong whenever the pipeline lags or replays.
Lambda versus Kappa. Which would you choose today?
Answer
Kappa by default: one codebase on an engine with checkpointed state and event-time windows, plus a replayable log (Kafka with long or tiered retention, or a lakehouse table) for reprocessing. Lambda only if history exists solely in the warehouse and the streaming path can tolerate being approximate.
Why is micro-batch latency bounded below?
Answer
Each micro-batch must be planned, scheduled, and committed. Results for data arriving just after a batch starts wait for the next trigger. Latency is roughly the trigger interval plus batch processing time.
When would you choose Kafka Streams over Flink?
Answer
When inputs and outputs are Kafka topics, the logic is per-key enrichment or aggregation, and you want a library deployed like any microservice. Choose Flink for heterogeneous sources and sinks, very large state, complex event-time joins, or a dedicated platform team.
What does exactly-once really mean in streaming?
Answer
Each input affects state and visible output exactly once, even if records are physically reprocessed after failures. It needs a replayable source, consistent snapshots of state and offsets, and an idempotent or transactional sink.
Give a use case where streaming is the wrong tool.
Answer
Monthly invoicing with late corrections from upstream ledgers. Freshness is irrelevant, inputs get restated, and a deterministic batch re-run is simpler to audit.
What metrics tell you a streaming job is healthy?
Answer
Consumer lag in records and in time, watermark lag versus wall clock, checkpoint duration, size, and failures, backpressure ratio per operator, late-record and side-output counts, and end-to-end latency at the sink.