skip to content

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%

answer

  1. ingress is the smallest number
  2. a write travels after it is accepted
  3. one send per extra copy
  4. disk write is copies times ingress
  5. rebuilds move retained bytes, not rates

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.

solid answer

~50 s

A write arrives once from a client and then has to reach the other copies. On platforms where one node accepts the write for a partition and forwards it, internal traffic is roughly `(copies - 1)` times ingress - with three copies, twice what the writers sent - and it shares interfaces and volumes with client traffic. Cluster-wide disk write demand is `copies x ingress`, since every copy stores the whole record, not a fragment of it. Designs differ in where that cost lands: some have the writer path send to several members directly, and some keep one shared copy in remote storage, which moves the multiplier onto the storage service's network instead of between nodes. Recovery is the part nobody sizes: rebuilding one lost copy streams a partition's entire retained bytes while normal traffic continues.

go deeper

for a junior

Recall that storing a record more than once means sending it more than once, so the cluster's own network carries more than the writers do. The exact multiplier can come later.

for a middle

Do the arithmetic: with three copies and one node forwarding, internal traffic is about twice ingress and cluster-wide disk write is three times it. Say that both share interfaces with client traffic.

for a senior

Size per node, not only per cluster, and include rebuild traffic - a lost node's copies mean moving every retained byte it held while the cluster is already short. Name which replication shape you assumed.

for a principal

Own the consequence: the copy count is a durability decision whose bill arrives as network and volumes, and whether that bill lands between your nodes or on a storage service is an architectural choice with very different cost curves.

## What a write costs after it has been accepted Client ingress is the smallest of the numbers in this question. A record that is stored more than once has to travel to wherever the other copies live, and that travel is between record-serving nodes, on the same interfaces that serve clients, competing with them. On the common shape - one node accepts writes for a partition and passes them to the nodes holding the other copies: ``` client ingress 40 MB/s copies kept per record 3 internal send (3 - 1) x 40 MB/s = 80 MB/s cluster-wide disk write 3 x 40 MB/s = 120 MB/s ``` So a cluster that clients push 40 MB/s into is moving roughly 120 MB/s across its own network and writing 120 MB/s to its volumes, before a single reader is served. The copy count is not being chosen here - it is a given, decided on durability grounds - but its byte consequences are exactly what sizing has to carry. ## Designs differ in where the multiplier lands This is the part to state carefully, because the leader-forwards shape is only one of several: - **One node forwards to the others.** Internal traffic is about `(copies - 1) x ingress`, concentrated on whichever nodes are accepting writes for the busiest partitions. - **The writer path sends to several members itself.** The multiplier moves onto the client-facing interfaces and the internal network carries less, but the ingress side of each node carries more. - **Storage is shared rather than per-node.** Where records are written once to remote or shared storage, there is little or no per-record copy traffic between nodes at all; the same multiplier reappears as traffic to the storage service, and as its charges. - **Copies are not all full copies.** Some designs keep a cheaper representation of older data, so the multiplier applies in full only to the recent part of the stream. The habit worth keeping is to ask *where does the second copy's bytes physically travel* rather than assuming the answer. ## The read side shares those interfaces Internal copy traffic does not have an interface to itself. On one node the outbound interface may be carrying, at the same moment: records being sent to peers for copies, records being served to readers, and records being sent to a peer that is catching up. The per-node interface has to hold the sum, which is why the arithmetic is done per node and not only cluster-wide. | traffic on a record-serving node | roughly how much | where it goes | |---|---|---| | client writes in | its share of ingress | inbound interface, then the append path | | copies sent to peers | `(copies - 1)` x its share of ingress | outbound interface, internal network | | copies received from peers | its share of others' ingress | inbound interface, then its own volume | | records served to readers | fan-out x its share | outbound interface, file cache or volume | | a peer rebuilding a copy | as fast as it is allowed | outbound interface, reads from the volume | ## Recovery is the traffic nobody sizes When a node is lost and its copies are rebuilt elsewhere, the cluster has to move **every retained byte of every partition that node held**, not the current rate. On a cluster holding days of records that is a large multiple of a second's traffic, and it arrives precisely when the cluster is already one node short. Two consequences follow, and both are sizing consequences rather than operational ones: 1. The interfaces have to leave room for it, or a recovery slows the live traffic it was meant to protect. 2. The time a rebuild takes is a *derived* number - retained bytes per node divided by the rate the rebuild is permitted - and it belongs in the sizing write-up, because it is how long the cluster runs with fewer copies than it is supposed to have. Whether that rebuild traffic can be limited, and at what granularity, varies between platforms; what does not vary is that an unbounded rebuild and a saturated interface are the same incident. ## Pulling it together For a per-node interface, size for the sum of its shares: client writes in, copies out to peers, copies in from peers, reads out to consumers, and a rebuild. Cluster-wide, remember that the volume total is `copies x sustained rate x retention span`, and that this is one of the few places where the copy count multiplies two different resources at once - bytes moved and bytes stored. A sizing exercise that stops at client ingress will be out by a factor of three on a three-copy cluster, and it will be out in the direction that fails.

  • Does the copy multiplier apply to the read side too?
    Not in the same way. Readers are normally served by one node per partition, so fan-out multiplies egress while the copy count does not. The copy count multiplies bytes moved between nodes and bytes written to volumes. The exception is a platform that allows reads from a non-leading copy, which spreads the same read bytes across more nodes rather than creating more of them.
  • Why is rebuild traffic a sizing input rather than an operational detail?
    Because it decides how long the cluster runs short of copies. Rebuild time is retained bytes per node divided by the rate the rebuild is allowed, and that rate has to fit in interfaces already carrying live traffic. Size without leaving room for it and the choice becomes: slow the recovery, or slow the customers.
  • Where does the multiplier go on a platform with shared storage?
    Onto the storage service. If records are written once to remote or shared storage rather than to a volume per node, node-to-node copy traffic largely disappears and reappears as bandwidth to, and charges from, that service. The arithmetic does not vanish; it changes which line of the bill and which network it lands on.

saying these in an interview costs you the question

  • Sizes interfaces from client ingress alone
  • Thinks extra copies store an index rather than the full record bytes
  • Assumes copy traffic has a dedicated network of its own
  • Never accounts for rebuilding a lost node's copies
  • Believes every platform replicates by forwarding from one node
  • Treats a rebuild as bounded by the current write rate