Stream Processing System: Kafka and Flink, and the Watermark You Will Misconfigure
Event time versus processing time, checkpointing and exactly-once state, backpressure, and choosing heap or disk-backed state with an explicit recovery budget.

Stream processing is batch processing where the batch is never finished. That sounds like a small difference and it changes almost every design decision: you cannot reprocess everything, you cannot assume a stable input, and the only thing that makes results correct is deciding how to treat time.
Most stream processing bugs are time bugs. Everything else — state, backpressure, skew — is real and worth understanding, but the failure you will actually debug at 3am is a window that closed before the data arrived.
The scale
200 producers
50,000 messages/sec
= 4.3 billion messages/day
2 KB average
= 8.6 TB/day into Kafka
7 days retention
= 60 TB
the Flink job
4 parallel instances
joins a 50M-key state space
1 minute tumbling windows keyed by customer
writes to a transactional sink
the state
50M customers x ~500 bytes
= 25 GB of keyed state
with a TTL of 7 days
against a 32 GB per-task state limit
The number that matters is the last one. State is what makes a stream job able to do things a batch job cannot, and it is also what makes it able to run out of memory in a way that requires a plan rather than a restart.
Event time, processing time, and why the difference is the whole job
event time when the thing happened
(in the event's own payload)
processing time when your job saw it
the case that separates them
14:00:00.000 event created, order placed
14:00:00.500 event in a network partition
14:01:03.000 event arrives at the consumer
14:01:05.000 job processes it
processing-time window [14:00, 14:01)
the event is processed at 14:01:05
-> it lands in window [14:01, 14:02)
-> counted in the WRONG minute
event-time window [14:00, 14:01)
the event belongs to 14:00
-> counted correctly
This is not a rounding detail. With processing-time windows, the result depends on when the job ran, so replaying the same data produces different answers. That makes the job impossible to test, impossible to compare against a batch recomputation, and impossible to reason about when someone asks why yesterday’s number changed after a restart.
Everything else in a good stream job follows from preferring event time.
Watermarks: a guess about lateness, and it will be wrong
A watermark is a statement from the job: “I have seen all events with event time up to T.” It is generated per partition, and the job’s watermark is the minimum across all partitions.
partition 0 watermark: 14:10:00
partition 1 watermark: 14:09:55
partition 2 watermark: 14:10:02
------------------------------
job watermark: 14:09:55 <- the minimum
a window closes when the job watermark
passes its end
window [14:09, 14:10) closes when the
job watermark exceeds 14:10:00
everything after that is LATE
The minimum rule has a consequence that surprises people: the slowest active input holds back downstream event time. An idle partition can stall it indefinitely unless the watermark strategy marks that input idle; inputs with very different rates may need watermark alignment to prevent the fast side from growing downstream state without bound. A stuck source or rebalance can still hold progress, so monitor watermarks per operator rather than only the final result rate.
The generated strategy is the part that has to be configured well.
bounded out-of-orderness (the usual choice)
watermark = max(event_time seen)
- out_of_orderness_bound
bound = 5 seconds
-> tolerates 5s of lateness
-> a 6s-late event is dropped or
sent to a late-data side output
bound = 0
-> drops any event that arrives even 1ms
after its window should have closed
-> wrong, in the presence of retries
or rebalances or a GC pause
bound = 1 hour
-> almost nothing is late
-> every window waits an hour to close
-> results are an hour late
-> state for open windows grows a lot
idle-source handling (essential)
if a partition has no data for 30s, do NOT
hold back the watermark forever
-> mark it idle, exclude it from the minimum
The bound is a real trade between correctness and latency, and the right value comes from a measurement of actual lateness rather than a default. Instrument your lateness distribution — a histogram of current_watermark - event_time — and set the bound from a percentile of it. A bound of zero on a system with retries and rebalances is a configuration that will drop data; a bound of ten minutes on a system with 200ms typical lateness is state and latency you pay for nothing.
And the consequence of getting it wrong is not a crash:
a record arrives after its window closed
default: DROPPED
no error, no log, no metric
the number is wrong and you do not know
better: a late-data side output
a separate stream of everything that
arrived after its window closed
-> reconciled in a correction pass
-> at minimum, counted so you can alert
on the late rate
worst: fire an alert and increment a counter
nobody fixes it, but it is visible
Silent drops are the failure mode that makes stream processing untrustworthy. A late-data side output turns “the number is wrong and I do not know” into “there are 400 late records, they are in this stream, and I can reprocess them”. That single feature is the difference between a pipeline you can operate and one you are afraid to change.
State, and the fact that it is unbounded
windowed aggregation
state = open windows x keys
bounded by the watermark lag IF lateness
is bounded
keyed state (a join, a dedup set)
state = every key ever seen
UNBOUNDED unless something expires it
The second is the trap. A join between an event stream and a slowly-changing dimension table is fine. A join between an event stream and a table that grows forever is a job that will eventually OOM its task managers, and the standard fix is a TTL on the state — but a TTL on a join is a silent correctness decision, because it means “after the TTL, keys are forgotten and a late event for a forgotten key produces a null instead of a match”.
The state stores, briefly, because the choice has real operational consequences.
| Store | Where | Fits in memory | Use when |
|---|---|---|---|
| HashMap (heap) | TaskManager heap | Yes | Small state, and you control the size |
| RocksDB | Local disk, on-heap cache | Effectively unbounded | Large state, or unknown size |
| Embedded / native | Depends | Varies | You have specific latency needs |
Flink ships a heap-backed HashMapStateBackend and an EmbeddedRocksDBStateBackend; current defaults are configuration- and version-dependent, so do not infer the backend from job code. RocksDB keeps working state on local disk and supports state larger than heap, at the cost of serialization, compaction, and variable I/O latency. Bloom filters can help some point-lookup workloads but are not a universal requirement. Choose using measured state size, access pattern, checkpoint duration, and restore time.
The better answer to unbounded state is usually not a bigger store. A job that keeps per-customer history forever is a job doing a batch workload in streaming clothing. The standard fix is incremental aggregation: emit a running total per key per interval rather than retaining all history, so state is bounded by the number of keys rather than the number of events. groupBy(key).sum() with a short window emits periodic partial results that a downstream aggregation combines — bounded state, and the same final answer.
Checkpointing and savepoints are what let you resize, upgrade, and recover without losing state, and they are the operational reason to be careful about state in the first place: a savepoint is a consistent cut, and it is only useful if the state behind it is coherent.
Exactly-once: real, and narrower than it sounds
Flink’s exactly-once comes from two things together:
1. checkpointing
periodically snapshot all operator state
to a durable store, and record the
source offsets in the same snapshot
on failure: restore state AND offsets
from the checkpoint
-> reprocesses from that point
-> no gap, no duplicate in state
2. transactional / idempotent sinks
the sink participates in the checkpoint
-> a checkpoint commits the writes
-> a failure before the commit rolls back
The guarantee this gives is: exactly-once state changes, and exactly-once writes to sinks that participate in the transaction. Both conditions matter.
The failure that breaks it is writing to something outside the transaction — a direct database write from an operator function, an HTTP call to an external API, a print to a file that something else reads. Those are side effects at a point where Flink cannot roll them back, so a replay duplicates them. This is the single most common way a “exactly-once” job produces duplicates, and the fix is either a transactional sink for that destination or accepting the duplicate and making the destination idempotent.
With multiple inputs, downstream event time normally advances at the minimum watermark of the active inputs. Flink’s watermark alignment feature addresses a different problem: it can pause a source split that advances too far ahead so the fast input does not force unbounded buffering or state growth while waiting for the slow one. Idleness detection excludes inputs that genuinely have no data. Configure and test these separately.
Backpressure propagates, but Kafka lag still grows
producer -> Kafka partition -> consumer
consumer slower than producer
-> messages accumulate in the partition
-> the partition is a durable but retention-bounded buffer
-> consumers fall further behind
-> latency grows without bound
-> nothing errors
-> the data is simply late
where it actually shows up
consumer lag grows
end-to-end latency grows
eventually: retention expires data
the consumer catches up to a partition
whose earliest messages are GONE
Inside a Flink job, unavailable output buffers propagate backpressure upstream and eventually slow the Kafka source task. That does not automatically slow independent Kafka producers: they can continue appending while the consumer group accumulates lag. The broker log is bounded by time or size retention, so sustained lag becomes a recoverability deadline rather than an infinite queue.
The controls, in order of usefulness:
Measure consumer lag per partition and alert on rate of growth. A growing lag is the leading indicator. By the time latency is visible to a user, the lag has been growing for a while.
Rate-limit the producer when the consumer cannot keep up, so the backpressure is visible at the source rather than accumulating in the partition. A producer that backs off is a system with feedback; a producer that never backs off is a queue that grows until retention solves it by deleting data.
Partition enough. A partition can be read by only one consumer in a group, so consumer parallelism is bounded by partition count. Under-provisioned partitions is the most common reason a consumer group cannot scale, and the symptom is that adding consumers does nothing.
Buffer deliberately, with a bound. Sometimes a burst of messages should be absorbed. That is fine, but the buffer must have a bound and the overflow behaviour must be defined — shed, backpressure the producer, or drop with a counter. A buffer that grows until the node runs out of memory has defined the overflow behaviour as “the consumer dies”.
State store sizing, and the cluster maths
4 parallel tasks
25 GB keyed state
32 GB state limit per task
vs
8 parallel tasks
25 GB keyed state
16 GB per task
-> more parallelism, smaller budget each
-> rescaling redistributes state and
requires a savepoint
rescaling
1. stop with savepoint
2. change parallelism
3. restore
-> state is redistributed
-> a savepoint of 25 GB takes time
to write and time to restore
The operational consequence is that a stateful streaming job is expensive to change, and teams under pressure avoid changing it, which means parallelism drifts from what the traffic needs. Taking periodic savepoints specifically so a resize is a five-minute operation rather than a project changes that behaviour.
The other cluster maths worth checking: RocksDB state is on local disk, so the constraint is disk per node and eviction pressure, not just memory. A node whose disk is full will evict RocksDB blocks, and the job slows down rather than failing, which is easy to miss.
Failure stories worth testing
Delay one partition by two minutes
Confirm the whole job’s windows stall. This is the minimum-watermark effect and it should be demonstrated deliberately once, because it explains a whole class of “the job just stopped” reports.
Stop a partition entirely and keep producing
Confirm the job stalls forever without idle-source handling. This is the same failure, more severe.
Set out-of-orderness to zero and introduce a 500ms delay
Confirm records are dropped silently. Then confirm the late-data side output catches them. The silent version is the important one to see.
Replay the same data twice into Kafka
Confirm the job’s output is unchanged, assuming a transactional sink. Then break the sink by writing directly to a database from an operator and confirm the output duplicates. That contrast is the entire lesson about the guarantee’s boundary.
Kill a TaskManager during a checkpoint
Confirm recovery from the last completed checkpoint, and measure how much reprocessing happens. Checkpoint interval is a direct trade between recovery time and overhead.
Grow a keyed state set without a TTL
Confirm the task manager eventually fails or the disk fills. Unbounded state is a scheduled incident.
Add a TTL to a join’s state
Confirm what a late event for an expired key produces. The answer is a null, and that is a silent correctness change you just made.
Take a savepoint of a 25 GB state
Confirm the time it takes and that a restore is possible. This is the operation you need under pressure, so its duration should not be a surprise.
Scale from 4 to 16 instances
Confirm the state redistributes and nothing is lost. This requires a savepoint and it is a planned change, not an emergency one.
Saturate a consumer and watch lag
Confirm lag grows on the dashboard and that nothing errors. Then confirm the alert fires. Silent lateness is the failure.
Reduce Kafka retention below the time it takes the slowest consumer to catch up
Confirm data is deleted unread. This is the endpoint of unbounded lag and it is silent until someone notices a gap.
A production-ready architecture
200 producers
|
v
+-------------------------------+
| Kafka |
| - partitioned for consumer |
| parallelism |
| - 200+ partitions so |
| consumers can scale |
| - retention > worst-case |
| consumer catch-up time |
| - lag alerting per group |
+---------------+---------------+
|
+---------+---------+
| |
v v
+-----------+ +-----------+
| source | | source |
| operators| | operators |
+-----+-----+ +-----+-----+
| |
| watermark alignment: min
| (never emit until both
| inputs have passed)
v
+-------------------------------+
| Flink job |
| - event time, never |
| processing time |
| - generated watermarks with |
| measured lateness bound |
| - idle source handling |
| - late data -> side output |
| |
| state |
| - RocksDB for large keyed |
| state, on local disk |
| - TTL on anything unbounded, |
| and aware of what the TTL |
| means for correctness |
| - incremental aggregation |
| where history is not |
| actually needed |
+---------------+---------------+
|
checkpointing: state + offsets
together, on a schedule
|
v
+-------------------------------+
| transactional sink | the guarantee ends
| - participates in the | at anything written
| checkpoint commit | outside the txn
| - exactly-once here, or |
| idempotent destination |
+-------------------------------+
observability: lag, watermark lag,
lateness histogram, backpressure,
late-record rate, checkpoint duration
A sensible delivery checklist:
- Use event time when the result is defined by when events occurred. Processing time is valid for operational controls and arrival-time semantics, but replaying the same input may produce a different grouping and must be tested accordingly.
- Instrument the lateness distribution and set the out-of-orderness bound from a percentile of it. Zero is wrong on any system with retries or rebalances.
- Handle idle sources, or one silent partition stalls every window in the job.
- Send late records to a side output and count them. Silent drops make the pipeline untrustworthy and undebuggable.
- Bound state. Use incremental aggregation instead of retaining history, and put a TTL on anything that would otherwise grow forever — while knowing what the TTL means for a join.
- Use RocksDB for large state, and size the cluster on local disk as well as memory, because eviction shows up as latency rather than as an error.
- Checkpoint state and offsets together, and take periodic savepoints so a resize is a five-minute operation.
- Use a transactional sink, and audit every side effect written outside the transaction. That is where exactly-once ends.
- Alert on consumer lag growth rate, not on lag level. The rate is the leading indicator.
- Keep producer backpressure real: a producer that cannot push when the consumer cannot pull will lose data to retention, not error.
- Partition for the consumer parallelism you need, and confirm adding consumers actually increases throughput.
- Verify that retention exceeds the worst realistic catch-up time, and know by how much.
Common mistakes
| Mistake | What actually happens | Better decision |
|---|---|---|
| Processing-time windows | Results depend on when the job ran | Event time, always |
| Out-of-orderness bound of zero | Late records dropped silently | Bound from the measured lateness distribution |
| No idle-source handling | One silent partition stalls every window | Mark idle partitions, exclude from the min |
| Late records dropped with no counter | Numbers wrong, nobody knows | Side output plus a late-rate metric |
| Watermark not monitored | Job stalls and nobody knows why | Alert on watermark lag |
min alignment not understood |
A slow second input holds the whole job | Know which inputs gate which windows |
| Unbounded keyed state | TaskManager OOM, then disk | TTL, or incremental aggregation |
| TTL on a join with no analysis | Late events silently produce nulls | Understand what forgetting a key means |
| Heap state store for large state | OOM, or GC pauses that stall watermarks | RocksDB for anything big |
| No savepoints | Rescaling is a project, so it never happens | Periodic savepoints |
| Direct writes outside the transaction | Exactly-once silently becomes at-least-once | Transactional sink, or idempotent destination |
| Checkpoint interval too long | Minutes of reprocessing on failure | Balance interval against overhead |
| Lag alerting on level, not growth | You are paged after the damage | Alert on rate of growth |
| No producer backpressure | Lag grows until retention deletes data | Rate-limit producers, bounded buffers |
| Too few partitions | Adding consumers does nothing | Partition for target parallelism |
| Retention below worst catch-up time | Data silently deleted unread | Retention must exceed it, with margin |
| No lateness histogram | The bound is a guess forever | Measure and derive it |
| Assuming a restart resumes exactly | A restore with a stale savepoint loses state | Savepoint hygiene, verify before replacing |
| Unbounded join with a growing table | Job grows until it dies | Bound the right side, aggregate upstream |
| Ignoring GC pauses on the task | Long pauses stall watermark advancement | Size heap, move state off-heap |
| One giant job for everything | A single bug takes out all processing | Split by domain, isolate blast radius |
The complete story in one minute
A streaming job is correct or incorrect based on whether its time semantics match the business contract. Processing-time windows deliberately describe arrival and execution time; they are testable, but replay is not deterministic with respect to original event time. Event-time windows are required when the result is meant to describe when events occurred, and they bring an explicit lateness and correction policy.
Watermarks estimate event-time progress, and the bound is a real latency-versus-completeness trade. A record behind the watermark is late; whether it updates a retained window, goes to a side output, or is dropped depends on allowed lateness and the operator. Set the policy from a measured lateness distribution, monitor late-record outcomes and operator watermarks, mark idle inputs, and consider alignment where fast inputs would otherwise inflate state.
Late records must not vanish quietly. A side output for everything that arrived after its window closed turns an unknowable wrong number into a countable, correctable set of records, and a late-rate metric is the cheapest correctness signal in the system.
State is the other half. Unbounded keyed state is a job that will eventually exhaust some resource, and the answer is usually a bounded contract or incremental aggregation rather than a bigger heap. Heap state is fast and GC-sensitive; embedded RocksDB trades CPU and I/O for state beyond heap. Size local disk and remote checkpoint storage, measure restore time, and rehearse savepoint compatibility before calling a resize routine.
Exactly-once is real and narrower than it sounds: exactly-once state changes and writes to sinks that join the transaction. The moment you write to a database from an operator or call an external API, the guarantee ends, and that is the most common way a correctly-configured job still produces duplicates.
event time + watermarks, bound measured not guessed
idle sources handled, late data to a side output
state bounded deliberately, savepoints taken
exactly-once ends at the transaction boundary
The hard part was never reading the stream. It was deciding when to stop waiting for data that might still be on its way.

