skip to content

Consumer Operations & Lag

The operational state of the side that reads: how far behind it is, what happens when its members change, and what to do when it stalls. Probed because a broker incident is often a reader incident.

part ofBroker & streaming operationsoverview, primer and where to startread it →
on this pageshow

questions

23

After an outage, what must a reader group's reading rate exceed before its backlog of unread records starts to shrink?

level: juniorimportance: must knowfreq 62%

answer

  1. two rates, one difference
  2. surplus, not throughput
  3. steady-state capacity only holds the line
  4. divide by the difference
  5. peak arrival rate, not average

basics

~20 s

The arrival rate. Only the surplus between reading rate and arrival rate drains a backlog of unread records, so capacity sized to keep pace in steady state merely holds the backlog at whatever depth the outage left it.

solid answer

~40 s

Writers do not pause while the reading side recovers, so a catch-up plan works on two rates, not one. If records arrive at `R` per second and the reader group completes `D` per second, the backlog shrinks at `D - R` per second and the time to close it is the backlog divided by that surplus — never by `D`. A group provisioned for normal traffic has `D` roughly equal to `R`, a surplus near zero, and will sit at the same depth indefinitely while looking perfectly healthy. So the first question of any drain is where extra capacity above the arrival rate is going to come from, and the second is whether the arrival rate you used is the average or the peak, because a daily peak can turn the surplus negative again.

go deeper

for a junior

Remember there are two rates, and only their difference empties anything. Backlog divided by (reading rate minus arrival rate) is the catch-up time, and matching the arrival rate exactly means never recovering.

for a middle

Be able to run both units and show they agree: the unread count falls by the surplus each second, and the age of the oldest unread record falls by the ratio of the two rates minus one. Disagreement means an input is wrong.

for a senior

Quote drain rate as completed work, not records fetched, and compute against the arrival rate present during the drain window rather than the daily mean. State the assumption you used out loud when you give an estimate.

for a principal

Decide whether surplus capacity for recovery is a standing cost or something bought at the time, and make sure the answer exists before an incident. A group with no reachable surplus has no recovery plan, only a hope.

A **backlog of unread records** is the set of records that exist in a stream and have not yet been handled by the reader group responsible for them. It appears whenever the reading side stops, slows or crashes while the writing side carries on, and it is the normal aftermath of a deployment gone wrong, a downstream outage or a night of undersized capacity. The moment the readers are healthy again, the operational question is not *are they running* but *are they running fast enough*, and that turns out to be arithmetic rather than judgment. ## Three numbers describe the whole situation - **Arrival rate (`R`)** — records written into the stream per second. This is the number people forget, because writers are unaffected by how far behind the reading side is. - **Drain rate (`D`)** — records the reader group actually completes per second, measured end to end, including the downstream work each record triggers. Not what the readers could fetch; what they finish. - **Backlog (`B`)** — the unread count at the moment the drain starts. Where a stream is measured in time rather than records, the equivalent is the **unread age** of the oldest unread record, and `B` is roughly that age multiplied by `R`. The backlog falls at `D - R` records per second. That difference is the drain; `D` on its own is not. ``` surplus = D - R time to close = B / (D - R) example: R = 4,000/s D = 5,000/s B = 18,000,000 surplus = 1,000/s time = 18,000,000 / 1,000 = 18,000 s = 5 hours ``` Read with `D` instead of the surplus, the same example gives one hour — a five-fold under-estimate, and the reason so many recovery estimates are announced and then quietly revised. ## Why steady-state capacity never recovers Most reader groups are sized so that they keep up with normal traffic and no more, sometimes with a modest safety factor. That sizing gives `D` approximately equal to `R` and a surplus of approximately zero. Such a group, restarted after a two-hour outage, is not broken and is not falling further behind — it is simply frozen at the depth the outage created. Every record it handles is matched by a new arrival. Nothing on a dashboard looks like an error; the gap is flat rather than rising, and it stays flat for as long as you leave it. Recovery therefore requires deliberately creating capacity that does not exist in normal operation, which is what makes a drain a plan rather than a wait. ## The same arithmetic in time units Where the gap is tracked as the age of the oldest unread record rather than a count, the drain shows up differently. While the group is behind, it reads historical records, advancing through `D / R` seconds of history for every second of wall clock. The unread age therefore changes at `1 - D / R` per second: it falls only when `D > R`, holds when `D = R`, and climbs at up to one second per second when the readers are stopped altogether. | Quantity | What it is | Behaviour while `D > R` | |---|---|---| | unread count | records not yet handled | falls by `D - R` each second | | unread age | age of the oldest unread record | falls by `D / R - 1` each second | | arrival rate | records written each second | unchanged by the drain | Both units give the same answer for when the group is level, which is a useful cross-check on a plan: if they disagree, one of the two inputs is wrong. ## Peak, not average A surplus computed against a daily average is optimistic for the hours that matter. If traffic doubles for three hours each evening, a group with a small surplus at the average rate has a **negative** surplus during the peak and gives back part of what it gained. Compute the plan against the rate that will actually be arriving during the drain window, and if that window spans a peak, say so in the estimate rather than discovering it. ## What varies between platforms The arithmetic is universal, but what limits `D` is not. Where a stream is split into a fixed set of parts and each reader takes some of them, the number of readers that can usefully work at once is capped by the stream's shape, and the surplus is capped with it. Where readers merely compete for records from a shared queue, there is no such cap, and the ceiling appears somewhere else — a broker-side rate ceiling on the client, or, more often than anyone expects, the downstream system each record is written into.

  • Arrivals are not flat across the day. How does that change the estimate?
    The surplus is a function of time, not a constant. Compute it against the arrival rate that will be present during the drain window: a group with a thin surplus at the daily average can have a negative surplus through a peak and lose ground for those hours. Quote catch-up times against peak, and if the window spans one, say which part of the estimate assumes recovery only happens off-peak.
  • What if the reading side can only ever match the arrival rate exactly?
    Then the backlog freezes at its current depth and the age of the oldest unread record holds flat. That is stable, not recovering, and it can persist for days without triggering anything. The options are to create surplus capacity, to reduce what arrives, or to take a decision about the group's recorded read position — which is not an operator's call alone and belongs with the stream's owner.
  • Why measure the drain rate end to end rather than at the broker?
    Because the surplus is set by completed work, not by records handed to the reader. If each record triggers a write to a downstream store, the store's throughput is the real `D`, and a group that fetches far faster than it finishes only builds an in-memory queue of its own. Measure where the work actually ends, and the plan stops being optimistic.

Bailing out a boat that is still taking water: what empties the hull is buckets per minute minus litres per minute coming in, and a crew that bails exactly as fast as the leak keeps the boat afloat forever without ever getting it dry. There is also a limit on how many people physically fit at the gunwale, which is why the answer is rarely just 'more hands'.

saying these in an interview costs you the question

  • Divides the backlog by total reading rate and promises an hour
  • Assumes writers pause while the reading side catches up
  • Says any running reader group must eventually catch up
  • Calls a flat, non-zero gap recovery in progress
  • Uses the daily average arrival rate for an evening drain
  • Measures the drain rate at the fetch instead of at completion
open as a page

When a reader group shares one stream's work, which membership events cause its shares to be reassigned?

level: juniorimportance: must knowfreq 68%

basics

~20 s

Three events: a member joins, a member leaves cleanly, or a member stops proving it is present — by missing its liveness signal or by holding work past its progress deadline. Designs where readers merely compete reassign nothing.

open as a page

A reader group's lag shows two million unread records, so why does that count alone not say whether anyone is waiting?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Unread count states the gap in records; unread age states it in time. Two million records may be seconds old on a fast stream or half a day old on a slow one, so only the age reading shows staleness.

open as a page

What happens to the waiting unread records when an operator moves a reader group's read position forward to the newest record?

level: juniorimportance: must knowfreq 60%

basics

~20 s

Moving a reader group's read position forward abandons every record between the old position and the new one, so nothing processes them. The gap drops to zero at once because work was skipped, not done.

open as a page

A reader group member still holds its share and makes no progress: what separates a stuck member from a merely slow one?

level: juniorimportance: must knowfreq 62%

basics

~10 s

A slow member still completes records, just fewer than arrive, so its read position keeps advancing. A stuck member completes nothing: its position is parked. Only the slow case is answered with capacity.

open as a page

Which levers actually raise a reader group's drain rate, and where does adding more reader instances stop helping?

level: middleimportance: must knowfreq 66%

basics

~20 s

Three levers: more readers, more records per request, and temporarily waiving in-order handling. Extra readers stop helping at the parallelism ceiling the stream's shape fixes, or wherever the shared downstream system saturates first — both walls arrive sooner than most plans assume.

open as a page

Why does a reader group stop making progress while its shares are being reassigned, and how much of it stops?

level: middleimportance: must knowfreq 60%

basics

~20 s

A share may be held by only one member at a time, so the old holder must stop before the new one starts. How much of the reader group stops depends on the design: many stop every member until the division settles, others stop only the shares that move.

open as a page

A reader group's gap is flat, rising steadily, or sawtoothing over a one-hour observation window — what does each shape mean?

level: middleimportance: must knowfreq 58%

basics

~20 s

Flat means arrival and drain rates match, at whatever height. A steady rise means the drain rate is below the arrival rate, and the slope is the difference. A sawtooth means arrival or reading happens in bursts rather than continuously.

open as a page

Why can a reader group's recorded read position vanish after an idle weekend while every record it had not read is still there?

level: middleimportance: must knowfreq 55%

basics

~20 s

A stored read position and the records it points at live under two independent lifetimes. Many platforms discard the position after the owning reader group has been inactive long enough, while the records stay inside the retained window.

open as a page

A reader group is reassigned every ninety seconds and pauses about twenty each time, with no deploys running. What is happening?

level: seniorimportance: must knowfreq 56%

basics

~20 s

Handling is slower than the progress deadline, so a live member is declared gone; its share is reassigned, the group pauses, the member rejoins and starts the same slow work, and the cycle repeats. The group loses most of its time to the loop, not to the work.

open as a page

Before rewinding a reader group's read position by six hours to reprocess, what side effects must an operator account for first?

level: seniorimportance: must knowfreq 58%

basics

~10 s

Every effect downstream of the new position fires again — including the ones that leave the system, such as customer email, payments, third-party calls and records republished onward, which no reader can take back.

open as a page

A reader is frozen on one record it cannot handle: what are the exits, and what does each one cost in correctness?

level: seniorimportance: must knowfreq 52%

basics

~20 s

Four exits: fix the handler, set the record aside, discard it, or wait for a transient cause to clear. Fixing loses nothing but is slowest; discarding is fastest and loses the record silently. Adding capacity is not an exit.

open as a page

Why does adding a reader to a group pause work on some brokers and cost nothing on others?

level: middleimportance: should knowfreq 44%

basics

~20 s

Because only one of the two shapes assigns anything. Where a stream is split into parts and each part has one holder, a new member forces the division to be recomputed. Where readers compete for records from a shared queue, nothing is held, so nothing is reassigned.

open as a page

Why is every member of a reader group stopped before its recorded read position is moved, and what breaks if one keeps running?

level: middleimportance: should knowfreq 52%

basics

~20 s

A running member holds its own place in memory and writes it back, so it overwrites the move or carries on past it. Stopping the whole group first makes the stored position the only writer of record.

open as a page

On a reader's assigned share, does one record that cannot be handled block the records behind it or only itself?

level: middleimportance: should knowfreq 55%

basics

~20 s

It depends on the unit of recorded progress. Where a share's progress is one advancing read position, the record blocks everything behind it on that share. Where each record is acknowledged on its own, it blocks only itself.

open as a page

A reader group was down six hours on a stream taking 10,000 records a second and can now drain 12,000 — when is its backlog of unread records gone, and what could still be lost?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Thirty hours: 216 million unread records divided by a 2,000-a-second surplus. Nothing is lost while that surplus holds and the six-hour unread age sits inside the retained window — loss arrives only if the surplus goes to zero or negative and the oldest unread records age out unhandled.

open as a page

A reader group's total unread count is low and flat, yet one downstream report is hours stale — what reading was missing?

level: seniorimportance: should knowfreq 47%

basics

~20 s

The group total is a sum, and a sum hides its distribution. One share of the work can be hours behind while the others are current, leaving the total small. The missing reading is the per-share breakdown, and specifically the maximum unread age across shares.

open as a page

On a broker that stores no read position and deletes each record once acknowledged, what does resetting forward mean, and is there a way back?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Forward means purging the waiting records — discarding them for every reader at once — because there is no bookmark to move. There is no way back: reprocessing requires a second copy written when the records were first produced.

open as a page

After a restart with no valid stored read position, a reader group silently skips a day of records, so which rule decided that and why was no error raised?

level: seniorimportance: should knowfreq 45%

basics

~20 s

The start-from rule ran. A reader with no valid position begins where that rule says, typically the newest record, the oldest retained record or a point in time, and beginning at the newest skips everything between. Applying a configured default is not an error.

open as a page

Which catch-up levers should an estate pre-authorise for whoever is on call, and which should still require the stream's owner?

level: principalimportance: should knowfreq 38%

basics

~20 s

Pre-authorise the levers that change only cost and latency — adding members to the ceiling, enlarging batches within a stated bound. Anything that changes what the system computes, such as waiving in-order handling, needs a per-stream answer recorded by its owner before the incident, not improvised during one.

open as a page

On a broker that deletes each record once it is acknowledged and stores no read position, what plays the role of lag?

level: middleimportance: nice to knowfreq 40%

basics

~20 s

Queue depth — the count of records still waiting to be handed out — plus the age of the oldest waiting record. It is the same signal read without a position to subtract from, so it describes the whole queue rather than any one reader.

open as a page

Across an estate, how would you set stored position lifetime against the retained window, and what changes where readers keep the position themselves?

level: principalimportance: nice to knowfreq 26%

basics

~20 s

Make the stored position lifetime at least as long as the retained window plus the longest reader outage you intend to survive unaided. Where the reading side keeps its own position, that store's cleanup and restore policy becomes the lifetime, and someone must own it.

open as a page

Across an estate of streams, what should a standing policy say about who may discard a record to unblock a frozen reader, and what must be captured first?

level: principalimportance: nice to knowfreq 34%

basics

~20 s

Settle it before the incident: classify each stream by tolerable loss, pre-authorise the cheap exits, require the record to be captured durably before any discard, name an owner for anything set aside, and set a deadline before the retained window decides for you.

open as a page