Stream and batch
Count clicks per page every 10 seconds. A batch job waits for all the data and is always right. A stream job answers within seconds, so it must handle clicks that arrive late and crashes in the middle.
- 1Batch sees all the data at once. A stream must decide when a window is complete.
- 2Window by event time. The watermark delay trades waiting against dropped late data.
- 3Exactly once needs three parts: a replayable source, offsets saved with the state, and an idempotent or transactional sink.
- 4With a log that keeps history, one stream job also does the backfills (Kappa).
- Phones go offline, queues back up, and retries arrive late.
- Processing-time windows are simple, but they count late events in the wrong window.
- Event-time windows are right, but the job must decide how long to wait.
I window by event time, because a click that arrives late still belongs to the moment it happened.
| method | where | status |
|---|---|---|
| Batch SQL every night | Postgres | Hours of delay are fine |
| Windows by processing time | Service | Metrics of the job itself |
| Event-time windows, watermark | Flink | Approved |
| Micro-batches, every trigger | Spark | 100 ms or more |
| Counts in memory, offsets committed | Service | Not approved |
| Checkpoint, transactional sink | Postgres | Approved |
| Batch and stream, merged (Lambda) | Service | Two code paths |
| Replay the log (Kappa) | Kafka | Approved |
Counts in memory with committed offsets lose or repeat events after a crash. The lab run below counted 8 clicks twice.
If an hour of delay is fine, I run batch. If not, I stream with event-time windows and keep the log so I can replay.
| shape | defined by | a click is in | use for |
|---|---|---|---|
| Tumbling | Size | Exactly 1 window | Counts per minute, billing per hour |
| Sliding | Size and step | Size ÷ step windows (2 here) | Moving averages, alerts over the last 5 minutes |
| Session | A gap of silence | 1 session, which can merge | User visits, rides, calls |
| Global | A custom trigger | 1 window for the key | Every 100 events, or on a signal |
Tumbling windows partition time, sliding windows overlap so each event counts more than once, and session windows close after a gap.
on each click c, in arrival order1:
FOR EACH window w that c's event time falls in:
IF w.end + lateness <= watermark: skip w // state gone: c is dropped2
count[w] += 1
IF w already fired: emit count[w] as an update3
watermark = max(watermark, largest event time seen - delay4)
FOR EACH window w not fired, with w.end <= watermark:
emit count[w]; mark w fired // the first result5
forget every w with w.end + lateness <= watermark
at the end of the input: watermark = infinity; fire the rest- 1The job sees events in processing time. Their event times can be out of order.
- 2Past the lateness, the window is forgotten. The click goes to a side output or is lost.
- 3A window that already fired sends a corrected count. Downstream must accept updates.
- 4The watermark: a bound on how late events can be. It moves only when events arrive.
- 5The window is complete as far as the watermark knows.
Tested source Go: the aggregator
// Run feeds events in arrival order through the windows. The watermark is the largest event time
// seen minus Delay. A window fires when the watermark reaches its end. A late event for a fired
// window updates it within Lateness; after that the event is dropped.
func Run(events []Event, cfg Config) (Result, error) {
if err := cfg.Window.valid(); err != nil {
return Result{}, err
}
in := slices.Clone(events)
slices.SortStableFunc(in, func(a, b Event) int { return cmp.Or(cmp.Compare(a.Arrive, b.Arrive), cmp.Compare(a.ID, b.ID)) })
res := Result{Status: map[int]EventStatus{}}
var open []*state
watermark := math.MinInt
maxSeen := math.MinInt
for _, e := range in {
var touched []*state
late := false
if cfg.Window.Kind == Session {
seed := assign(cfg.Window, e.At)[0]
if seed.End+cfg.Lateness > watermark {
st := mergeSession(&open, seed)
touched = append(touched, st)
late = seed.End <= watermark || st.fired
}
} else {
for _, w := range assign(cfg.Window, e.At) {
if w.End+cfg.Lateness <= watermark {
continue // past its lateness: this window's state is gone
}
st := find(open, w)
if st == nil {
st = &state{w: w}
open = append(open, st)
}
st.count++
touched = append(touched, st)
late = late || w.End <= watermark
}
}
switch {
case len(touched) == 0:
res.Status[e.ID] = Dropped
case late:
res.Status[e.ID] = Kept
default:
res.Status[e.ID] = OnTime
}
// A window that already fired emits its corrected count at once.
for _, st := range touched {
if st.fired {
res.Emissions = append(res.Emissions, Emission{st.w, st.count, e.Arrive, true})
}
}
maxSeen = max(maxSeen, e.At)
watermark = max(watermark, maxSeen-cfg.Delay)
res.Marks = append(res.Marks, Mark{e.Arrive, watermark})
res.Emissions = append(res.Emissions, fire(&open, watermark, cfg.Lateness, e.Arrive)...)
}
// The end of the input moves the watermark to the end of time: every window fires.
res.Emissions = append(res.Emissions, fire(&open, math.MaxInt-cfg.Lateness, cfg.Lateness, -1)...)
return res, nil
}
// fire emits every window whose end the watermark has reached, once, and forgets windows whose
// lateness has run out.
func fire(open *[]*state, watermark, lateness, at int) []Emission {
var out []Emission
slices.SortFunc(*open, func(a, b *state) int { return cmp.Compare(a.w.Start, b.w.Start) })
for _, st := range *open {
if !st.fired && st.w.End <= watermark {
st.fired = true
out = append(out, Emission{st.w, st.count, at, false})
}
}
*open = slices.DeleteFunc(*open, func(st *state) bool { return st.w.End+lateness <= watermark })
return out
}
| setting | larger value gives | costs |
|---|---|---|
| Watermark delay | Fewer late events; first results are complete more often. | Every result waits longer. |
| Allowed lateness | Late events update their window instead of being dropped. | Window state stays longer; results change after they were sent. |
| Side output | Dropped events are kept for a repair job. | A second path to run. |
- recorded
- Session windows dropped 9 of 43 clicks at a delay of 0 s and 1 at 10 s.
- wait
- The first result came 3.6 s after the session ended at 0 s delay, and 13.5 s at 10 s delay.
- defaults
- Flink allows 0 s of lateness by default. Kafka Streams discards a record after the window's grace period.
The watermark says no more events older than this should come. A window fires when the watermark passes its end; later events update it or are dropped.
| window | batch | first result | updates | status |
|---|---|---|---|---|
| [0, 10) | 8 | 8 at 17 s | · | ✓ matches batch |
| [10, 20) | 7 | 7 at 26 s | · | ✓ matches batch |
| [20, 30) | 7 | 6 at 37 s | · | ✕ 1 missing |
| [30, 40) | 8 | 5 at 46 s | · | ✕ 3 missing |
| [40, 50) | 8 | 6 at 54 s | · | ✕ 2 missing |
| [50, 60) | 5 | 5 at end of input | · | ✓ matches batch |
With no delay and no lateness, several windows came out wrong. A short delay plus some allowed lateness gave every window the batch answer.
Step 1: Read
- Each task reads its partitions in offset order.
- Event time comes from the event, not from the clock of the job.
If it fails
The job restarts: it seeks back to the offsets in the last checkpoint and reads again.
The job reads the log, keys and windows each click by event time, keeps counts as state, and checkpoints offsets and state together.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Flink | Window assigners: tumbling, sliding, session, global | The four shapes, by event time, per key. | Sessionizing, rollups |
| Flink | Watermark strategy with bounded out-of-orderness | The watermark: the largest event time seen, minus a set delay. | Event-time joins |
| Flink | allowedLateness, sideOutputLateData | Update a fired window, or route the late event aside. Lateness is 0 by default. | Repair jobs |
| Flink | Checkpoints by barriers | A consistent snapshot of offsets and state, taken while the job runs. Off by default. | Savepoints for upgrades |
| Flink | Transactional Kafka sink | Output commits only when the checkpoint completes. | Exactly-once pipelines |
| Kafka Streams | processing.guarantee = exactly_once_v2 | Reads, state and writes commit together in a Kafka transaction. At least once is the default. | Microservice aggregates |
| Kafka Streams | commit.interval.ms | 30 s by default; 100 ms under exactly once, so output is not held back long. | |
| Spark | Micro-batches, withWatermark, checkpointLocation | Latency as low as 100 ms with exactly once. Continuous mode: about 1 ms, at least once. | ETL into a lakehouse |
| Spark | Output modes: append, update, complete | Append writes a window once it is final; update rewrites changed rows. | Dashboards |
| Kafka | Retention and offsets | The replayable source: restore seeks back; a backfill starts from an old offset. | Queues, change data capture |
| Postgres | Transactions | The dashboard rows and the checkpoint row commit together. | Outbox, ledgers |
| Postgres | INSERT ... ON CONFLICT DO UPDATE | An idempotent sink: a replay writes the same count to the same key. | Upserts, materialized counters |
Flink gives me event-time windows, watermarks and barrier checkpoints. Kafka gives me a log to replay. Postgres gives me a transaction that holds output and progress.
- Delivery is still at least once inside the job: events after the checkpoint are read twice.
- The effect is exactly once because state and output roll back to the same point.
- A sink that cannot roll back or overwrite breaks the chain, such as an email.
restore():
(next, counts) = load the checkpoint row // offset and state match1
read the log from offset next
process(click):
counts[window, page] += 1
in the open transaction2: write counts[window, page]
next += 1
IF next is a multiple of 103: checkpoint()
checkpoint():
in the same transaction: save (next, counts)
COMMIT // output and progress together4- 1One row holds both, so the counts always describe exactly the events before the offset.
- 2The dashboard write is not visible until the checkpoint commits.
- 3The checkpoint interval. Shorter means less replay and fresher output, and more commits.
- 4A crash before this line rolls back both. A crash after it loses nothing.
Tested source Go: restore, process, checkpoint · SQL: sink and checkpoint · SQL: tables
// Restore loads progress after a start or a crash. With auto-commit only the offset survives;
// a checkpoint brings back the offset and the state that matches it.
func (j *Job) Restore(ctx context.Context) error {
j.next, j.state = 0, map[key]int{}
var raw []byte
err := j.DB.QueryRow(ctx, jq["load"], j.Name).Scan(&j.next, &raw)
if errors.Is(err, pgx.ErrNoRows) {
return nil
}
if err != nil {
return fmt.Errorf("load progress of %s: %w", j.Name, err)
}
if j.Mode == AutoCommit || raw == nil {
return nil
}
var saved []struct {
K key `json:"k"`
N int `json:"n"`
}
if err := json.Unmarshal(raw, &saved); err != nil {
return fmt.Errorf("decode state of %s: %w", j.Name, err)
}
for _, s := range saved {
j.state[s.K] = s.N
}
return nil
}
// Process handles the click at offset j.Next() and, every Every events, commits or checkpoints.
func (j *Job) Process(ctx context.Context, c Click) error {
k := key{c.At - c.At%WindowSize, c.Page}
j.state[k]++
var err error
switch j.Mode {
case AutoCommit, CheckpointAdd:
_, err = j.DB.Exec(ctx, jq["add"], k.Window, k.Page)
case CheckpointPut:
_, err = j.DB.Exec(ctx, jq["put"], k.Window, k.Page, j.state[k])
case Transactional:
if j.tx == nil {
if j.tx, err = j.DB.Begin(ctx); err != nil {
return fmt.Errorf("begin: %w", err)
}
}
_, err = j.tx.Exec(ctx, jq["put"], k.Window, k.Page, j.state[k])
default:
return fmt.Errorf("unknown mode %q", j.Mode)
}
if err != nil {
return fmt.Errorf("write %+v: %w", k, err)
}
j.next++
if j.next%j.Every == 0 {
return j.Checkpoint(ctx)
}
return nil
}
// Checkpoint saves progress. In Transactional mode the sink writes since the last checkpoint
// commit in the same transaction as the offset and the state.
func (j *Job) Checkpoint(ctx context.Context) error {
var state []byte
if j.Mode != AutoCommit {
type entry struct {
K key `json:"k"`
N int `json:"n"`
}
var all []entry
for k, n := range j.state {
all = append(all, entry{k, n})
}
var err error
if state, err = json.Marshal(all); err != nil {
return fmt.Errorf("encode state: %w", err)
}
}
if j.Mode != Transactional {
if _, err := j.DB.Exec(ctx, jq["save"], j.Name, j.next, state); err != nil {
return fmt.Errorf("save progress at %d: %w", j.next, err)
}
return nil
}
if j.tx == nil {
return nil
}
if _, err := j.tx.Exec(ctx, jq["save"], j.Name, j.next, state); err != nil {
return fmt.Errorf("save progress at %d: %w", j.next, err)
}
err := j.tx.Commit(ctx)
j.tx = nil
if err != nil {
return fmt.Errorf("commit at %d: %w", j.next, err)
}
return nil
}
-- Idempotent: writes the count from the job's state. A replay writes the same value again.
INSERT INTO window_counts (window_start, page, clicks) VALUES ($1, $2, $3)
ON CONFLICT (window_start, page) DO UPDATE SET clicks = EXCLUDED.clicks;
-- The checkpoint: the next offset and the state, in one row, so they always match.
INSERT INTO job_progress (job, next_offset, state) VALUES ($1, $2, $3)
ON CONFLICT (job) DO UPDATE SET next_offset = EXCLUDED.next_offset, state = EXCLUDED.state;
SELECT next_offset, state FROM job_progress WHERE job = $1;-- The output: clicks per page per 10-second window, read by a dashboard.
CREATE TABLE window_counts (
window_start int NOT NULL,
page text NOT NULL,
clicks int NOT NULL,
PRIMARY KEY (window_start, page)
);
-- Where the job is in the log. An auto-committing consumer stores only the offset. A
-- checkpointing job stores the offset and its state together, as one row.
CREATE TABLE job_progress (
job text PRIMARY KEY,
next_offset bigint NOT NULL,
state jsonb
);On a crash, the job restores offsets and state from one checkpoint and replays from there. The sink commits with the checkpoint, so a replay never writes twice.
After this step
- saved
- next offset 10
- job
- next offset 10
- dashboard
- 10 clicks
- log read
- 10 clicks
Saving only offsets, or adding to the sink, showed 68 clicks for 60. Overwriting or one transaction showed exactly 60.
| sink | after a replay | status |
|---|---|---|
| Add 1 per event | Replayed events count twice. | Not approved |
| Overwrite the count by key | Same final values. A count can fall, then climb back. | Approved |
| Commit with the checkpoint | Nothing doubled. Output appears once per interval. | Approved |
| Insert with a unique event id | The repeat is refused by the key. | Approved |
| Send an email or a webhook | Sent again. | Receiver dedupes |
I make the sink idempotent by key, or commit it with the checkpoint. A sink that adds is never safe under replay.
| join | state kept | example |
|---|---|---|
| Stream with table | The latest row per key of the table | Click plus the user's country |
| Stream with stream, windowed | Both sides, for the window plus the lateness | Ad view plus click within 1 hour |
| Interval join | Events of B within a range around each A | Order plus payment within 10 minutes |
| Table with table | Both tables in full | Users plus accounts |
- Both inputs must be partitioned by the join key, so matching events meet in one task.
- A join without a time bound keeps every event forever.
A stream join keeps state for each side. I bound it by time, and partition both inputs by the join key.
Lambda keeps a batch path to correct the stream. Kappa replays the log through the same job, so there is one code path to keep right.
| step | do |
|---|---|
| 1 | Start job v2 at the offset of the first bad day. Write to a new table. |
| 2 | Limit its read rate, so the live job and the sink keep up. |
| 3 | When v2's lag nears 0, compare v1 and v2 on recent windows. |
| 4 | Switch reads to the v2 table. Stop v1, then drop its table. |
- Event-time logic gives the same answer on replay; processing-time logic does not.
- In the lab, a delay as long as the latest click gave every window the batch answer.
- The log must keep the whole range. Otherwise, read old data from the batch store.
To backfill, I run a new version of the job from an old offset into a new table, then switch reads when it catches up.
| event | result | why it is safe | saved by |
|---|---|---|---|
| A click arrives after its window fired | Dropped, or the window fires again. | Allowed lateness updates the window; a side output keeps the rest. | Lateness |
| One partition goes quiet | The watermark stops; windows never fire. | Mark the quiet source as idle, so it does not hold the watermark back. | Idle source |
| The job crashes between checkpoints | Events since the checkpoint are read again. | State and sink roll back to the same checkpoint. | Transaction |
| The sink adds per event | 68 clicks shown for 60. | Not safe. Overwrite by key, or commit with the checkpoint. | Upsert |
| An overwrite sink replays | 3 counts fell, then recovered. | Final values are right. Readers that need monotonic counts read committed output. | Transaction |
| Checkpoints are far apart | A long replay after a crash; late output with a transactional sink. | Correct, but slow. Pick the interval from the replay time you accept. | Interval |
| State grows without limit | Memory or disk runs out. | Bound joins by time; expire idle keys with a TTL. | TTL |
| A device clock is wrong | Clicks from the future move the watermark. | Reject or clamp event times far ahead of the receive time. | Validation |
| One page gets most clicks | One task does most of the work. | Pre-aggregate in each task, then combine the partial counts. | Two-phase sum |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | Batch SQL every few minutes: GROUP BY page and window over the raw table. | Always complete for closed windows. Delay is the schedule plus the run time. | Results must come within seconds, or the scan takes longer than the schedule. |
| 2 | One consumer with event-time windows and a transactional sink in Postgres. | About 2,700 events a second committing each one; 15,000 committing every 1,000. | State outgrows one process, or one consumer cannot keep up. |
| 3 | A stream processor: keyed state, parallel tasks, barrier checkpoints. | One task per partition. The windows alone ran about 2,900,000 events a second in memory here. | Backfills and bug fixes need a second code path. |
| 4 | Replay from the log for backfills, with long retention. | One code path for live data and history. | Top of the ladder. |
Demand example: 1 billion clicks a day is 11,574 a second; a 10× peak is 115,741. The job rates are one client on one laptop with the WAL flushed at each commit, so read them as orders of magnitude.
I start with batch SQL. When results must come within seconds, I run a consumer with a transactional sink, and move to a stream processor when state outgrows one process.
0 of 10 known
A click happens at 10:00:07 and reaches the job at 10:00:25. Which 10-second window counts it?
What does the watermark delay trade?
A click arrives after its window fired. What can the job do?
With a delay of 0 s, first results still came about 6 s after each tumbling window ended. Why?
The job saves offset and counts at offset 30 and crashes after 37. The sink adds 1 per click. What does the dashboard show?
What makes the output exactly once, not only the state?
An overwrite sink and a transactional sink both end with the right counts. What does a reader see differently?
What is the difference between Lambda and Kappa?
You fix a bug in the counting. How do you backfill a month?
What bounds the state of a join between two streams?
- windows
- About 2,900,000 to 3,600,000 events a second on one core, in memory.
- commits
- About 2,700 events a second committing each; 15,000 committing every 1,000.
- replay
- A crash replays up to one checkpoint interval of events.
- Kafka Streams
- Commits every 30 s by default; every 100 ms under exactly once.
- Spark
- Micro-batches as low as 100 ms; continuous mode about 1 ms, at least once.
- Flink
- Checkpoints off and lateness 0 by default.
Go aggregator on one goroutine; job and Postgres 16 on an 8-core laptop, one client, median of 3 runs. Defaults are from the Kafka, Spark and Flink documentation.