metric

Replay Throughput Ceiling

also called Reprocessing Rate Limit, Backfill Throughput Budget

The rate at which retained history can actually be re-read and reprocessed - set by the slowest storage tier and any configured quota rather than by retention - which decides whether a correction run takes hours or a week.

replaytiered-storagebackfillobject-storageplanning

Retention is quoted as a capability: "we keep 90 days, so we can always reprocess". The number that decides whether that is true is a different one. How fast can those 90 days actually be read back?

A platform ingesting 200 MB/s keeps 90 days. Local disk would serve a replay at a large multiple of ingest, so the team projects 3 hours for a 30-day correction run. The segments older than six hours live in object storage, where reads do not come from the broker's page cache, arrive a whole segment at a time, and pass through a path that is rate-limited on purpose. The run takes 14 hours.

The ceiling is a property of the architecture and often of a single configuration value, and almost nobody measures it until the day it is load-bearing.

Why it matters

Every argument for a log-centric design rests on replay: a bug is survivable because history can be reprocessed, a new consumer can be added because it can bootstrap, a schema change is safe because the old bytes are still there. Those claims are all claims about throughput, and they are made with retention numbers.

It also decides the shape of incident response. A correction that completes in 3 hours is run during the incident. One that takes 14 hours is a scheduled project with a communications plan, a decision about whether to serve wrong data meanwhile, and a second decision about whether to throttle the replay so it does not starve live consumers.

Implementation patterns

  • Measure it, publish it next to retention. "90 days retained and replayable at 400 MB/s, so a full replay is 14 hours." Re-measure after any storage change.
  • Keep local retention as long as the replay you run routinely. If a weekly correction covers 48 hours, hold 48 hours on broker disk and let the remote tier serve the rare deep reads.
  • Know the quota. Kafka's tiered storage quotas (KIP-956, in the 3.9 release where tiered storage reached general availability) expose remote.log.manager.fetch.max.bytes.per.second as a global broker cap across every partition reading remotely. Unlimited by default, and set by anyone who has watched a replay evict the page cache from under live consumers.
  • Give replay its own capacity. A separate consumer group, its own compute, and a throttle, so catch-up cannot become a second incident.
  • For full-history recomputation, leave the log. Export once to a columnar table on object storage and reprocess there with warehouse parallelism rather than pulling petabytes back through brokers built for streaming.

Industry example

Kafka's tiered storage, specified in KIP-405 and generally available in the 3.9 release, is explicit about its own assumption: the operations documentation describes data as mostly consumed as tail reads, with older data read infrequently for backfill or failure recovery. Two companion proposals in the same release add quotas for remote copy and remote fetch, because the uncapped version degrades producer latency and starves live consumers. The feature was designed to make long retention affordable, not to make long replays fast - and the consequence seen in production is that replay rate is now a number an operator chose.

Failure scenarios

  • The estimate from the wrong tier. Hours projected from local throughput, days delivered from remote.
  • Replay starves production. An uncapped backfill evicts the page cache, so live consumers start reading from disk and lag across unrelated topics.
  • The silent quota. Somebody set the cap during a previous incident. Months later a replay runs at a fraction of the expected rate and nobody connects it to a broker config.
  • Retention that is not replayable at all. History exists but the oldest schema cannot be deserialised by current code, so the ceiling is effectively zero for the part of the window that matters.

Trade-offs

Cheap retention and fast replay pull in opposite directions. Object storage is an order of magnitude cheaper per GB-month and gives up page cache, per-fetch latency and parallelism. Local disk gives the throughput and charges for provisioned capacity times the replication factor. Choose by how often you actually read old data, and keep the two windows separately configured so the choice stays reversible.

When not to use it

As a metric it is wasted effort on a topic nobody replays - compliance archives, audit trails kept because a regulator requires them. Measure it where a replay is part of the recovery plan, and where a correction run has an actual deadline. For a small topic whose whole retention fits on broker disk, the ceiling is high enough not to be worth a dashboard.

Interview question

Q: You inherit a platform that advertises 90-day retention as its recovery story. What is the first number you ask for, and how would you turn the answer into an operational commitment?

What a strong answer covers: asking for measured replay throughput from the tier the data actually sits on, not retention; converting it into a time for the longest realistic correction; checking whether a remote-fetch quota is set and who set it; separating local from remote retention so routine reprocessing never touches object storage; confirming current code can deserialise the oldest retained schema; and writing the resulting hours into the recovery plan so nobody promises three.

Quick check

Quiz: Why can 90-day retention be useless as a recovery story? Because the recovery time is set by replay throughput from the tier the data sits on, which can be several times slower than local disk and is often capped by a configured quota.

Flashcard: Which number should always be published next to retention? - Measured replay throughput, converted into the wall-clock time of a full replay.