skip to content

How do you estimate the network bandwidth a Kafka cluster needs, including replication, from throughput metrics?

level: seniorimportance: should knowfreq 45%

answer

  1. replication out = BytesIn × (RF-1)
  2. ReplicationBytesIn/OutPerSec metrics
  3. BytesOut already includes consumer fan-out
  4. catch-up/reassign bursts > steady state
  5. leader/follower.replication.throttled.rate

basics

~20 s

Outbound bandwidth per broker ≈ consumer egress (BytesOutPerSec, multiplied by consumer fan-out) plus replication egress to followers (BytesInPerSec × (replicationFactor − 1)). Inbound ≈ producer ingress plus replica fetch traffic. Replication is often the dominant, easily-forgotten term.

solid answer

~40 s

Model each NIC direction separately. Producer ingress = BytesInPerSec. Replication adds the biggest hidden cost: a leader must ship every produced byte to (RF−1) followers, so leader outbound replication ≈ BytesInPerSec × (RF−1), tracked as ReplicationBytesOutPerSec; correspondingly a broker's inbound includes replica-fetch traffic for partitions it follows, ReplicationBytesInPerSec. Consumer egress = BytesOutPerSec, which already embeds fan-out (each group re-reads the data). So per broker: outbound ≈ consumer_egress + replication_out; inbound ≈ producer_in + replication_in. On a balanced cluster each broker leads ~1/N of partitions and follows others, so totals spread out, but you must size NICs for the busiest broker and for catch-up bursts (a recovering broker or reassignment saturates links, which is why follower.replication.throttled.rate / leader.replication.throttled.rate exist). Always size for peak BytesIn with replication, then add headroom.

go deeper

for a junior

Know replication copies data to other brokers and adds network traffic beyond produce/consume.

for a middle

Compute replication out as BytesIn × (RF−1) and separate inbound vs outbound.

for a senior

Account for fan-out in BytesOut, hottest-broker sizing, and catch-up bursts with replication throttles.

for a principal

Architect NIC/topology and throttle policy across clusters, modeling failure/reassignment scenarios against link budgets and SLOs.

**Why network is usually the first bottleneck.** Disk is large and cheap; the NIC is fixed. Replication and consumer fan-out *multiply* the base produce rate, so a modest ingress can produce large network demand. **The traffic components.** Define base ingress B = BytesInPerSec (compressed on-disk bytes/sec) for the partitions a broker leads. Replication factor RF. *Per leader broker, outbound:* - **Replication out:** the leader must send each produced byte to every follower, i.e. `B × (RF − 1)`. Metric: `ReplicationBytesOutPerSec`. - **Consumer egress:** `BytesOutPerSec`. This already includes consumer fan-out — if K consumer groups read the topic, the same bytes go out ~K times, so BytesOut can be several times BytesIn. - Total outbound ≈ `B×(RF−1) + BytesOutPerSec`. *Per broker, inbound:* - **Producer ingress:** `BytesInPerSec` for led partitions. - **Replica fetch (as a follower):** for partitions this broker follows, it pulls data from their leaders — `ReplicationBytesInPerSec`. - Total inbound ≈ producer_in + replication_in. **Cluster vs per-broker.** On a well-balanced cluster of N brokers, leadership and followership spread so that aggregate replication traffic across the cluster is `total_produce × (RF−1)`, divided roughly evenly. But you size hardware for the **hottest broker** plus failure scenarios, not the average. **Catch-up and reassignment bursts — the dangerous edge.** Steady state is the easy part. When a broker restarts or you reassign partitions, followers fetch from the *earliest needed offset* as fast as the link allows, briefly demanding far more than steady-state replication bandwidth and competing with live consumer/produce traffic. Kafka provides replication throttles — `leader.replication.throttled.rate` and `follower.replication.throttled.rate` (set via kafka-configs / kafka-reassign-partitions --throttle) — precisely to cap this so catch-up doesn't starve the workload. The metrics `ReplicationBytesInPerSec`/`OutPerSec` and the throttle metrics let you see whether replication is approaching link limits. **Worked example.** Topic at B = 100 MB/s ingress, RF=3, read by 2 consumer groups (fan-out 2), balanced across 5 brokers. - Cluster replication out = 100 × (3−1) = 200 MB/s. - Cluster consumer egress = 100 × 2 = 200 MB/s. - Cluster producer in = 100 MB/s; replica fetch in = 200 MB/s. - Total cluster outbound ≈ 400 MB/s ⇒ ~80 MB/s per broker outbound average; size the NIC (and headroom) for the peak broker, which on a 25 Gbps (~3.1 GB/s) link is comfortable but on 1 Gbps (~125 MB/s) would be tight once bursts and skew are added. **Key takeaways / gotchas:** - The `(RF−1)` replication multiplier is the most commonly forgotten term; people size for produce + consume and get surprised. - BytesOutPerSec already includes fan-out; don't multiply it by group count again. - Replication catch-up can momentarily dwarf steady-state — provision headroom and use throttles. - Inter-broker traffic (replication) and client traffic may share or split NICs depending on `inter.broker.listener.name`; if they share one NIC, sum both directions against the same link budget.

  • Why is the (RF−1) replication term so often underestimated?
    Teams size for produce + consume and forget the leader ships every byte to each follower. At RF=3 replication alone equals twice the ingress, often the single largest network term.
  • Why might replication bandwidth spike far above steady state?
    When a broker restarts or partitions are reassigned, followers fetch backlog as fast as the link allows to catch up, briefly demanding much more than steady-state replication — hence replication throttles to cap it.

saying these in an interview costs you the question

  • Omitting the (RF−1) replication multiplier from network sizing.
  • Multiplying BytesOutPerSec by consumer-group count (fan-out is already in it).
  • Assuming replication traffic equals steady-state during broker recovery/reassignment.
  • Sizing NICs for average rather than the hottest broker plus burst headroom.

context