Event Time, Watermarks & Checkpoints
Part 0. Start here
Counted twice, and nothing was added twice
It is Super Bowl night, and an ad network is counting 100,000 clicks a second. One worker of the stream job that counts them dies. The job restores its last checkpoint and replays what came after it. Every sink write is an absolute SET, never an ADD, and the job's state is exactly-once.
Yet one ad's hour row now counts a click twice: once in fast_valid (the minute it belongs to) and once in late_valid (the late total for its hour). Nothing was added twice. The replay simply decided differently about that click.
The numbers
The replay is not small, and it is not rare: every crash, rescale and upgrade causes one.
- How much is replayed. The job checkpoints every 30 s and takes about 50 s to restart (the ad-click loop's assumption). If the last checkpoint completed on time, a crash replays up to 30 + 50 = 80 s of clicks; at the Super Bowl rate that is 80 × 100,000 = 8,000,000 clicks. On an average day the loop's trace replays 70 s × 11,574 clicks/s ≈ 810,000. When checkpoints stall, the replay is longer. The small example on this page replays 40 s of input after a 50 s restart: 40 + 50 = 90 s.
- A second harm hides in the first 200 ms. After a restore, the job's watermark (its sense of "how far event time has got", Part 2) starts from nothing, and it is recomputed only every 200 ms. The restored job reads its backlog at its catch-up rate, about 64,000 clicks/s (the loop's 32 KPUs × ~2,000/s), so 64,000 × 0.2 = 12,800 clicks are judged with no watermark at all. If 2% of them belong to minutes that closed before the checkpoint, that is 12,800 × 0.02 = up to 256 clicks that each re-open a finished minute (several can share a row), and each such minute's row is overwritten with a count of about 1 (say 812 → 1). If one quiet shard holds the watermark down for the whole 30 s idle timeout, it is 64,000 × 30 × 0.02 = 38,400, up to about 38,000 rows. (An estimate at the loop's normal-day sizing; on a Super Bowl night at 64 KPUs the rate doubles, and so do these numbers. Not a measurement.)
- What the cleanup costs. The loop's cleanup after a restore rewrites about 300,000 keys at 2,800 write units a second: 300,000 ÷ 2,800 ≈ 107 s, about 2 minutes (Part 8 shows how to make it smaller).
The job restored state and stream positions from the same checkpoint and replayed the same records. Every sink write sets an absolute value. How can the replay still count a click twice, and what would you change?
The big picture
Synthesizing vector architecture diagram...
What to notice: the watermark is computed at the sources and flows with the data; lateness is decided in one place, before the window; every box with state, and the shards' read positions (dashed arrows), is saved in each checkpoint. The sinks are outside the checkpoint, which is why Part 8 exists.
What you'll be able to do after this page
- Tell event time from processing time, and name the three questions every streaming design answers (Part 1).
- Compute a job's watermark across shards, and say what idleness buys and costs (Part 2).
- Choose tumbling, sliding or session windows, and build hour rows from closed minutes without double counting (Part 3).
- Choose between allowed lateness and a side output, and make early results safe for their consumers (Part 4).
- Say what lives in keyed state, why event-time timers (not TTL) clear dedup state, and what a stuck watermark does (Part 5).
- Walk a checkpoint barrier through the job, aligned and unaligned, and say what backpressure does to it (Part 6).
- List what a restore brings back and what it doesn't, and explain why a replay can decide differently (Part 7).
- Make a sink safe after a restore: per-writer tags, the flush, the closed-row condition, or a transactional sink with timeouts that fit (Part 8).
- Rescale and upgrade a job without losing or duplicating output (Part 9).
- Compare continuous streams, micro-batches and a batch recount on equal terms (Part 10).
- Trace one click through every layer, and name the four numbers to watch (Part 11).
- Map all of it to AWS (Part 12).
You may have arrived from step 2.2 ("When Is a Minute Complete?") or 2.4 ("A Worker Crash Replayed a Minute and Double-Billed") of the ad-click aggregation loop, step 2.5 of the metrics and alerting loop ("Late and Out-of-Order Samples Are Rejected"), or step 2.5 of the YouTube loop ("A Counter Row per Video Melts on Viral Videos"). This page is the "why" behind all of them.
Part 1. The core mental model: two clocks
A click happens at one moment and reaches the job at another. Everything on this page follows from keeping those two moments apart.
Two clocks
| Clock | What it measures | Who sets it |
|---|---|---|
| Event time | When the click happened | The record: a timestamp set at ingest |
| Processing time | When the job sees the click | The job's wall clock |
Take one click from our example, c15. A phone clicked an ad at 10:13:40, lost signal, and uploaded the click at 10:15:04. The job reads it at 10:15:04.
| Counted by | c15 goes into minute | Why |
|---|---|---|
| Event time | 10:13 | That's when the user clicked |
| Processing time | 10:15 | That's when the job happened to read it |
c15 was clicked at 10:13:40 and read at 10:15:04. Which minute should it count in, and who decides that the minute is finished?
Why processing time moves clicks
Suppose the job falls 2 minutes behind (a slow sink, a restart). With processing-time minutes, every click read during the catch-up lands in the minute the job happened to read it, not the minute it happened. The same input gives different counts depending on the job's health. In Part 3 we count the same clicks both ways: processing-time minutes give one answer before the crash and another after it, and neither is the event-time answer.
Three questions every streaming design answers
| Question | Answered by | Where on this page |
|---|---|---|
| Which window does a record belong to? | Its event time | Parts 1, 3 |
| When is a window complete? | The watermark | Part 2 |
| What if a record comes after that? | A lateness policy: allowed lateness, a side output, a batch recount, or rejection past the acceptance horizon | Part 4 |
Until the third question's answer has run out, every result is provisional.
Not this watermark
The word "watermark" means other things in the loops. They share the word, not the mechanism:
- Kafka's high watermark: the offset up to which a partition's messages are replicated (Message Queues vs Event Streams). Nothing about event time.
- Read-receipt watermarks in chat: a per-conversation "read up to here" cursor.
- A scheduler's watermark: the last fully processed time bucket, advanced by a job.
- A watermark on a video: a mark on the picture.
On this page a watermark is always an event-time progress marker.
The minute we follow
Real jobs read billions of clicks, too many to watch. So one small story runs through the whole page: one minute of clicks from three shards. Clicks on two ads, ad_7 and ad_9, of one campaign, for the event-time minute 10:14, read by one job from a Kinesis stream with three shards, S1, S2 and S3. Every trace on this page was produced by running a private reference implementation of this job, not worked out by hand.
| Setting | Tiny example | At real scale (the ad-click loop) |
|---|---|---|
| Source | Kinesis stream, 3 shards; a random partition key per record | 64 shards, enhanced fan-out |
| Event time | clamp(device_time, impression_time, arrival_time), set at ingest; each click's impression is 10 s before it | Same |
| Watermark per shard | Highest event time read − 5 s allowance (− 1 ms more, Part 2) | Same 5 s |
| Watermark emission | Every 200 ms of processing time, on the 200 ms marks of the wall clock | Same |
| Job watermark | Minimum over active shards; never goes back | Same |
| Idleness | A shard with no records for 30 s is left out of the minimum | 30 s |
| Windows | 1-minute tumbling, by event time, keyed by ad; allowed lateness 0; late clicks to a side output | Same |
| Early firings | Running totals to Valkey once a second | Same |
| Late path | Late totals per (ad, event hour): SET late_valid on the hour row, at most every 10 s per row (a new row at once) | Same |
| Hour rows | fast_* = the sum of the hour's closed minutes | Same |
| Pre-aggregation | 500 ms buffers, before the shuffle by ad | 64 subtasks |
| Dedup | Keyed state by click_id, cleared by an event-time timer at impression + 2 h | Same |
| Checkpoints | Every 30 s: c39 at 10:14:00, c40 at 10:14:30, c41 at 10:15:00, c42 at 10:15:30 | 30 s |
| Restart | 50 s from crash to running (an assumption) | Same |
| Parallelism | 2, max parallelism 128; rescaled to 4 in Part 9 | 32 → 64 KPUs |
| Sinks | DynamoDB ad_metrics (absolute SETs, each writer tags its own attributes); Valkey for running totals | Same, 2,800 WCU |
| Job generation | job_gen 7 (8 in Part 9) | 41 in the loop |
Before the story starts, each shard has already read some clicks; among them are six clicks on ad_7 in minute 10:13 (so minute 10:13 will close with 6), and S3's last one has event time 10:13:52. The clicks:
| Click | Shard | Event time | Read at | Ad | Note |
|---|---|---|---|---|---|
c1 | S1 | 10:14:05 | 10:14:05.3 | ad_7 | |
c2 | S2 | 10:14:12 | 10:14:12.5 | ad_9 | |
c3 | S3 | 10:14:20 | 10:14:20.4 | ad_7 | S3's last record for a while |
c4 | S1 | 10:14:31 | 10:14:31.2 | ad_7 | |
c5 | S2 | 10:14:40 | 10:14:40.3 | ad_7 | |
c2′ | S1 | 10:14:12 | 10:14:45.0 | ad_9 | The SDK resent c2: same click_id |
c7 | S1 | 10:14:58 | 10:14:58.4 | ad_9 | |
c9 | S2 | 10:15:02 | 10:15:02.5 | ad_9 | Minute 10:15 |
c15 | S2 | 10:13:40 | 10:15:04.0 | ad_7 | A phone upload, 84 s after the click |
c8 | S2 | 10:14:55 | 10:15:05.0 | ad_7 | Out of order, but its minute is still open |
c10 | S1 | 10:15:06 | 10:15:06.2 | ad_7 | |
c11 | S2 | 10:15:07 | 10:15:07.1 | ad_9 | |
c12 | S3 | 10:14:57 | 10:15:20.0 | ad_9 | A phone upload on the quiet shard |
c13 | S1 | 10:15:31 | 10:15:31.1 | ad_7 | |
c14 | S1 | 10:16:45 | 10:16:45.2 | ad_9 | Clicked after the crash, once the job runs again |
c16 | S2 | 10:16:50 | 10:16:50.3 | ad_7 | Same |
(Click numbers follow a first draft of the story, so c6 isn't used; c14 and c16 are live traffic after the restart.)
The story in seven beats: the clicks arrive out of order (Parts 1 and 2), a shard goes quiet (Part 2), the minute closes (Part 3), two clicks come too late (Part 4), the checkpoints are taken (Parts 5 and 6), a worker crashes and the replay decides differently (Parts 7 and 8), and the job is rescaled and upgraded (Part 9). Each Part shows only its slice; the full table is in Part 13.
What to remember from Part 1
- Event time decides where a record belongs; processing time only says when we saw it.
- A window is complete when the watermark says so, not when the wall clock does.
- Every streaming answer is provisional until its lateness policy has run out.
Part 2. Watermarks
A window must close sometime, so its row can be written and its memory freed. Closing by the wall clock fails after any backlog; waiting forever never writes anything. We need a way for the data itself to say "event time has got this far".
What a watermark promises
A watermark is a timestamp that flows through the job with the records and says: "event time has reached X; I don't expect anything older." It is a promise the job makes to itself, based on what it has read. When the watermark passes the end of a window, the window closes.
One shard
Each shard's reader tracks the highest event time it has read and, every 200 ms (Flink's default pipeline.auto-watermark-interval), emits a watermark of highest event time − allowance. Records in a shard are a little out of order (phones, retries, batching), and the allowance (5 s here) is how much disorder we wait for. This strategy is called bounded out-of-orderness.
Flink subtracts one more millisecond: its generator emits highest − allowance − 1 ms. So after c1 (10:14:05) on S1, S1's watermark is 10:13:59.999. The millisecond makes the arithmetic come out cleanly: a window covers [start, end), its last millisecond is end − 1 ms, and it fires when the watermark reaches that last millisecond. So minute 10:14 closes exactly when every active shard has read an event time of 10:15:05 or later.
Three shards and the minimum
A job reading three shards can't take any single shard's watermark: S1 may be at 10:14:53 while S3 still has clicks from 10:14:15 in flight. So the job's watermark is the minimum over the shards, and it only moves forward:
From the reference implementation (watermarks shown with Flink's − 1 ms; "tick" is the 200 ms emission that makes a change visible):
| # | Wall clock | Event | S1 | S2 | S3 | Job watermark |
|---|---|---|---|---|---|---|
| 1 (start) | 10:14:05.3 | c1 read | 10:13:59.999 | 10:13:51.999 | 10:13:46.999 | 10:13:46.999 (S3) |
| 2 | 10:14:20.4 | c3 lifts S3; minute 10:13 fires with 6 | 10:13:59.999 | 10:14:06.999 | 10:14:14.999 | 10:13:59.999 (S1) |
| 1 (end) | 10:14:40.4 | tick after c5 | 10:14:25.999 | 10:14:34.999 | 10:14:14.999 | 10:14:14.999 (S3) |
| 4 | 10:14:50.8 | S3 marked idle | 10:14:25.999 | 10:14:34.999 | idle | 10:14:25.999 (S1) |
| 5 | 10:14:58.4 | c7 read | 10:14:52.999 | 10:14:34.999 | idle | 10:14:34.999 (S2) |
| 7 | 10:15:02.6 | tick after c9 | 10:14:52.999 | 10:14:56.999 | idle | 10:14:52.999 (S1) |
| 10 | 10:15:06.2 | c10 read | 10:15:00.999 | 10:14:56.999 | idle | 10:14:56.999 (S2) |
| 11 | 10:15:07.2 | tick after c11 | 10:15:00.999 | 10:15:01.999 | idle | 10:15:00.999: minute 10:14 closes |
Two things to see. At event 2, minute 10:13 was held open by S3 alone (its last record was 10:13:52, so its watermark was 10:13:46.999); c3 lifted S3, the minimum became S1's 10:13:59.999, which is minute 10:13's last millisecond, and the minute fired at once. And at event 11, minute 10:14 closed at 10:15:07.2, the first 200 ms tick after c11 was read at 10:15:07.1: min(10:15:00.999, 10:15:01.999) = 10:15:00.999, past 10:14:59.999.
Synthesizing vector architecture diagram...
Snapshot W2, 10:15:07.2: S3 is idle, so the minimum covers S1 and S2 only; it passes minute 10:14's last millisecond and the minute closes with 5 and 2. (W1, the checkpoint about 6 seconds earlier, is drawn in Part 6.)
Out of order is not late
c8 was clicked at 10:14:55 and read at 10:15:05.0, 3 s behind c9 (10:15:02), which S2 had already read. It is out of order. But it is not late: a record is late only when the watermark has already passed the end of its window. At 10:15:05.0 the job watermark is 10:14:52.999, and c8's minute ends at 10:15:00, so c8 joins minute 10:14: ad_7 goes from 4 to 5, and Valkey shows 5 at the next once-a-second firing (10:15:06).
A quiet shard
S3 carries phone uploads, and after c3 it is silent. Without help, its watermark would sit at 10:14:14.999 and hold the whole job there.
S3 has been silent for 40 seconds. Should minute 10:14 wait for it?
Flink's withIdleness works this way, and its Kafka connector generates watermarks per partition and combines them the same way. (When every input is idle at once, recent Flink versions emit the highest of their watermarks; our story never lets that happen, and your design shouldn't rely on it.)
The clamp
Event time comes from the device. A phone whose clock says 2031 would lift its shard's watermark by years. The ad-click loop clamps at ingest: event_time = clamp(device_time, impression_time, arrival_time) (impression ≤ event time ≤ arrival), and keeps the raw value for fraud analysis.
From the reference implementation, c7's phone reporting a time in 2031, no crash:
| Without the clamp | With the clamp | |
|---|---|---|
c7's event time | 2031 | 10:14:58.4 (its arrival time) |
| S1's watermark | 2031: S1 never holds the minimum again | 10:14:53.399 |
Minute 10:14 for ad_9 | 1 (c7 sits in a window in 2031 that won't fire for years) | 2 |
| When S2 and S3 go idle (10:15:37.6 and 10:15:50.4) | The job watermark jumps to 2031: minute 10:15 closes at once, and every click read afterwards, on any shard, is late | The job watermark is S1's 10:15:25.999 |
Watermarks are guesses
A watermark is a heuristic. The Dataflow paper (Akidau et al., 2015) names both ways it fails: too fast (late data exists after it passed) and too slow (one slow source holds everyone back). A bigger allowance trades delay for completeness, but the loop's clicks that are really late are minutes late, so past a few seconds a bigger allowance buys almost nothing. Flink can also align watermarks: pause reading sources that run too far ahead of the slowest one, so a fast source doesn't fill memory with windows that can't close yet.
Watching the watermark
Watermark lag = wall clock − job watermark. At 10:15:07.2 it is 10:15:07.2 − 10:15:00.999 ≈ 6.2 s: the 5 s allowance plus about a second of disorder and emission timing. If it climbs, something is stuck: a shard with no data and no idleness, a stalled reader, a stuck operator. The ad-click loop publishes it as data_complete_through and alarms when the lag exceeds 60 s for 5 minutes. It is a different number from consumer lag (how far behind the stream's tip the reader is): a job can be fully caught up with a stuck watermark, or far behind with a healthy one.
What to remember from Part 2
- The job's watermark is the minimum over active sources, each its highest event time minus the allowance, and it never goes back.
- Idleness keeps a silent source from freezing the job, and makes that source's next records likely to be late.
- Clamp device time before it reaches the watermark; alarm on watermark lag as well as consumer lag.
Part 3. Windows
Records now carry event times and the job knows how far event time has got. Next: what to count over. Reports, rolling fraud rules and user sessions need three different shapes, and each has a trap.
Tumbling: the minute closes
A tumbling window is a fixed, non-overlapping bucket: [10:14:00, 10:15:00), then [10:15:00, 10:16:00). Each record belongs to exactly one, found by rounding its event time down to the minute.
textTUMBLING WINDOW (per ad, size 1 minute, allowed lateness 0) 1. window = floor(event time / 1 minute) 2. add the record to that window's running count -- one number per ad and window 3. register an event-time timer at the window's last millisecond 4. when the watermark reaches the timer: emit the count (SET fast_valid), add it to the hour roll-up, delete the window's state
Event 11 from the reference implementation:
| At | Watermark | Window | Count | Written |
|---|---|---|---|---|
| 10:15:07.2 | 10:15:00.999 | ad_7, 10:14 | 5 (c1, c3, c4, c5, c8) | SET fast_valid = 5 on ad_7 / M10:14 |
| 10:15:07.2 | 10:15:00.999 | ad_9, 10:14 | 2 (c2, c7; c2′ was a duplicate) | SET fast_valid = 2 on ad_9 / M10:14 |
An aggregate function (a count, a sum) keeps one number per window; a function that needs every record of the window at the end (Flink's ProcessWindowFunction alone) keeps them all. At 100,000 clicks a second, that difference is the state size.
The hour from closed minutes
Hour rows must equal the sum of their minutes. The obvious design is a second, 1-hour window.
Why not compute the hour row with its own 1-hour window? It's simpler.
Sliding windows
A sliding window has a size and a slide: 60 s, every 10 s. A fraud rule like "more than 5 clicks from one device in a minute" needs it, because tumbling minutes miss a burst that straddles a boundary.
From the reference implementation, device d_5 clicks ad_7 at 10:14:40, :48, :55, 10:15:02, :09 and :15 (six clicks in 35 s):
| Windows | What they see | Flag (more than 5)? |
|---|---|---|
| Tumbling minutes | 10:14 has 3, 10:15 has 3 | No |
| 60 s sliding every 10 s | [10:14:20, 10:15:20), [10:14:30, 10:15:30) and [10:14:40, 10:15:40) each hold all 6 | Yes: the first to fire is [10:14:20, 10:15:20) |
The cost: each click belongs to 60 ÷ 10 = 6 windows, so state and output are 6 times larger (or the job keeps a count per 10 s slice and adds the last 6 slices).
Session windows
A session window groups a key's activity separated by gaps. Each record opens [t, t + gap); overlapping windows of one key merge; the session fires when the watermark passes its end. Device d_8, gap 30 s, from the reference implementation:
| After the click at | Sessions |
|---|---|
| 10:14:05, 10:14:20 | [10:14:05, 10:14:50) |
| 10:15:00 | [10:14:05, 10:14:50) and [10:15:00, 10:15:30): a 40 s gap |
| 10:14:40 (arrives out of order) | Its window [10:14:40, 10:15:10) overlaps both, so they merge: [10:14:05, 10:15:30) |
A late or out-of-order element can join two sessions into one. Sessions can stay open as long as activity continues.
Synthesizing vector architecture diagram...
What to notice: tumbling minutes never overlap; three sliding windows each contain d_5's whole burst, and slide A (10:14:20 to 10:15:20) is the first to fire with 6; d_8's two sessions become one after the out-of-order click at 10:14:40.
Processing-time windows
A processing-time window buckets by the job's wall clock. The same clicks, from the reference implementation:
| Minute (by read time) | Before the crash | After the restore |
|---|---|---|
| 10:14 | ad_7 4, ad_9 2 | (closed before the crash) |
| 10:15 | ad_7 4 (c15, c8, c10, c13), ad_9 3 (c9, c11, c12) | Lost with the crash: its state was never checkpointed and it never fired |
| 10:16 | Every replayed click, plus c14 and c16: ad_7 5, ad_9 4 |
The event-time answer for minute 10:14 is ad_7 5 and ad_9 3 (Part 8), with c15 in 10:13. Processing-time counts depend on lag and on crashes, so they are fine only for metrics about the pipeline itself (records read per minute), never for billing.
Pre-aggregation without losing clicks
The ad-click loop's hot ad gets 40,000 clicks a second, and the shuffle by ad would send them all to one subtask. So each upstream subtask pre-aggregates: it adds up its own clicks per (ad, minute) for 500 ms and forwards one partial sum. The window adds partials just as it would add clicks. Two rules keep this exact (a third, at the checkpoint barrier, is in Part 6):
textPRE-AGGREGATE (per upstream subtask, before the shuffle by ad) on each click: 1. if its window ended at or before the current watermark: -- decided per click send the click, alone, to the late path; never fold it into a partial 2. else add +1 to the buffered partial for (ad, window) every 500 ms: forward the buffered partials on a new watermark: 3. forward every buffered partial FIRST, then forward the watermark
Rule 1: a partial is a sum, so once a late click is folded in, nobody downstream can separate it. Rule 3: without it, a partial can reach the window after the watermark that closed the window.
Branch P1 (from the reference implementation): c8 is read at 10:15:06.9 instead of 10:15:05.0. It is on time (its minute ends after the watermark 10:14:56.999), so its +1 sits in the 500 ms buffer when c11 lifts the watermark at 10:15:07.2.
| Without rule 3 | With rule 3 | |
|---|---|---|
| 10:15:07.2 | Watermark forwarded; minute 10:14 fires with ad_7 = 4 | Buffer forwarded first; minute 10:14 fires with ad_7 = 5 |
| 10:15:07.4 | The buffered +1 arrives after the minute closed: dropped at the window | |
c8 | Lost: on time by the pre-aggregator's check, dropped by the window, and in no late total | Counted |
What to remember from Part 3
- Tumbling for reports, sliding for rolling rules, sessions for activity; all by event time.
- Build hours from closed minutes, each added once, or a click late for its minute is counted twice.
- A pre-aggregator must decide lateness per click and forward its partials before forwarding a watermark.
Part 4. Late data and early answers
Minute 10:14 is closed and written. Two clicks came too late for their minutes, and dashboards wanted numbers long before the minute closed. Both are about the same thing: results that change after we first showed them.
Two late clicks
| # | At | Click | Its minute ended | Watermark it was judged against | Result |
|---|---|---|---|---|---|
| 8 | 10:15:04.0 | c15 (10:13:40) | 10:14:00 | 10:14:52.999 (about 53 s past its minute's end) | Late → late total (ad_7, hour 10) = 1: SET late_valid = 1 on ad_7 / H10, tagged late_ckpt = 41 |
| 12 | 10:15:20.0 | c12 (10:14:57) | 10:15:00 | 10:15:00.999 (its minute closed 12.8 s earlier) | Late → late total (ad_9, hour 10) = 1 on ad_9 / H10, tagged 41 |
Both late totals are new attributes written after checkpoint c41. Remember that; it is the whole of Part 8. This is snapshot W3, 10:15:20.0:
| Item | Attributes |
|---|---|
ad_7 / M10:13 | fast_valid 6, fast_ckpt 39 |
ad_7 / M10:14 | fast_valid 5, fast_ckpt 41 |
ad_9 / M10:14 | fast_valid 2, fast_ckpt 41 |
ad_7 / H10 | late_valid 1, late_ckpt 41 |
ad_9 / H10 | late_valid 1, late_ckpt 41 |
Allowed lateness or a side output
A late record must go somewhere. Flink's default allowed lateness is 0: a late record is dropped, unless you route it. Two ways to keep it:
- Allowed lateness: keep each window's state for a while after it closes; a late record inside that time updates the window and fires it again. Flink's documentation says late firings "should be treated as updated results of a previous computation".
- A side output: a second output of the operator for records late for their window (Flink's
sideOutputLateData). A separate operator keeps late totals; the minute row is written once.
Branch A (from the reference implementation), allowed lateness 60 s, no crash:
| Allowed lateness 60 s | Side output (our design) | |
|---|---|---|
c15 at 10:15:04.0 | Minute 10:13 is kept until the watermark passes 10:14:59.999, so it fires again: ad_7 / M10:13 6 → 7 at 10:15:04.5 | Late total ad_7 / H10 = 1; M10:13 stays 6 |
c12 at 10:15:20.0 | Minute 10:14 fires again: ad_9 / M10:14 2 → 3 at 10:15:20.5 | Late total ad_9 / H10 = 1; M10:14 stays 2 |
| State held | Every window 60 s longer | One late total per (ad, hour) |
| Rows | Change after they were first written | Minute rows written once; late totals change |
| What consumers must do | Treat a second firing as an update, by key | Add on-time and late: dashboards show fast + late |
| Where the numbers are reconciled | Nowhere: one number per minute | Per hour: final = fast + late − A + B in the loop's recount |
| Clicks later than the allowance | Still need a path (drop or side output) | Same path, whatever the delay |
The ad-click loop picks the side output because its late clicks come minutes late: 45 minutes of allowed lateness would keep 45 times the window state.
The acceptance horizon and the recount
Late can't mean "forever". The acceptance horizon (2 hours after the impression in the ad-click loop) is the point after which a record is rejected as stale; Idempotency & Effectively-Once Processing (Part 8) states the rule and ties dedup expiry to it. Inside the horizon, the stream's numbers are provisional, and an exact batch recount from the lake writes final_* attributes once the hour can no longer change. The S3 loop's metering does the same with a day: a late batch "marks that hour dirty", and "hours are final after a day".
Synthesizing vector architecture diagram...
What to notice: three writers touch the hour row, each only its own attributes, and each streaming writer with its own tag. That separation is what lets Part 8 clean up one writer's values without touching another's.
Early firings
The ad-click dashboards need numbers within 3 seconds, not when the minute closes. A trigger decides when a window emits; here it fires every second of processing time with the running total (an early firing), again at the watermark, and the side output handles anything later. From the reference implementation, the running totals Valkey received for minute 10:14 up to 10:14:41:
| At | ad_7 | ad_9 |
|---|---|---|
| 10:14:06 | 1 | |
| 10:14:13 | 1 | |
| 10:14:21 | 2 | |
| 10:14:32 | 3 | |
| 10:14:41 | 4 |
Your dashboard sums the per-ad early results into a campaign total. Why does it show 11 when the campaign has 5 clicks?
What to remember from Part 4
- Late records either update the window (allowed lateness) or go to a separate path (side output); choose per consumer.
- Early results are provisional: a consumer must overwrite them by key or receive retractions.
- Past the acceptance horizon a record is rejected; before it, a batch recount can still correct the answer.
Part 5. Keyed state and timers
Windows, dedup, roll-ups and late totals all need memory that belongs to a key, survives restarts, and is cleared at the right moment. How it is cleared turns out to matter as much as how it is kept.
State that belongs to a key
Keyed state is memory scoped to the current key: when the job processes a click for ad_7, it sees only ad_7's entries. It is partitioned by key (Part 9 shows how), so each entry lives on exactly one subtask. Operator state belongs to a subtask instead, whatever keys pass through it.
| State | Kind | Key | In our example |
|---|---|---|---|
| Dedup ids | Keyed | click_id | c1 … c7 (and the six from minute 10:13) |
| Window counts | Keyed | ad, window | ad_7 10:14 = 4, ad_9 10:14 = 2 at c41 |
| Hour roll-up | Keyed | ad, hour | ad_7 hour 10: minute 10:13 = 6 |
| Late totals | Keyed | ad, event hour | ad_7 hour 10 = 1 after c15 |
| Pre-aggregation buffer | Operator | (subtask) | Partials for 500 ms (Part 6: flushed at each barrier) |
| Watermark floor | Operator | (subtask) | 10:14:34.999 at c41 (Part 7) |
Event 3: at 10:14:45.0, c2′ arrives on S1 with c2's click_id. Dedup state already holds c2, so c2′ is dropped as a duplicate and written to the exceptions log. Minute 10:14 for ad_9 stays at 1 until c7.
Timers on two clocks
A timer asks the job to call back at a time. Event-time timers fire when the watermark passes them: a window closes, a dedup id expires. Processing-time timers fire on the wall clock: the once-a-second early firing, the 500 ms buffer.
Dedup ids are cleared by an event-time timer at impression + 2 h, the acceptance horizon. For c1: impression 10:13:55 + 2 h = 12:13:55, in event time. However far behind the job falls, the timer can't fire before the watermark reaches 12:13:55, so every copy the horizon still accepts finds its original in state.
Why TTL can't be the dedup window
Flink also has state TTL: an entry expires a set time after it was last written. But Flink's TTL counts processing time only, and expired entries are removed lazily on read and in the background (for RocksDB, during compaction). So a 2-hour TTL forgets a click id 2 hours after the job saw it, however far behind the job is; after a long outage, copies that the horizon still accepts find nothing and are counted twice. Idempotency & Effectively-Once Processing (Part 8) traces exactly that. Keep a much longer TTL only as a guard (the loop uses 26 h) in case the watermark gets stuck.
Shard S3's producer stops for 3 hours, and idleness is off. What happens to window state, dedup state and data_complete_through?
Where state lives
Keyed state bigger than memory lives in RocksDB on the worker's local disk: reads go through its caches and Bloom filters, and snapshots copy its immutable files (the LSM-Trees & Compaction loop primitive explains how RocksDB stores and compacts them). The ad-click loop's numbers:
| Size | |
|---|---|
| Typical state (mostly dedup ids: 83.3M × ~60 B) | ~5 to 6 GB |
| Worst case (2 hours at full Super Bowl peak) | ~43 GB |
| Running application storage: 32 KPUs × 50 GB | 1.6 TB |
That local disk is shared by the state itself, RocksDB compaction, and the files each checkpoint stages for upload.
What to remember from Part 5
- Keyed state lives with the key, on local disk, and is restored with the job.
- Event-time timers fire on the watermark; TTL fires on the wall clock and can't define a dedup window.
- A stuck watermark stops timers and state grows: alarm on watermark lag, not just on consumer lag.
Part 6. Checkpoints
The job's state is spread over many subtasks, all changing thousands of times a second, and its input is three shards moving independently. To restart after a crash we need a snapshot of all of it that matches one moment of the input, taken without stopping the job.
A consistent cut
A checkpoint saves every operator's state and every source's read position, such that the state reflects exactly the records before those positions. After a crash, the job loads both and replays from the positions; every record then affects the state exactly once. This is the same rule as a write-ahead log's "log first, then act" (Write-Ahead Log, fsync & Group Commit): the state and the position that produced it are saved together, and the metrics loop's ingesters follow it by hand (upload the block, then checkpoint the stream position). Flink's method comes from the Chandy-Lamport snapshot algorithm.
Barriers, step by step
textCHECKPOINT n (aligned) 1. the coordinator tells every source: inject barrier n each source saves its read position and sends barrier n downstream, in line with its records 2. an operator with several inputs waits until barrier n has arrived on ALL of them, holding back records from inputs that already delivered it 3. then it snapshots its state and forwards barrier n 4. snapshots are uploaded to checkpoint storage in the background 5. when every task has acknowledged, checkpoint n is complete; operators are notified
Synthesizing vector architecture diagram...
What to notice: records that arrived on S1 after its barrier wait until S2's barrier arrives, so the snapshot holds everything before barrier 41 on both inputs and nothing after it. In our example both barriers arrive at 10:15:00.0; under load, the wait for the slower input is the alignment time.
Event 6, checkpoint c41, from the reference implementation:
| Part of the cut | Value |
|---|---|
| Barriers injected | 10:15:00.0: after c7 on S1, c5 on S2, c3 on S3 |
| Source positions | S1 after c7, S2 after c5, S3 after c3 |
| Window counts | ad_7 10:14 = 4 {c1, c3, c4, c5}; ad_9 10:14 = 2 {c2, c7} |
| Dedup ids | c1 … c7 (and minute 10:13's six), with their timers |
| Late totals | None yet |
| Job watermark at the barrier | 10:14:34.999 = min(S1 10:14:52.999, S2 10:14:34.999), S3 idle |
| Watermark floor (Part 7) | 10:14:34.999 |
| Complete | 10:15:01.5 |
Synthesizing vector architecture diagram...
Snapshot W1, 10:15:01.5: everything the restore in Part 7 will come back to. Note what is not in the checkpoint storage box: the watermarks themselves, the idle flag on S3, and anything in the sink.
Buffers outside state: flush at the barrier
The pre-aggregator's 500 ms buffer is memory, not state. If a barrier passes while a partial sits in it, the snapshot has neither the partial nor the click (the source position is already past the click, and dedup already remembers its id).
Branch P2 (from the reference implementation): c7 is read at 10:14:59.9, so its +1 for ad_9 is still buffered when barrier 41 arrives at 10:15:00.0.
| Buffer not flushed at the barrier | Flushed before the snapshot | |
|---|---|---|
c41's window state for ad_9 10:14 | 1 | 2 |
c41's dedup state | Holds c7 | Holds c7 |
| S1's saved position | After c7 | After c7 |
| After the crash and replay | c7 is never read again: ad_9 10:14 ends at 2 | ad_9 ends at 3 (correct) |
So the pre-aggregation rule's third line: on a barrier, forward every buffered partial before snapshotting (or snapshot the buffer itself). Any buffer outside checkpointed state needs a rule like this.
When the sink is slow
| # | At | Event |
|---|---|---|
| 13 | 10:15:25.0 | Valkey fails over (assumed); the window operator can't emit, and its input queues fill: backpressure |
| 14 | 10:15:30.0 | Checkpoint c42: the sources inject barrier 42, but at the window operator it waits behind the queue |
| 15 | 10:15:31.1 | c13 read on S1 |
| 16 | 10:15:40.0 | A worker dies. c42 never completed; the last completed checkpoint is c41 |
In the toy the queue holds only other ads' output; in the real job it holds thousands of records, and an aligned barrier must wait for all of them.
Checkpoints take 45 s at a 30 s interval during a backlog. What breaks, and in what order?
Unaligned barriers
An unaligned checkpoint lets the barrier overtake the queued records: the operator forwards the barrier at once and stores the records it overtook (the in-flight data) as part of its snapshot; on restore they are replayed from the checkpoint. Flink requires exactly-once checkpointing for it and doesn't take unaligned checkpoints concurrently. The setting is execution.checkpointing.unaligned.enabled (the older name, execution.checkpointing.unaligned, still appears in the managed service's list of settings).
Branch U (from the reference implementation): same run, but c42 is unaligned.
Aligned c42 (main line) | Unaligned c42 | |
|---|---|---|
c42 under backpressure | Waits behind the queue; never completes before the crash | Barrier overtakes the queue; complete at 10:15:31.0 (assumed 1 s) |
| Checkpoint size | State only | State plus the overtaken in-flight records |
| Restore uses | c41 | c42 |
| Input replayed | 10:15:00 → 10:15:40 = 40 s | 10:15:30 → 10:15:40 = 10 s (only c13) |
c12 | Replayed, and decided differently (Part 7) | Already in c42's state as a late total: not replayed at all |
| Restrictions | None | Exactly-once mode; no concurrent checkpoints |
Incremental uploads, interval, pause and timeout
- Incremental checkpoints upload only the RocksDB files created since the last checkpoint, not the whole state. At the loop's scale, about 11,574 clicks/s × 30 s × 60 B ≈ 20.8 MB of new dedup state per checkpoint, plus whatever compaction rewrote, instead of 6 to 43 GB.
- Interval (how often a checkpoint starts), minimum pause (the gap after one completes before the next may start) and timeout (when an unfinished checkpoint is abandoned; 10 minutes by default in Flink). Managed Service for Apache Flink's defaults are a 60,000 ms interval and a 5,000 ms minimum pause; the ad-click job sets 30 s.
- What a failure costs: everything since the last completed checkpoint, plus the restart: 30 + 50 = 80 s when checkpoints complete on time.
- Checkpoints complete even if the sink is "ahead" of them: output is not part of a checkpoint.
Savepoints
| Checkpoint | Savepoint | |
|---|---|---|
| Who triggers it | The job, periodically | You ("manually triggered checkpoints") |
| Used for | Recovery from failures | Upgrades, rescaling, moving a job |
| Kept | Replaced as newer ones complete | Until you delete it |
| On the managed service | Automatic | Called a snapshot; updates and scaling take one with stop-with-savepoint |
What to remember from Part 6
- A checkpoint is a consistent cut: every operator's state at the barrier plus every source's position; any buffer outside state must be flushed at the barrier.
- Backpressure delays aligned barriers; unaligned barriers overtake the queue at the price of storing it.
- A failure replays everything since the last completed checkpoint, plus the restart time.
Part 7. Crash, restore and replay
The worker died at 10:15:40.0, and the last completed checkpoint is c41. A restore must bring the job back to a state that matches its input, and then catch up. What it brings back, and what it doesn't, decides whether the output stays right.
What comes back, and what doesn't
Restored from c41? | After the restore | |
|---|---|---|
| Keyed state (windows, dedup ids, roll-ups, late totals) | Yes | As at the barrier: ad_7 10:14 = 4, ad_9 10:14 = 2 |
| Event-time timers | Yes | Fire again when the watermark passes them |
| Source positions | Yes | S1 after c7, S2 after c5, S3 after c3 |
| Operator state you declared (the watermark floor) | Yes | 10:14:34.999 |
| The watermark | No | Every input starts at Flink's minimum value (Long.MIN_VALUE) until it reports one (FLINK-5601, "Window operator does not checkpoint watermarks", is still open) |
| Idleness clocks, idle flags | No | Every shard starts active, with a fresh 30 s clock |
| In-memory buffers (the pre-aggregator's 500 ms) | No | Empty (so Part 6 flushes them at every barrier) |
| Anything written outside: DynamoDB, Valkey, an external Redis | No | Still holds whatever the first run wrote, including writes made after the checkpoint |
The replay, from the reference implementation:
| # | At | Event |
|---|---|---|
| 17 | 10:16:30.0 | The job runs again from c41 (50 s restart). The sink's epoch starts at the restored id, 41 |
| 18 | 10:16:30.0 | Full flush before any new output (Part 8) |
| 19 | 10:16:30.00 → 10:16:30.06 | The backlog is read at once: c9, c15, c8, c10, c11, c12, c13, before the first watermark emission |
| 20 | 10:16:30.01 | c15 reaches the pre-aggregator, which judges it against max(floor 10:14:34.999, watermark none yet) = 10:14:34.999. Its minute ended 10:14:00: late again. Late total ad_7 / H10 = 1, rewritten with late_ckpt 41 |
| 21 | 10:16:30.05 | c12 is judged against the same 10:14:34.999. Its minute ends 10:15:00, after it: on time. ad_9 10:14 = 3. At the first emission (10:16:30.2) S3 has reported, and its watermark (10:14:51.999 after c12) holds the job below 10:15:00: minute 10:14 stays open. The replay decided differently. |
The replay decides differently
Synthesizing vector architecture diagram...
What to notice: the same click, the same state at c41, two different decisions. In run 1, S3 had been idle for 30 s and minute 10:14 had closed without it; in the replay, S3's record is waiting at the start, S3 is active, and its watermark holds the minute open.
Lateness is not a function of the input alone. It depends on:
- How records from different shards interleave when they reach the operator. In run 1, S1 and S2 had moved on by the time
c12arrived; in the replay, every shard's backlog arrives together. - Idleness, which is processing time: S3 was idle in run 1 and active in the replay.
- When watermarks are emitted: every 200 ms of wall clock, which a replay doesn't reproduce.
- The watermark itself, which restarts from nothing.
None of these is in the checkpoint. So each run makes a valid, but possibly different, split into on time and late. Two consequences: the first run's decisions left values in the sink that the replay won't rewrite (Part 8), and a window that closed before the checkpoint can re-open (next).
After the restore, c12 is on time. It was late before the crash. Which run was right?
A closed minute re-opens
Right after the restore, the pre-aggregator has no watermark until every shard reports one. Branch W (from the reference implementation): the same restore without a floor.
| Without a floor, no closed-row condition | Without a floor, with the closed-row condition (Part 8) | With the floor (main line) | |
|---|---|---|---|
c15 at 10:16:30.01 | Judged against the minimum value: on time | Same | Judged against 10:14:34.999: late |
Its minute, ad_7 10:13 | Closed and purged before c41, so the window is re-created with a count of 1 and fires at the first emission, 10:16:30.2 | Same | Never re-created |
ad_7 / M10:13 | SET fast_valid = 1 over the correct 6 | The write is refused (stored fast_ckpt 39 < 41): stays 6 | Stays 6 |
c15's late total | Removed by the flush, never rewritten | Same | Rewritten: 1 |
| Hour roll-up | Already holds minute 10:13 = 6; the re-fired minute is not added again | Same | Unchanged |
| Result | A wrong row, 1 instead of 6: the flush has already run, and nothing will rewrite it | c15 missing: an under-count of 1 until the recount, never a wrong row | Correct |
Had hour 10 already closed before the checkpoint, the re-created minute would also have re-created the hour's roll-up, and its row would have been written with fast_valid = 1 over the whole hour.
How long is the exposure? Until every shard has reported a watermark. In the main line all three shards had backlog, so the first emission at 10:16:30.2 covered them: 200 ms. In branch U (restore from c42), S3 had nothing to replay, and the job's watermark stayed at the minimum value from 10:16:30.0 until S3 was marked idle at 10:17:00.4: 30.4 s, the whole idle timeout (the timer starts at the first check that sees no records, 10:16:30.2). That is the hook's 38,000-row case.
The watermark floor
textWATERMARK FLOOR (in the operator that decides lateness, before the window) on a new watermark w: floor_candidate = w on a barrier: save floor_candidate in operator state -- union list state on restore: floor = the MINIMUM of the saved values -- safe after rescaling on each record: judge = max(floor, current watermark) if the record's window ended at or before judge: late path (and the horizon check uses judge too) else: on to the window
It lives in the operator that decides lateness before the window: in the ad-click job, the pre-aggregation operator (the loop's step 2.4, item 3a); in a job without one, a small filter right before the window. Flink's built-in window operator can't apply it, because it judges lateness against its own input watermark; a custom keyed process function that does the windowing itself can. Timers need no floor: after a restore they can only fire late, never early. The acceptance-horizon check (Idempotency & Effectively-Once Processing, Part 8) uses the same max(floor, watermark). The floor is a design choice, not a Flink feature.
External state isn't rolled back
The restore rewinds Flink's state, not the world's. If dedup ids lived in an external Redis instead of keyed state, Redis would still hold c9 … c13 from run 1; the replay would find them there and drop them as duplicates, and the restore would under-count (the ad-click loop's step 2.6).
What to remember from Part 7
- A restore rewinds state and positions, not watermarks, idle flags or anything written outside.
- The same records can be late in one run and on time in the next.
- Restore a watermark floor in the operator that decides lateness, so nothing closed before the checkpoint can re-open.
Part 8. Sinks: making repeated and changed output harmless
Everything the job wrote between barrier 41 and the crash is written again by the replay, and some of it differently. The sink must end up as if only the replay had run. This Part is the mechanism behind the rule that page 04 states: effectively-once = exactly-once state + a sink that makes repeats harmless.
Output repeats, and changes
What the replay can do to each kind of value, if nothing else is done:
| Value written after the checkpoint | What the replay does | Without a cleanup |
|---|---|---|
| A value the replay writes again with the same result | Overwrites it | Fine (absolute SET) |
A value the replay writes again with a different result (ad_9 10:14: 2, then 3) | Overwrites it later | Wrong until the replay reaches it |
A value the replay never writes (ad_9's late total for c12) | Nothing | Wrong forever: c12 in both the minute and the late total |
An additive write (ADD fast_valid :n) | Adds again | Double counted |
Absolute values per key
Every write sets an absolute value per key: SET fast_valid = 5 on ad_7 / M10:14, never ADD. That handles the first row. For the others, the sink needs to know which values were written after the checkpoint.
Each writer tags its own attributes
Several writers share one item. The hour row carries the stream's fast_* (the roll-up), the late path's late_* and the recount's final_*. So a tag per item isn't enough: after a restore we must remove the late path's values without touching the recount's. Each writer tags only its own attributes, with the same names as the ad-click loop:
| Attributes | Written by | Tag |
|---|---|---|
fast_valid, fast_spend_micros, … on minute rows and hour rows | The window operator | fast_ckpt |
late_valid, late_spend_micros on hour rows | The late path | late_ckpt |
final_* on hour rows | The batch recount | None: the stream never touches them |
job_gen | Every stream write | (Part 9) |
textTAGGED WRITE (late path) UpdateItem ad_metrics key (entity, bucket) SET late_valid = :n, late_spend_micros = :s, late_ckpt = :epoch, job_gen = :g, late_ix = ":epoch#:shard" -- :shard = hash(entity, bucket) mod 8 CONDITION attribute_not_exists(job_gen) OR job_gen <= :g FULL FLUSH (after restoring checkpoint c, before any new output) for each epoch e >= c, for each shard n in 0..7: query each writer's tag index for "e#n" minute item, fast_ckpt >= c: delete the item -- its window is open or absent in restored state hour item, late_ckpt >= c: if restored state holds that late total: rewrite it, tag c else REMOVE late_valid, late_spend_micros, late_ckpt, late_ix hour item, fast_ckpt >= c: REMOVE the fast_* attributes, their tag and fast_ix never delete an item that carries another writer's attributes rewrite Valkey's running totals from restored state; delete totals written since c that state lacks then replay
The full flush, and the index that finds it
Event 18, from the reference implementation, at 10:16:30.0:
| Item | Found through | Action |
|---|---|---|
ad_7 / M10:14 (5, fast_ckpt 41) | 41#6 | Deleted: restored state holds this window open |
ad_9 / M10:14 (2, fast_ckpt 41) | 41#4 | Deleted |
ad_7 / H10 (late_valid 1, late_ckpt 41) | 41#4 | REMOVE late_*; the item stays (with any fast_* or final_* other writers put there) |
ad_9 / H10 (late_valid 1, late_ckpt 41) | 41#5 | REMOVE late_*; the item stays |
ad_7 / M10:13 (6, fast_ckpt 39) | Not found: tagged before c41 | Untouched |
Valkey ad_7 10:14 | Restored state | Rewritten to 4: a visible dip from 5 |
Valkey ad_9 10:14 | Restored state | Rewritten to 2 |
| Valkey 10:15 totals | Written since c41, not in state | Deleted |
Synthesizing vector architecture diagram...
Snapshot W4, 10:16:30.0: the sink now holds nothing the first run decided after c41. The minute row finished before the checkpoint (M10:13) is untouched.
Why the index is sharded. An index keyed only by the checkpoint id puts every write of one checkpoint interval on one index partition. In the loop, that is up to 2,800 write units a second, and one DynamoDB partition takes at most 1,000; a throttled index throttles writes to its table. With a suffix 0 to 7, each index partition takes about 2,800 ÷ 8 = 350. The flush queries all 8 shards for each epoch ≥ c. (Every sink-written item carries a tag, so the index isn't sparse in any useful sense; querying by tag value is what keeps the flush's read small.)
Why exact tags shrink the flush. With an exact tag, only values tagged ≥ c need attention: values written earlier were written from state at or before the checkpoint, so they can never be ahead of restored state. The ad-click loop also rewrites every key restored state holds (about 300,000 late totals, about 107 s): safe, but more than necessary. The cost that remains: an index entry with every sink write, the flush's own writes competing with live output (rate-limit them), and the dip.
The tag must be set at the barrier
Which id goes into the tag? The obvious answer, "the last completed checkpoint", has a gap: a checkpoint completes some time after the sink passed its barrier, and a write in between carries the older id.
Branch G (from the reference implementation): c41's upload is slow and completes at 10:15:20.5 instead of 10:15:01.5. c12's late total is written at 10:15:20.0.
| Tag rule | c12's late total tagged | Closed-row condition | After the replay: ad_9 minute + late |
|---|---|---|---|
| Last completed | 40 | Off | Flush (epochs ≥ 41) misses it; the replay writes minute 10:14 = 3: 3 + 1 = 4, c12 billed twice |
| Last completed | 40 | On | Flush misses it; the replay's minute 10:14 = 3 is refused, because the minute row was also tagged 40 and looks older than the checkpoint: the sink keeps run 1's split, 2 + 1 = 3, correct here only by luck, and the minute row disagrees with the job's state |
| Last barrier passed | 41 | Off or on | Flush removes it; the replay writes 3: 3 + 0 = 3 |
So the tag is the id of the last barrier this sink subtask has snapshotted. Flink hands it to the sink when it snapshots: Sink V2's StatefulSinkWriter.snapshotState(checkpointId), or FunctionSnapshotContext.getCheckpointId() for older sink functions. On start after a restore, initialise it from the restored checkpoint's id (getRestoredCheckpointId(), an optional value). Never use notifyCheckpointComplete: it arrives later, and Flink calls these notifications "best effort", so they can be skipped. Checkpoint ids never fall below the restored id. A failover continues the counter, and restoring a savepoint sets it to the savepoint's id + 1, so "≥ c" finds everything written since. (After a savepoint restore, two job versions can use overlapping ids: branch B.)
Closed rows are final
The flush handles values written after the checkpoint. Branch W showed the other danger: a replay rewriting a row finished before it. So closed-window attributes are written with a condition:
textCLOSED-WINDOW WRITE UpdateItem ad_metrics key (ad, "M10:14") SET fast_valid = :n, fast_ckpt = :epoch, job_gen = :g, fast_ix = ":epoch#:shard" CONDITION (attribute_not_exists(fast_ckpt) OR fast_ckpt >= :c) -- :c = the restored checkpoint (before any restore, the id the job started from) AND (attribute_not_exists(job_gen) OR job_gen <= :g) a refused write is counted, not retried; the recount repairs what it leaves out
Closed-row condition (fast_ckpt >= :c) | Window fence (last_window < :w) | |
|---|---|---|
| Write | Absolute value per window key | ADD into a longer total, e.g. a per-hour or per-video counter |
| Refuses | Any rewrite of a row finished before the restored checkpoint | Any window at or before the last one added |
| Used in | The ad-click loop's fast_ckpt condition (item 3a) | YouTube step 2.5, URL shortener step 2.2 |
| Late totals | Not conditioned: they are running values, and the flush handles them |
Transactional sinks and their timeouts
The other way to make output safe is to not show it until the checkpoint that covers it completes. A two-phase-commit sink (Flink's Kafka sink in exactly-once mode, or its file sink) writes each checkpoint interval's output inside a transaction:
textTWO-PHASE-COMMIT SINK 1. after barrier n: open transaction T(n); write output into it 2. at barrier n + 1: flush and pre-commit T(n); open T(n + 1) 3. when checkpoint n + 1 completes: commit T(n) 4. after a failure: commit or abort the open transactions according to the restored state readers use isolation.level = read_committed
Branch K (from the reference implementation): the job also writes closed minutes to a Kafka topic in exactly-once mode.
Synthesizing vector architecture diagram...
What to notice: nothing to flush, because run 1's output was never committed; but readers see every row only when the next checkpoint completes. Without the crash, the 10:15:07.2 rows would have become visible when c42 completed at 10:15:31.5, 24.3 s later.
A read_committed reader stops at the partition's last stable offset, the first message of any open transaction, so an open transaction also holds back everything written after it. And a transaction the broker times out before the sink commits it is lost ("data loss may happen when Kafka expires an uncommitted transaction", in Flink's words). Flink's documentation puts the requirement as "maximum checkpoint duration + maximum restart duration"; the checkpoint timeout bounds the duration:
With our numbers: we set the checkpoint timeout to 60 s, so 30 + 60 + 50 = 140 s; choose a transaction timeout of 600 s: 140 < 600 ≤ 900 s, the broker's 15-minute default maximum (MSK doesn't change it). With Flink's default 10-minute checkpoint timeout the left side would be 30 + 600 + 50 = 680 s, still under 900 but with little room. And one trap: Flink's KafkaSink sets the producer's transaction.timeout.ms to 1 hour by default, 3,600 s > 900 s, and the broker refuses a producer whose timeout is above its maximum when it initialises (InitProducerId). A KafkaSink left at its default fails on a default broker: lower it.
Why not make the DynamoDB sink transactional and skip the flush?
| Idempotent sink + flush (DynamoDB, Valkey) | Two-phase-commit sink (Kafka, files) | |
|---|---|---|
| When readers see output | At once | When the checkpoint after the output completes: up to one interval plus the checkpoint's duration (1.5 to 31.5 s here; 24.3 and 30.9 s in branch K) |
| After a restore | Flush every writer's values tagged ≥ c; a dip | Nothing to flush for uncommitted output |
| Reader requirements | Overwrite by key; tolerate the dip | read_committed, or it also sees aborted output (the default is read_uncommitted) |
| Timeout risk | None | Transaction timeout must fit the timeout rule above and the broker's maximum |
| What it costs every write | A tag and an index entry | Nothing extra, but open transactions hold back readers |
| Which stores qualify | Any store with keyed overwrites, conditions and an index | Only stores with transactions committed later |
The replay then continues, from the reference implementation:
| # | At | Event |
|---|---|---|
| 22 | 10:16:45.2 and 10:16:50.3 | Live clicks c14 (S1) and c16 (S2), clicked after the restart |
| 23 | 10:17:00.0 | Checkpoint c43; the floor saved is now 10:14:51.999 |
| 24 | 10:17:00.6 | S3, silent since c12 was re-read, is marked idle; the job watermark jumps to min(S1 10:16:39.999, S2 10:16:44.999) = 10:16:39.999. Minutes 10:14 and 10:15 close: SET fast_valid ad_7 / M10:14 = 5, ad_9 / M10:14 = 3, tagged fast_ckpt 43; the condition passes (the items were deleted). No late_valid for ad_9: the flush removed run 1's. ad_9 for minute 10:14: 3 + 0 = 3, correct |
Synthesizing vector architecture diagram...
Snapshot W5, 10:17:00.6: minute 10:14 closed again, with c12 on time this time and no trace of run 1's late decision. Every click is counted exactly once: ad_7 5 in minute 10:14 plus c15 in hour 10's late total; ad_9 3.
What to remember from Part 8
- Checkpoints make state exactly-once; the sink must make repeated and changed output harmless.
- After a restore, remove each writer's attributes written after the restored checkpoint, found by a per-writer tag set at the sink's barrier; never delete another writer's data; condition closed rows so a replay can't rewrite them.
- A transactional sink needs no flush but shows output only at checkpoints, and its timeout must fit inside the broker's.
Part 9. Rescaling and upgrades
The Super Bowl traffic arrives, and the job must go from 2 subtasks to 4. Later, a new version must replace the old one. Both move keyed state, and both are restores with everything Parts 7 and 8 require.
Key groups
Keyed state can't be rehashed key by key on every rescale: that would mean reading every key. Instead, Flink splits the key space into a fixed number of key groups, the unit that moves:
- Key group = murmurHash(key's hash code) mod max parallelism (128 here).
- Owner = key group × parallelism ÷ max parallelism, rounded down. So each subtask owns one contiguous range of key groups.
Branch Q, computed with Java's string hash and Flink's murmur hash in the reference implementation:
| Key | Key group | Parallelism 2: owner | Parallelism 4: owner |
|---|---|---|---|
ad_7 | 117 | 117 × 2 ÷ 128 = 1.83 → 1 | 117 × 4 ÷ 128 = 3.66 → 3 |
ad_9 | 44 | 44 × 2 ÷ 128 = 0.69 → 0 | 44 × 4 ÷ 128 = 1.38 → 1 |
Synthesizing vector architecture diagram...
What to notice: whole ranges of key groups move, and each new subtask reads from exactly one old one. The number of key groups (128) never changes, so no key is rehashed.
Max parallelism is therefore the ceiling for rescaling and must be chosen before the first deployment: changing it later means the job can't restore from its old snapshots. The default differs: open-source Flink uses min(max(roundUpToPowerOfTwo(1.5 × parallelism), 128), 32,768), while Managed Service for Apache Flink sets 128 for every operator if the parallelism is at most 128 and nothing is set.
Rescaling is a restore
A rescale is a savepoint plus a restart with the new parallelism. Everything from Parts 7 and 8 applies: the replay from the savepoint's positions, the flush (every writer's values tagged ≥ the savepoint's id), and the floor: the floor is saved as union list state, so after a rescale every new subtask receives all the old subtasks' values and takes their minimum, which is safe whichever keys it now owns.
The managed service's automatic scaling does exactly this: when the maximum containerCPUUtilization stays at or above 75% for 15 minutes, it doubles the parallelism (and halves it after 6 hours below 10%), and the application has downtime while it restarts from a snapshot.
You double the parallelism during the Super Bowl. What does the job do in the next two minutes, and what do dashboards show?
Blue/green and job_gen
An in-place update on the managed service uses stop-with-savepoint, so two versions never overlap. The risk appears when you run two applications side by side: v2 starts from a savepoint of v1 and both write to the same sink until v1 is stopped. Each version gets a higher generation number, job_gen, and every write carries attribute_not_exists(job_gen) OR job_gen <= :g.
Branch B (from the reference implementation): v1 is job_gen 7; v2 is job_gen 8, started from v1's savepoint c51 taken at 10:20:00; v1 keeps running until 10:23:00. v2 excludes one click v1 counts.
| At | Write | With job_gen | Without it |
|---|---|---|---|
| 10:21:35 | v1 creates a late total ad_9 / H10 (late_ckpt 54) | Passes the condition: v2 hasn't written this item | Same |
| 10:22:05.1 | v2: ad_7 / M10:21 = 3, :g = 8 | Stored: 3, job_gen 8 | Stored: 3 |
| 10:22:05.2 | v1: ad_7 / M10:21 = 4, :g = 7 | 8 ≤ 7 is false: refused; stays 3 | 4: the old version wins by arriving last |
| Flush | If v2 flushed at 10:20:00, that late total, created afterwards, survives | v2 flushes at restore (10:20:00) as always; once v1 stops (10:23:00), a second sweep over values tagged ≥ 51 removes only those v1 wrote (job_gen 7): the late total tagged 54 goes; v2's ad_7 / M10:21 = 3 (job_gen 8) stays |
So the rule, the ad-click loop's too: the new version flushes at its restore as always, and once the old one has stopped writing, sweeps once more to remove only what the old one wrote since the savepoint. Until then v2 writes normally. A caveat: job_gen here is per item, so if v2 writes an item's fast_* after v1 wrote its late_*, the item's job_gen becomes 8 and v1's late_* can no longer be told apart; a fully general sweep needs a generation per writer (fast_gen and late_gen, or gen#epoch#shard in the index key).
Two different numbers, two jobs: job_gen says which version wrote a value (fencing between versions, like a fencing token in Leases, Fencing Tokens & Distributed Locks); fast_ckpt/late_ckpt say which checkpoint interval wrote it (for the flush). The two versions' checkpoint ids overlap (both continue from 51), and that is fine: the flush asks "≥ 51" and finds both versions' writes, and job_gen decides between them.
What to remember from Part 9
- Max parallelism is the number of key groups and the ceiling for rescaling; choose it before the first deployment.
- Every rescale and upgrade is a restore: the replay, the flush and the floor run again.
- Two versions writing one sink need a generation on every write; once the old one stops, the new one sweeps away only the old one's values written since the savepoint.
Part 10. Compared on equal terms: stream, micro-batch, batch recount
A continuous stream is one way to count. Micro-batch engines, Kafka Streams and a nightly batch are others, and each answers the three questions of Part 1 differently.
Four ways to count minute 10:14
From the reference implementation, the same clicks (run 1's read order, dedup applied to all), with comparable settings: a 5 s allowance, grace or threshold everywhere.
| Engine | ad_7 10:14 | ad_9 10:14 | c15 | c12 |
|---|---|---|---|---|
| Flink, no crash | 5 | 2 | Late path (hour 10) | Late path (hour 10) |
| Flink, after the crash and replay | 5 | 3 | Late path (hour 10) | On time |
| Kafka Streams, grace 5 s (stream time = highest event time the task has seen) | 5 | 2 | Dropped (no side output) | Dropped |
| Spark Structured Streaming, threshold 5 s, a 10 s trigger | 5 | 2 | Dropped at the 10:15:10 batch (watermark 10:14:53) | Dropped at the 10:15:20 batch (watermark 10:15:02) |
| Batch recount | 5 | 3 | In minute 10:13 (7) | In minute 10:14 |
Two things to see. Only the batch recount and the Flink replay put c12 in minute 10:14, and Flink's run-1 answer plus its late total gives the same hour total. And Kafka Streams and Spark lose c15 and c12 from the stream result entirely unless something else counts them.
The table
| Continuous stream (Flink) | Micro-batch (Spark Structured Streaming) | Kafka Streams | Batch recount | |
|---|---|---|---|---|
| Latency | Milliseconds to seconds | One trigger interval or more | Milliseconds to seconds | Hours |
| How completeness is decided | Watermark = minimum over sources of (highest event time − allowance) | Watermark = maximum event time seen by the query − threshold, set for the next batch (the minimum across several watermarked inputs by default) | Stream time = highest event time the task has seen; a window closes at its end + grace | The input is complete |
| Lateness knob | Allowance, allowed lateness, side output | Threshold: data within it is guaranteed to be counted; later data "may or may not" be | Grace period; later records dropped | None needed |
| What a restart replays | Since the last checkpoint | The failed batch, from the offset log | Since the last commit, state from its changelog | The whole input |
| Can a restart reclassify? | Yes: watermarks aren't checkpointed | No for a re-run batch: the watermark is stored with each batch in the offset log | Yes, it can: the window operator's stream time is kept in memory and rebuilt from the replayed records | No |
| Exactly-once covers | State; output only with a transactional or idempotent sink | State; output only with a replayable source and an idempotent or transactional sink | State; output with its transactional mode | The whole result, if the job overwrites its output |
| Cost | Always-on workers | Always-on, bursty | Runs inside your service | Compute for the batch only |
| Code paths | One | One | One | A second one, next to a stream |
Lambda and kappa, honestly
Lambda is a stream for fast, provisional answers plus a batch recount for exact ones. The honest cost: the rules (dedup, fraud, lateness) live in two code paths that must agree, and the ad-click loop reconciles them every hour with final = fast + late − A + B. Kappa is one stream job, and reprocessing means replaying the stream. The honest cost: a replay needs the input kept long enough, and replaying raw input brings back records that were later rejected, so the source of truth must be stored with its verdicts. The ad-click loop keeps lambda for billing because a bug in the stream job would cost money, and the recount is the second opinion.
Finance wants the dashboard number to be the billed number. What do you tell them?
Reprocessing a day
A bug found after a day can't be fixed by replaying the stream: Kinesis keeps records 24 hours by default. So:
- Read the day from the S3 lake as a bounded source, with the same event-time code.
- Run it in batch execution mode: the input is complete, so there are no watermarks and no lateness to decide.
- Write under a new
job_gen, intofinal_*-style attributes or a separate table that is swapped in when checked; never over the live stream's attributes. - Never replay raw input without the exceptions log (duplicates, invalid clicks, late clicks, each with its verdict), or records that were rejected come back.
What to remember from Part 10
- A stream answers in seconds and is provisional; a recount answers exactly, hours later.
- Micro-batch trades latency for a watermark that moves once per batch and is stored with it; batch trades more latency for no lateness at all.
- Keep a second computation when a bug in the stream job would cost money.
Part 11. End to end through the layers
"Flink handles it" hides six places where a click lives, each with different rules about what survives a crash. Following c12 through all six turns the page into a checklist.
One click through six layers
Synthesizing vector architecture diagram...
What to notice: only the dashed arrows reach the checkpoint. The watermark generator and the sink are outside it, and they are exactly where run 1 and the replay differed.
| Layer | What it holds for c12 | Checkpointed? | Lost on restart | May repeat or change | Rule |
|---|---|---|---|---|---|
| Source | S3's position (after c3 at c41) | Yes | Nothing | c12 is read again | Positions and state in the same cut |
| Watermark generator | S3's highest event time; the idle flag | No | Watermark, idle flag, idle clock | Run 1: S3 idle, c12 late; replay: S3 active | Assume it will differ; the floor covers the start |
| Operator | Lateness decision; floor; 500 ms buffer | Floor yes; buffer only if flushed at the barrier | The buffer, unless flushed | The decision can change | Decide lateness per click with max(floor, watermark); flush at the barrier |
| State backend | ad_9 10:14 window count; late total | Yes | Nothing after c41 is lost; nothing before c41 is duplicated | Rebuilt by the replay | Keep dedup and totals in state, never in an external store |
| Checkpoint storage | c41: positions, state, floor | It is the checkpoint | Nothing | Durable, off the worker (S3) | |
| Sink | ad_9 / H10 late_valid 1 (run 1) | No | Nothing: it keeps run 1's writes | Run 1's values stay unless cleaned | Absolute writes, per-writer tags, flush, closed-row condition |
The four numbers to watch
| Number | Managed Flink metric | What it tells you | The ad-click loop's alarm |
|---|---|---|---|
| Watermark lag (wall clock − watermark) | currentOutputWatermark | Windows aren't closing: a stuck or quiet source, a stuck operator | > 60 s for 5 minutes |
| Consumer lag | millisBehindLatest (Kinesis) | The job can't keep up with the stream | > 30 s for 5 minutes |
| Checkpoint duration | lastCheckpointDuration, numberOfFailedCheckpoints | Backpressure or state growth; the replay after a crash is growing | > 10 s, or 3 failures in a row |
| State size | lastCheckpointSize, containerDiskUtilization | Timers not firing (a stuck watermark), or key growth | Trend, against the 50 GB per KPU |
What to remember from Part 11
- Only the source position and operator state are checkpointed; watermarks, when timers fire, and sink writes are recomputed.
- Every layer that holds something outside the checkpoint needs a rule: the pre-aggregator flushes at the barrier, the sink flushes by tag.
- Watch four numbers: watermark lag, consumer lag, checkpoint duration, state size.
Part 12. On AWS
The mechanism is Flink's, whoever runs it. What AWS adds is where it runs, where the state and checkpoints live, and which sources and sinks it uses.
Managed services that use it
| Service | What it provides | What the documentation says |
|---|---|---|
| Amazon Managed Service for Apache Flink (formerly Kinesis Data Analytics for Apache Flink) | Runs Flink: checkpoints, snapshots, scaling. Your code sets the watermark strategy, windows, lateness policy and the sink's idempotence | Default checkpointing: enabled, 60,000 ms interval, 5,000 ms minimum pause; with configuration type DEFAULT these apply even if set in code. Updates, scaling and stops take snapshots with stop-with-savepoint. With snapshots enabled, the service "provides exactly-once processing semantics during application updates, or during service-related scaling or maintenance". A KPU is 1 vCPU and 4 GB with 50 GB of running application storage, plus one KPU per application for orchestration; 0.10 per GB-month of running storage in us-east-1. Default limit 64 KPUs per application. Autoscaling doubles parallelism after 15 minutes at ≥ 75% CPU and halves it after 6 hours below 10%, with downtime while it scales. Max parallelism is 128 when parallelism ≤ 128 and nothing is set; changing it prevents restoring older snapshots. The state backend type, incremental checkpoints and the unaligned-checkpoint setting are among the settings you change through a support case |
| Amazon Kinesis Data Streams | The replayable source: per-shard positions in every checkpoint, per-shard watermarks | 1 MB/s or 1,000 records/s of writes per shard; 24-hour default retention. Flink's Kinesis source reads a shard only after its parent shards are fully read, so order per key survives resharding |
| Amazon MSK | A replayable source (per-partition watermarks; offsets committed to Kafka when checkpoints complete) and the transactional sink | Broker transaction.max.timeout.ms is Kafka's default, 900,000 ms (15 minutes); a producer asking for more is refused at InitProducerId. Consumers default to read_uncommitted; read_committed consumers read only up to the last stable offset |
| Amazon EMR | Runs Flink (as a YARN application) or Spark Structured Streaming on clusters you size | The self-managed engine with AWS-managed instances; checkpoints in S3 |
| Amazon S3 | Checkpoint and savepoint storage for self-run Flink; the lake the recount and reprocessing read | Durable storage off the workers; the managed service keeps its checkpoints for you |
| Amazon DynamoDB | The idempotent sink: absolute SETs, per-writer tags, conditions (closed rows, job_gen), a tag index sharded <ckpt>#<0-7> | Each partition delivers at most 1,000 write units a second; if an index lacks write capacity, "the write activity on the table will be throttled" |
ElastiCache for Valkey appears on this page as the running-total store that the flush also rewrites; it is a plain store the job writes to, not a service that uses the mechanism.
Running it yourself
| Option | What it is | Sizing and notes |
|---|---|---|
| Amazon EC2 running Flink | You own JobManager high availability, checkpoint storage (S3), RocksDB disks, upgrades and rescaling | RocksDB on instance-store NVMe (fast, lost with the instance, which is fine because checkpoints are in S3) or on EBS gp3 (whose IOPS and throughput are shared by state reads, compaction and checkpoint staging). The loop's state: ~6 GB typical, 43 GB worst case, over all workers; incremental checkpoints ≈ 20.8 MB per 30 s at the average rate plus compaction output. The ad-click loop's build-or-buy step prices an m7g.large at about $0.041 per KPU-equivalent hour, about 2.7 times cheaper compute, and still chooses the managed service |
Shared limits
- The watermark is shared by every key: one slow or stuck shard holds back every window of every ad; idleness trades that for late data.
- Checkpoint uploads share the worker's network with source reads, sink writes and shuffles; RocksDB compaction and checkpoint staging share the local disk with processing.
- Backpressure delays aligned barriers: one slow sink lengthens checkpoints, which lengthens the replay, which lengthens the flush.
- The 50 GB per KPU holds state, compaction space and staging; a stuck watermark lets state grow until the guard TTL.
- The flush shares the sink's write capacity with live output (2,800 WCU in the loop); rate-limit it.
- The tag index takes every sink write: shard it (
<ckpt>#<0-7>: about 350 WCU per index partition against 1,000). - The replay shares processing capacity with live input: the backlog drains at capacity − arrivals (the loop: 64K − 11.6K ≈ 52K clicks/s on an average day; 28K/s at the Super Bowl).
- An open Kafka transaction holds back every
read_committedreader of its partitions.
Look-alikes that are not this mechanism
| Look-alike | Why it looks like it | Why it isn't |
|---|---|---|
| Kafka's high watermark (MSK) | The word | The offset up to which a partition is replicated; nothing about event time |
Kinesis iterator age / MillisBehindLatest | "How far behind are we?" | Processing lag behind the stream's tip. Watermark lag is a different number; alarm on both |
| KCL checkpoints | "Checkpoint" | A consumer's stored position only, no operator state; at-least-once. The metrics loop builds the state half by hand (upload, then checkpoint) |
| Amazon Data Firehose buffering | "It groups records by time" | Buffers by size and arrival time before delivery; no event-time windows, no watermarks |
| Lambda tumbling windows on Kinesis and DynamoDB Streams | "Tumbling windows" | Windows by the time records were inserted into the stream, per shard, at most 15 minutes, 1 MB of state per shard, no resharding; an idle window closes after up to 2 minutes. Ingestion time, not event time: no watermark, no late data |
| SQS FIFO "exactly-once processing" | The words | A 5-minute deduplication of sends; consumers still see redeliveries |
| DynamoDB TTL and Flink state TTL | "Expire after 2 hours" | TTL deletes within days (DynamoDB) or by processing time (Flink); neither is an event-time horizon |
| Kinesis Data Analytics for SQL | The old name of the managed service | A different, discontinued product: no new applications since 2025-10-15, deleted from 2026-01-27 |
What to remember from Part 12
- Managed Service for Apache Flink runs the checkpoints and snapshots; the watermark strategy, the lateness policy and the sink's idempotence are yours.
- Kinesis and MSK give the replayable positions; DynamoDB (and Valkey) are sinks that need the flush.
- Scaling and upgrades restart from a snapshot: plan the replay and the flush.
Part 13. What you've learned
Back to the double count
A worker died, the job restored c41 and replayed, and a click could have been counted twice without anything being added twice. Here is how each piece did its job in our minute:
- Event time put every click in the minute it happened (Part 1), and the watermark, the minimum over active shards, decided when minute 10:14 was complete: 10:15:07.2, with S3 left out as idle (Part 2, snapshot W2).
- Hours are built from closed minutes, and the pre-aggregator decides lateness per click and flushes before watermarks, so nothing is counted in two places or lost between them (Part 3).
c15andc12went to the side output, and the dashboards overwrite early results by key (Part 4, snapshot W3).- Checkpoint
c41cut state and positions together, with the pre-aggregator's buffer flushed at the barrier (Parts 5 and 6, snapshot W1). - The restore brought state back but not the watermark:
c12was on time in the replay, and the floor keptc15(and minute 10:13) from re-opening a closed minute (Part 7). - The flush removed every value tagged 41 or later, per writer, and the closed-row condition protected rows finished earlier; the replay rewrote minute 10:14 as 5 and 3, with no late total for
ad_9(Part 8, snapshots W4 and W5). - Rescales and blue/green switches are restores with the same rules, plus
job_gen(Part 9).
The whole story, event by event
| # | Wall clock | Event | Result |
|---|---|---|---|
| 1 | 10:14:05.3 → 10:14:40.3 | c1 to c5 read | Job watermark 10:14:14.999 (S3) at 10:14:40.4 |
| 2 | 10:14:20.4 | c3 lifts S3 | Minute 10:13 fires: ad_7 / M10:13 = 6, fast_ckpt 39 |
| 3 | 10:14:45.0 | c2′ | Duplicate, dropped |
| 4 | 10:14:50.8 | S3 idle | Job watermark 10:14:25.999 |
| 5 | 10:14:58.4 | c7 | Job watermark 10:14:34.999 |
| 6 | 10:15:00.0 → 10:15:01.5 | Checkpoint c41 | 4 and 2 in state; floor 10:14:34.999 (W1) |
| 7 | 10:15:02.5 | c9 (minute 10:15) | Job watermark 10:14:52.999 |
| 8 | 10:15:04.0 | c15 | Late: ad_7 / H10 late_valid 1, late_ckpt 41 |
| 9 | 10:15:05.0 | c8 | On time: ad_7 10:14 = 5 |
| 10 | 10:15:06.2 | c10 | Job watermark 10:14:56.999 |
| 11 | 10:15:07.1 → 10:15:07.2 | c11, then the tick | Minute 10:14 closes: 5 and 2, fast_ckpt 41 (W2) |
| 12 | 10:15:20.0 | c12 on S3 | Late: ad_9 / H10 late_valid 1, late_ckpt 41 (W3) |
| 13 | 10:15:25.0 | Valkey failover | Backpressure |
| 14 | 10:15:30.0 | Checkpoint c42 | Barrier stuck at the window operator |
| 15 | 10:15:31.1 | c13 | On time |
| 16 | 10:15:40.0 | Crash | Last completed: c41 |
| 17 | 10:16:30.0 | Restore from c41 | State 4 and 2; epoch 41; no watermark |
| 18 | 10:16:30.0 | Full flush | Two minute items deleted, late_* removed from two hour items; Valkey dip 5 → 4 (W4) |
| 19 | 10:16:30.00 → .06 | Backlog read | Seven clicks before the first emission |
| 20 | 10:16:30.01 | c15 against the floor | Late again: ad_7 / H10 = 1 |
| 21 | 10:16:30.05 | c12 against the floor | On time: ad_9 10:14 = 3 |
| 22 | 10:16:45.2, 10:16:50.3 | c14, c16 | Live clicks |
| 23 | 10:17:00.0 | Checkpoint c43 | Floor 10:14:51.999 |
| 24 | 10:17:00.6 | S3 idle | Minute 10:14 closes: 5 and 3, fast_ckpt 43; no late total for ad_9 (W5) |
The cheat card
| Topic | Remember |
|---|---|
| Two clocks | Event time says where; the watermark says when; the lateness policy says what happens after |
| Watermark | Per source: highest event time − allowance (Flink: − 1 ms more), emitted every 200 ms; job: minimum over active sources, never backwards |
| Idleness | Processing-time timeout; stops a quiet source from freezing the job; its next records are likely late |
| Windows | Tumbling (reports), sliding (size ÷ slide copies per record), sessions (merge, even by a late record) |
| Hours | The sum of closed minutes, each added once; never a separate hour window |
| Pre-aggregation | Lateness per click; forward partials before the watermark and at the barrier |
| Late data | Allowed lateness (updates) or side output (separate total); acceptance horizon; recount for the final answer |
| Early firings | Running totals: overwrite by key, or send retractions |
| Timers vs TTL | Event-time timers for dedup and windows; Flink's TTL is processing time, a guard only |
| Checkpoint | Barriers; state + positions in one cut; aligned waits, unaligned stores in-flight records |
| Restore | State, positions, timers back; watermark, idle flags, buffers and sink writes not |
| Floor | Last watermark, in operator state; after a restore judge with max(floor, watermark), in the operator before the window |
| Sink | Absolute writes; per-writer tag = last barrier passed; flush ≥ c per writer; closed-row condition; job_gen |
| 2PC sink | Visible at checkpoints; interval + checkpoint timeout + restart < transaction timeout ≤ the broker's maximum (15 minutes by default); lower KafkaSink's 1-hour default |
| Rescale | Key groups = max parallelism; owner = key group × parallelism ÷ max parallelism; a rescale is a restore |
Failure checklist
- Are windows assigned by event time, and is device time clamped before the watermark sees it?
- Is the job's watermark the minimum over sources, with idleness on, and is watermark lag alarmed?
- Do hour rows come from closed minutes, never a separate hour window?
- Does the pre-aggregator decide lateness per click, flush before forwarding a watermark, and flush at the barrier?
- Is dedup state cleared by event-time timers, with TTL only as a much longer guard?
- Does anyone believe "exactly-once" covers the output? Is every sink write absolute, never
ADD, or fenced withlast_window? - After every restore and rescale, does the flush remove every writer's values tagged ≥ the restored checkpoint, and only that writer's, before new output?
- Is each tag the last barrier the sink passed, initialised from the restored id, never the last completed checkpoint?
- Is a watermark floor restored and enforced before the window, and are closed rows conditioned on their tag?
- For a 2PC sink: interval + checkpoint timeout + restart < transaction timeout ≤ the broker's maximum;
KafkaSink's 1-hour default lowered; readersread_committed? - For side-by-side versions: a generation on every write, and a second sweep (the old version's values only) after the old one stops?
Think-first drills
Drill 1. A job reads 8 Kafka partitions with a 10 s allowance and no idleness. Partition 5 gets no data for 20 minutes every night. What do dashboards show, what do you change, and what does the change cost the first record partition 5 sends afterwards?
Drill 2. Checkpoints every 60 s, restart 90 s, checkpoint timeout 120 s, and a Kafka sink in exactly-once mode. Set the transaction timeout and check it against the broker. What else must be true of the readers?
Drill 3. After a rescale from 32 to 64, one ad's hour row shows the same click in fast_valid and late_valid. Name the two things the restore should have done, and the tag each sink write needed.
Interview questions
| Question | Model answer |
|---|---|
| Event time vs processing time: when does each give the wrong answer? | Processing time moves records between windows whenever the job lags or replays, so its counts depend on the job's health; use it only for metrics about the pipeline. Event time is right about where a record belongs but needs a watermark to decide completeness, and a clamp when device clocks are wrong. |
| What is a watermark, how is it computed across partitions, and what happens with an idle partition? | "Event time has reached X." Per partition: highest event time − allowance, emitted periodically; the job takes the minimum over partitions, and it never goes back. An idle partition would freeze the minimum, so a processing-time idleness timeout removes it; its next records behind the watermark are late. |
| A click arrives an hour late: design what happens to it. | Its window closed long ago, so it goes to a side output and a late total per ad and hour, written as an absolute value; dashboards add on-time and late. Past the acceptance horizon it is rejected; an exact recount decides the final number. Allowed lateness for an hour would keep an hour of window state. |
| "Flink is exactly-once": what does a checkpoint cover, and what must the sink do? | Operator state and source positions, cut together by barriers, so state is exactly-once. Output after the checkpoint is emitted again, sometimes differently. The sink writes absolute values, tags each writer's attributes with the last barrier it passed, and after a restore removes values tagged ≥ the checkpoint; or it is transactional and commits on checkpoint completion. |
| After a crash the job reprocesses data: why can the result differ, and how do you keep the output correct? | Lateness depends on shard interleaving, idleness and emission timing, and the watermark isn't checkpointed, so a replay can classify records differently. Keep only the replay's decisions: flush the first run's post-checkpoint values by tag, restore a watermark floor so closed windows don't re-open, and condition closed rows. |
| Streaming, micro-batch or a nightly batch for billing: pick and defend. | A stream for the dashboard (seconds, provisional) plus a batch recount for the bill (exact, hours later), reconciled every hour: lambda. Micro-batch is steadier on restart (the watermark is logged per batch) but slower. A stream alone can't be the bill: late data and replays keep changing its split. |
Where to go next
- Idempotency & Effectively-Once Processing: the rules this page implements (effectively-once, the flush, dedup expiry and the acceptance horizon), in Parts 4 and 8.
- Change Streams & the Transactional Outbox: how events get out of databases and into the streams this page counts.
- Leases, Fencing Tokens & Distributed Locks: fencing tokens in general;
job_genis one. - LSM-Trees & Compaction: how RocksDB, the state backend, stores and compacts state.
- Replication, Quorums & Read-Your-Writes: replication of the stores the job reads and writes.
- Background primitives: Message Queues vs Event Streams, Two-Phase Commit & Saga Orchestration, Event Sourcing & CQRS.
- Loops: the ad-click aggregation loop (steps 2.1 to 2.6, R2.5 to R2.8, R3.9), the metrics and alerting loop (steps 2.2, 2.3, 2.5), YouTube (step 2.5), the URL shortener (step 2.2), the rate limiter (step 3.2), search autocomplete (step 2.1) and the Google Maps loop.