2.5 Selecting hardware
Other factors: specific drive technology (SAS vs SATA) and the quality of the drive controller.
Retention sizing
- 29531 GB
- raw
- 88594 GB
- replicated
- 14766 GB
- per broker
- 24609 GB
- provision
41 hours to rebuild a replacement broker at ~100 MB/s of replication traffic. That is how long you run without full redundancy after a disk failure. More brokers with smaller disks shortens it.
"Selecting an appropriate hardware configuration for a Kafka broker can be more art than science." Kafka has no strict hardware requirement and runs fine on most systems. Once performance matters, the bottlenecks are:
| Dimension | What it governs |
|---|---|
| Disk throughput | producer latency |
| Disk capacity | retention window |
| Memory | consumer performance (page cache) |
| Networking | total cluster throughput ceiling |
| CPU | matters only at very large scale |
| Partition count | metadata update volume per broker |
5.1 Disk throughput → producer latency
The causal chain:
SSD vs HDD:
| SSD | HDD | |
|---|---|---|
| Seek/access time | Drastically lower → best performance | Higher |
| Cost per unit capacity | Expensive | More economical, more capacity |
| Improve by | — | more of them per broker: multiple data directories, or RAID |
Other factors: specific drive technology (SAS vs SATA) and the quality of the drive controller.
The observed rule of thumb:
HDD → clusters with very high storage needs but data not accessed often. SSD → when there is a very large number of client connections.
5.2 Disk capacity → retention
needed = daily_ingest × retention_days × (1 + replication overhead)
- + 10% overhead for other files
- + buffer for traffic fluctuation and growth
Example from the book: 1 TB/day traffic × 7 days retention = minimum 7 TB usable for log segments, plus at least 10% overhead, plus a growth buffer.
Total cluster traffic is balanced by multiple partitions per topic, letting additional brokers augment capacity when single-broker density is insufficient. The requirement is also shaped by the replication strategy (Ch. 7).
5.3 Memory → consumer performance (this is the subtle one)
The normal mode of operation: a consumer reads from the end of the partition, caught up, lagging very little or not at all. In that state:
The messages the consumer is reading are optimally stored in the system's page cache, resulting in faster reads than if the broker had to reread from disk. Therefore more memory available for page cache improves consumer performance.
“More memory available for page cache improves consumer performance.” Kafka itself needs little heap — “even a broker handling 150,000 messages/second and 200 megabits/second can run with a 5 GB heap.” Hence the colocation rule: this is the main reason it is not recommended to colocate Kafka with any other significant application — it would have to share the page cache, which decreases consumer performance for Kafka.
Heap requirements are tiny by comparison:
Kafka itself does not need much heap. Even a broker handling 150,000 messages/second and 200 megabits/second can run with a 5 GB heap. The rest of system memory is used by page cache and benefits Kafka by caching log segments in use.
Hence the colocation rule:
This is the main reason it is not recommended to colocate Kafka with any other significant application — it will have to share the page cache, which decreases consumer performance for Kafka.
Restated: Kafka's performance model is "the OS is my cache." Any neighbor process that touches a lot of memory is stealing your read path.
5.4 Networking → the throughput ceiling
Available network throughput specifies the maximum amount of traffic Kafka can handle. Combined with disk storage, it's a governing factor for cluster sizing.
The inherent asymmetry:
- Inbound: producer writes 1 MB/s for a topic.
- Outbound: 1 MB/s × (number of consumers) + replication traffic (Ch. 7) + mirroring traffic (Ch. 10).
⇒ Outbound can be many multiples of inbound. And should the network interface become saturated, “it is not uncommon for cluster replication to fall behind, which can leave the cluster in a vulnerable state.” Recommendation: at least 10 Gb NICs — older machines with 1 Gb NICs “are easily saturated and aren’t recommended.”
The dangerous consequence:
Should the network interface become saturated, it is not uncommon for cluster replication to fall behind, which can leave the cluster in a vulnerable state.
Read that as: network saturation converts into a durability risk, because under-replicated partitions mean you're one failure from data loss.
Recommendation: at least 10 Gb NICs. "Older machines with 1 Gb NICs are easily saturated and aren't recommended."
5.5 CPU → only at extreme scale
Processing power is not as important as disk and memory until you scale very large.
Where the CPU actually goes — and it's a specific, non-obvious path:
This decompress → validate → recompress cycle is where most of Kafka’s CPU requirement comes from. Still not the primary hardware factor unless clusters become very large — hundreds of nodes and millions of partitions in a single cluster — at which point better CPU can reduce cluster size.
Not the primary hardware factor unless clusters become very large — hundreds of nodes and millions of partitions in a single cluster. At that point better CPU can reduce cluster size.
(This is also the mechanism behind zero-copy being defeated: if the broker must recompress, it can't just sendfile() bytes through. Matching client and broker compression codecs is how you keep this cheap.)
5.6 Kafka in the cloud
General method: prioritize Kafka's performance characteristics, then pick the instance shape — each instance type is a different mix of CPU, memory, IOPS, and disk.
Microsoft Azure
- Disks are managed separately from the VM → storage needs are decoupled from VM type.
- Decision order: start with data retention required, then producer performance needed.
- Very low latency needed → I/O optimized instances with premium SSD.
- Otherwise → Azure Managed Disks or Azure Blob Storage may suffice.
| Cluster size | Instance |
|---|---|
| Smaller clusters, most use cases | Standard D16s v3 |
| High performance / larger clusters | D64s v4 |
- Build the cluster in an Azure availability set and balance partitions across Azure compute fault domains to ensure availability.
- Strongly prefer Azure Managed Disks over ephemeral disks — "If a VM is moved, you run the risk of losing all the data on your Kafka broker."
| Storage | Cost | SLA |
|---|---|---|
| HDD Managed Disks | relatively inexpensive | no clearly defined availability SLA from Microsoft |
| Premium SSD / Ultra SSD | much more expensive | much quicker, 99.99% SLA |
| Microsoft Blob Storage | — | option if not latency sensitive |
Amazon Web Services
- Very low latency needed → I/O optimized instances with local SSD.
- Otherwise → ephemeral storage (e.g. Amazon EBS) may suffice.
| Instance | Trade |
|---|---|
| m4 | greater retention periods, but lower disk throughput (elastic block storage) |
| r3 | much better throughput (local SSD), but drives limit retained data |
| i2 / d2 | best of both worlds, but significantly more expensive |