advanced 3 min answer

A topic holds 7 days of order events - about 3 TB of unique data at replication factor 3 across 24 partitions on 6 brokers - and has exactly one consumer group, which always reads within a second of the tail. At 14:00 an analytics team deploys a second consumer group configured to start from the earliest offset. What happens over the next hour?

event-drivenpage-cachebackfillquotasread-amplification
Show the full answer Hide the answer

Second by second, what happens

Before 14:00, the brokers do almost no disk reading. The only reader is at the tail, so every fetch is served from the operating system's page cache, and 1.5 TB on each broker's disks is cold and irrelevant. With, say, 64 GB of RAM per broker, roughly 40 GB of page cache holds something like the last half hour of the log.

At 14:00 the new group's first fetch misses the cache and becomes a disk read, and it keeps missing, because it is reading a week-old offset and walking forward. Two things now happen at once.

The historical reads fill the page cache with cold data, evicting the tail that the production consumer and the replication fetchers were being served from. Within minutes the existing consumer's reads also start going to disk. And the brokers' request-handler threads are now occupied with large historical fetches, competing with produce requests for the same fixed thread pool.

Where it amplifies

Produce latency rises, and that is how an analytics deployment reaches the checkout path. Application threads block in send, the producer's in-flight buffer fills, and callers upstream of it start timing out.

Replication fetchers are competing for the same disk and threads, so follower lag grows. If a follower falls out of the in-sync replica set and min.insync.replicas is 2, produces with acks=all fail outright — the first hard error of the incident, and it appears in the order service, not in analytics.

The existing consumer is now in a self-reinforcing loop: it falls behind, so it reads from further back, so it misses the cache more, so it falls further behind.

What the user sees

Checkout slows and some orders fail to be written, starting roughly ten minutes after a deployment that changed nothing in the write path and whose owners are not on the order service's on-call rotation.

What stops it

  • A client-level fetch byte-rate quota on the new consumer group. This is the mechanism: it bounds historical read throughput at the broker, where the contention is. Everything else is advice.
  • Move historical reads off the brokers — tiered storage, or a backfill from a copy of the log in object storage, so the week-old data is not served by the same device as live traffic.
  • Start from a snapshot plus the last few hours, which is usually what the analytics team actually needed.
  • Do not rely on it finishing quickly. 3 TB at an aggregate 600 MB/s is about 80 minutes of saturated reading, so this does not self-heal before the in-sync replica set is at risk. It would self-heal only if the backfill were short enough to finish inside the replication lag budget.

When this is the wrong answer

If the whole topic fits in page cache — a topic of 20 GB against 40 GB of cache per broker — a new group reading from earliest costs almost nothing, finishes in seconds, and quotas are needless ceremony. The threshold to carry is retained bytes per broker against page cache per broker: once the log is an order of magnitude larger than the cache, every new consumer group starting from earliest is a planned load test on the production write path.