System Design
B7

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.

Not startedSaved in this browser only.
  1. 1Batch sees all the data at once. A stream must decide when a window is complete.
  2. 2Window by event time. The watermark delay trades waiting against dropped late data.
  3. 3Exactly once needs three parts: a replayable source, offsets saved with the state, and an idempotent or transactional sink.
  4. 4With a log that keeps history, one stream job also does the backfills (Kappa).
B7
    A

    Two clocks

    event time against processing time
    Event timewhen the click happened, on the phoneProcessing timewhen it reached the job[0, 10)[0, 10)[10, 20)[10, 20)[20, 30)[20, 30)[30, 40)[30, 40)abcdc: phone offline, 18 s late[0, 10): 3 clicks by event time, 2 by arrival. Only event time is right.
    • 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.

    B

    Batch or stream, where

    for counts over time
    methodwherestatus
    Batch SQL every nightPostgresHours of delay are fine
    Windows by processing timeServiceMetrics of the job itself
    Event-time windows, watermarkFlinkApproved
    Micro-batches, every triggerSpark100 ms or more
    Counts in memory, offsets committedServiceNot approved
    Checkpoint, transactional sinkPostgresApproved
    Batch and stream, merged (Lambda)ServiceTwo code paths
    Replay the log (Kappa)KafkaApproved

    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.

    C

    Window shapes

    the same clicks, three ways
    clicks by event time0102030405060Tumbling10 s877885Sliding10 s, every 5 s48576778881053Sessiongap 4 s877885number in a window: clicks it counts
    shapedefined bya click is inuse for
    TumblingSizeExactly 1 windowCounts per minute, billing per hour
    SlidingSize and stepSize ÷ step windows (2 here)Moving averages, alerts over the last 5 minutes
    SessionA gap of silence1 session, which can mergeUser visits, rides, calls
    GlobalA custom trigger1 window for the keyEvery 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.

    D

    Watermarks and late data

    pseudo code, then the trade-offs
    event-time windowspseudo code
    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
    1. 1The job sees events in processing time. Their event times can be out of order.
    2. 2Past the lateness, the window is forgotten. The click goes to a side output or is lost.
    3. 3A window that already fired sends a corrected count. Downstream must accept updates.
    4. 4The watermark: a bound on how late events can be. It moves only when events arrive.
    5. 5The window is complete as far as the watermark knows.
    Tested source Go: the aggregator
    Go: the aggregatorgo
    // 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
    }
    
    settinglarger value givescosts
    Watermark delayFewer late events; first results are complete more often.Every result waits longer.
    Allowed latenessLate events update their window instead of being dropped.Window state stays longer; results change after they were sent.
    Side outputDropped 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.

    E

    See it: windows and watermarks

    43 recorded clicks, every setting
    window
    watermark delaythen
    0102030405060010203040506070arrival at the job, sevent time, son timelate, countedlate, droppedwatermarkno delaybands: 10 s stepsof event time
    6 dropped0 late, counted3 of 6 windows match batchfirst result 6 s after the window ends, on average
    windowbatchfirst resultupdatesstatus
    [0, 10)88 at 17 s·✓ matches batch
    [10, 20)77 at 26 s·✓ matches batch
    [20, 30)76 at 37 s·✕ 1 missing
    [30, 40)85 at 46 s·✕ 3 missing
    [40, 50)86 at 54 s·✕ 2 missing
    [50, 60)55 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.

    F

    The pipeline

    click a step; its path lights up
    Click logpartitions, retentionStream jobN tasks, keyedState storecounts per windowLate clicksside outputCheckpointsoffsets and statePostgresdashboard table

    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.

    G

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    FlinkWindow assigners: tumbling, sliding, session, globalThe four shapes, by event time, per key.Sessionizing, rollups
    FlinkWatermark strategy with bounded out-of-ordernessThe watermark: the largest event time seen, minus a set delay.Event-time joins
    FlinkallowedLateness, sideOutputLateDataUpdate a fired window, or route the late event aside. Lateness is 0 by default.Repair jobs
    FlinkCheckpoints by barriersA consistent snapshot of offsets and state, taken while the job runs. Off by default.Savepoints for upgrades
    FlinkTransactional Kafka sinkOutput commits only when the checkpoint completes.Exactly-once pipelines
    Kafka Streamsprocessing.guarantee = exactly_once_v2Reads, state and writes commit together in a Kafka transaction. At least once is the default.Microservice aggregates
    Kafka Streamscommit.interval.ms30 s by default; 100 ms under exactly once, so output is not held back long.
    SparkMicro-batches, withWatermark, checkpointLocationLatency as low as 100 ms with exactly once. Continuous mode: about 1 ms, at least once.ETL into a lakehouse
    SparkOutput modes: append, update, completeAppend writes a window once it is final; update rewrites changed rows.Dashboards
    KafkaRetention and offsetsThe replayable source: restore seeks back; a backfill starts from an old offset.Queues, change data capture
    PostgresTransactionsThe dashboard rows and the checkpoint row commit together.Outbox, ledgers
    PostgresINSERT ... ON CONFLICT DO UPDATEAn 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.

    H

    Exactly once

    figure, then pseudo code
    records flow right; barrier 3 follows offset 29Sourcereads the log3savesWindow countskeyed statesavesSinkto the dashboardsavesCheckpoint 3, in durable storageoffset 30counts after 0 to 29output, uncommittedAll parts saved: checkpoint 3 is complete, and the sink commits.A crash before that: restore checkpoint 2, replay from its offset.
    • 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, process, checkpointpseudo code
    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
    1. 1One row holds both, so the counts always describe exactly the events before the offset.
    2. 2The dashboard write is not visible until the checkpoint commits.
    3. 3The checkpoint interval. Shorter means less replay and fresher output, and more commits.
    4. 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
    Go: restore, process, checkpointgo
    // 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
    }
    
    SQL: sink and checkpointsql
    -- 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;
    SQL: tablessql
    -- 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.

    I

    See it: crash and restore

    recorded on Postgres
    Logoffsets 0 to 59saved: 10next: 10Jobnext offset: 10counts saved in each checkpointProgress tablenext offset: 10offset, counts, and the sink rowsDashboard tableclicks shown: 10clicks read from the log: 10✓ matches
    step 1 of 9

    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.

    J

    Sink choices

    what a replay does to the output
    sinkafter a replaystatus
    Add 1 per eventReplayed events count twice.Not approved
    Overwrite the count by keySame final values. A count can fall, then climb back.Approved
    Commit with the checkpointNothing doubled. Output appears once per interval.Approved
    Insert with a unique event idThe repeat is refused by the key.Approved
    Send an email or a webhookSent 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.

    K

    Joins on streams

    every join keeps state
    joinstate keptexample
    Stream with tableThe latest row per key of the tableClick plus the user's country
    Stream with stream, windowedBoth sides, for the window plus the latenessAd view plus click within 1 hour
    Interval joinEvents of B within a range around each AOrder plus payment within 10 minutes
    Table with tableBoth tables in fullUsers 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.

    L

    Lambda and Kappa

    one answer, one or two code paths
    Lambdatwo code paths for one answerKappaone code path; reprocess by replayEventsraw logBatch joball data, nightlyStream joblast hours, liveServingmerge bothQueryreadsLoglong retentionStream job v1liveStream job v2replays from 0Table v1serving nowTable v2catching upQueryswitch to v2Kappa needs a log that keeps enough history, and a job fast enough to catch up.

    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.

    M

    Backfills and reprocessing

    fix a bug, recompute a month
    stepdo
    1Start job v2 at the offset of the first bad day. Write to a new table.
    2Limit its read rate, so the live job and the sink keep up.
    3When v2's lag nears 0, compare v1 and v2 on recent windows.
    4Switch 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.

    N

    Failure cases

    what breaks, and what keeps the counts right
    eventresultwhy it is safesaved by
    A click arrives after its window firedDropped, or the window fires again.Allowed lateness updates the window; a side output keeps the rest.Lateness
    One partition goes quietThe 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 checkpointsEvents since the checkpoint are read again.State and sink roll back to the same checkpoint.Transaction
    The sink adds per event68 clicks shown for 60.Not safe. Overwrite by key, or commit with the checkpoint.Upsert
    An overwrite sink replays3 counts fell, then recovered.Final values are right. Readers that need monotonic counts read committed output.Transaction
    Checkpoints are far apartA 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 limitMemory or disk runs out.Bound joins by time; expire idle keys with a TTL.TTL
    A device clock is wrongClicks from the future move the watermark.Reject or clamp event times far ahead of the receive time.Validation
    One page gets most clicksOne task does most of the work.Pre-aggregate in each task, then combine the partial counts.Two-phase sum
    O

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Batch SQL2+ consumer, checkpoints3stream processor4+ replay (Kappa)more load →
    Measured rates and demand1k10k100k1M10MDemand: 1B clicks a day: 11,574 events per secondDemand: 1B clicks a day11,574Demand: 10× peak: 115,741 events per secondDemand: 10× peak115,741Job, commit every event: 2,654 events per secondJob, commit every event2,654Job, commit every 100: 12,479 events per secondJob, commit every 10012,479Job, commit every 1,000: 14,823 events per secondJob, commit every 1,00014,823Windows in memory: 2,908,925 events per secondWindows in memory2,908,925events per second, log scale
    stepaddit handlesmove up when you see
    1Batch 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.
    2One 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.
    3A 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.
    4Replay 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.

    P

    Drill

    predict, then reveal

    0 of 10 known

    1. A click happens at 10:00:07 and reaches the job at 10:00:25. Which 10-second window counts it?

    2. What does the watermark delay trade?

    3. A click arrives after its window fired. What can the job do?

    4. With a delay of 0 s, first results still came about 6 s after each tumbling window ended. Why?

    5. The job saves offset and counts at offset 30 and crashes after 37. The sink adds 1 per click. What does the dashboard show?

    6. What makes the output exactly once, not only the state?

    7. An overwrite sink and a transactional sink both end with the right counts. What does a reader see differently?

    8. What is the difference between Lambda and Kappa?

    9. You fix a bug in the counting. How do you backfill a month?

    10. What bounds the state of a join between two streams?

    Q

    Numbers to say

    measured, then documented
    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.