skip to content

Sizing from Throughput

Turning a measured write rate, record size and fan-out into node, disk and network capacity instead of a node count somebody guessed. Asked because the headroom nobody reserved is what fails first.

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

questions

4

A stream takes 20,000 records per second averaging 2 KB and four independent reader sets read every record - what bandwidth must the cluster carry?

level: middleimportance: must knowfreq 62%

answer

  1. a rate is not a size
  2. records per second times bytes
  3. the read side is its own number
  4. count independent reader sets
  5. fan-out multiplies egress, not ingress

basics

~20 s

Rate times average record size gives the write side: 20,000 x 2 KB is about 40 MB per second in. Each independent reader set that reads every record adds the same again, so four make roughly 160 MB per second out.

solid answer

~50 s

Size the write side first: sustained rate multiplied by average record size. Twenty thousand records a second at 2 KB is about 40 MB/s of ingress. The read side is a separate number, and where a stream is retained and consumed in full by several independent reader sets it is fan-out times ingress - four sets means roughly 160 MB/s of egress, which is why the outbound interface usually saturates before the inbound one. Platform shape decides whether that multiplier exists at all: where consumers compete for one shared work list, total read bytes stay near the write rate however many attach. Keep the record rate too, not just the byte rate, because per-record work lands on request handling rather than on the interfaces. Both numbers are then scaled by the peak-to-average ratio of the busiest interval, and the bytes copies cost between nodes are added on top.

go deeper

for a junior

Recall that a cluster is sized in bytes per second, not in records per second, and that you get bytes by multiplying the rate by the average record size. Knowing that reading costs bandwidth too is already worth saying.

for a middle

Explain both sides of the arithmetic: ingress as rate times size, egress as that multiplied by how many independent reader sets consume every record. Say why the outbound interface is usually the one that saturates first.

for a senior

Show where the averages lie. Check the size distribution, name the reader sets nobody counted, and say which resource each input lands on so the node count comes from the worst ratio rather than an average of them.

for a principal

Frame it as which number the organisation will be held to. Decide whether fan-out is bounded by policy or allowed to grow with every new consumer, because an unbounded read multiplier turns every new reader into a hardware purchase.

## What sizing actually starts from A record rate on its own is not a size. **Three measured inputs** turn it into one: - **sustained write rate** - records accepted per second, averaged over an interval that represents ordinary operation rather than a number from a design document; - **average record size** - the bytes of a record as it lands, taken over that same interval; - **fan-out** - how many independent reader sets consume each record end to end. Two more scale the result: the **peak-to-average ratio** of the busiest interval, and - where a byte total rather than a byte rate is wanted - the **retention span** the records are kept for, which is taken here as a given. ## The write side Rate multiplied by average size is the raw ingress the cluster must accept: ``` sustained write rate 20,000 records/s average record size 2 KB raw ingress 20,000 x 2 KB = 40 MB/s ``` Two things are worth keeping separate even at this first step. The **byte rate** lands on network interfaces and on the storage the append path writes to. The **record rate** lands somewhere else: on request handling, on per-record bookkeeping, and on whatever per-record index or metadata the platform maintains. A stream of 200,000 tiny records a second and a stream of 2,000 large ones can carry similar bytes and cost very different amounts of processor, which is why both numbers are carried forward instead of being collapsed into one. The average deserves a look at its distribution as well. A mean of 2 KB produced by a population that is mostly 500 B with a daily tail of 5 MB sizes buffers, request ceilings and the busiest interval quite differently from a population that genuinely is 2 KB. ## The read side, which is where fan-out lives Read demand is a separate number, and on a retained stream with several independent consumers it is the larger of the two. **Where each independent reader set reads every record**, egress is fan-out multiplied by ingress: ``` independent reader sets 4 raw egress 4 x 40 MB/s = 160 MB/s ``` **Platform shape decides whether that multiplier exists at all.** Where consumers compete to drain one shared work list, adding consumers divides the work rather than duplicating it, and total read bytes stay near the write rate however many attach. Where a stream is retained and split into partitions, both models are present at once: several independent reader sets each read the whole stream, while inside one set the partitions are divided among its members. The arithmetic follows the model, so establishing which model is in play is the first move, not a detail. Where those reads are served from matters as much as their volume. Readers that stay near the newest record are usually served from memory the records are still sitting in - the operating system's file cache on platforms that lean on it, in-process memory on platforms that keep records inside the process, remote storage on platforms that offload older data. **Readers that are far behind, or replaying from the start, are served from the volume instead**, which turns what looked like a network number into a disk-read number. ## What each input lands on | measured input | what it multiplies | the resource it sizes | |---|---|---| | sustained write rate | records per second | request handling, per-record bookkeeping | | average record size | bytes per second in | ingress network, the append path on the volume | | fan-out across independent reader sets | bytes per second out | egress network, file cache, reads from the volume | | copies kept per record | bytes between nodes, bytes written | internal network, total volume written | | peak-to-average ratio | every rate above | every resource above | | retention span (given) | bytes stored | volume size | ## From bandwidth to machines Each demand is divided by what one record-serving node honestly sustains for that resource - not what a specification sheet claims, but what it holds with the append path, its copies and its readers all active at once. The node count is the **worst** of those ratios, never their average: a cluster with spare processor and a saturated outbound interface is a saturated cluster. Only then is the design target applied, because a count derived exactly from these numbers leaves nothing for the day a reader replays, a writer retries or a node dies. ## Where the arithmetic usually goes wrong 1. **Sizing from records per second alone**, then discovering the records are ten times larger than anybody assumed. 2. **Assuming read demand equals write demand.** With four independent reader sets it is four times larger, and the interface that gives way is the outbound one. 3. **Counting only the reader sets that exist today** and forgetting the reporting job that re-reads the stream every night. 4. **Taking one average across a whole day**, which hides the interval in which the cluster actually fails.

  • Why keep the record rate after you have the byte rate?
    Because they land on different resources. Bytes per second size interfaces and the volume; records per second size request handling and whatever per-record bookkeeping the platform keeps. Two streams carrying identical bytes, one as many tiny records and one as few large ones, cost very different amounts of processor, and a cluster can run out of per-record work long before it runs out of bandwidth.
  • A nightly job re-reads the whole stream. How does that enter the numbers?
    As an extra reader set on the read side, but not one served the same way. Live readers are usually served from memory the records still sit in; a job replaying from the start reads mostly from the volume, so it adds disk-read demand rather than cache hits. Size it as its own fan-out contribution during the interval it runs.
  • Does fan-out change the storage number?
    No. Stored bytes are the sustained write rate integrated over the retention span, multiplied by how many copies are kept. Reading a record many times does not store it again. Fan-out moves the egress network and the read path; it leaves the volume size alone.

saying these in an interview costs you the question

  • Quotes a record rate as though it were a bandwidth number
  • Assumes read bandwidth equals write bandwidth
  • Thinks fan-out divides the bytes among the readers rather than multiplying them
  • Sizes from a mean record size without looking at the tail
  • Counts only the reader sets running today
  • Believes reading a record more often increases stored bytes
open as a page

A cluster sized from a daily average of 8,000 records per second sees 40,000 in its busiest ten minutes - which number should have sized it?

level: middleimportance: should knowfreq 55%

basics

~20 s

The busiest interval, not the daily mean. A peak-to-average ratio of five means the machines must carry 40,000 records per second, while the volume size still comes from the sustained rate accumulated over the retention span.

open as a page

Every record on a stream is kept on three record-serving nodes - why does the network between them carry more than the writers send?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Because a write is not finished when it is accepted. Where one node takes the write and sends it to the other two, internal traffic is about twice client ingress, and every stored copy also writes the full record bytes to its own volume.

open as a page

A broker cluster is sized exactly to its measured peak rate - what routinely arrives that the machines have no room for, and what should the design target have been?

level: principalimportance: should knowfreq 44%

basics

~20 s

Three things reliably arrive on top of measured traffic: a reader replaying from the start, a writer retry surge re-sending records already accepted, and a dead node's share landing on its surviving peers. The design target is a fraction of what the machines can do, not all of it.

open as a page