A field guide to throughput

How Kafka swallows seven trillion messages a day

0
messages per day — LinkedIn's production clusters, 2019, and the number has only grown

That is a sustained average of ≈ 81,000,000 messages every second. Not a queue. Not a database. An append-only log, sharded to the horizon. What follows is the anatomy of that throughput — six decisions, one scroll at a time.

prototyped 2010 · open-sourced 2011 · apache top-level 2012
scroll ↓ the log is append-only
Chapter 01 · The Flood

Every pipeline, hand-braided

LinkedIn, 2010. Clickstream into Hadoop. Databases into the warehouse. Metrics into alerting. Every new source × every new sink meant another bespoke connector — its own protocol, its own failure modes, its own 3 a.m. page.

Five sources and five sinks already means 25 integrations. LinkedIn's real inventory ran far beyond that.

complexity grows as M × N — quadratic, and it never stops
01 · step 2 of 3

Replace the mesh with one log

The fix is topological. Every source publishes into a single, durable log; every sink subscribes to it. M × N collapses to M + N, and a new consumer never touches a producer — it just reads, from wherever it likes.

That log became Kafka.

01 · step 3 of 3

Then the numbers got absurd

By 2019 the clusters absorbed seven trillion messages a day — a sustained average near 81 million per second, with headroom above it.

That is not queue territory; it is physics territory. Six design decisions get you there. The first is about lying to a hard disk.

Chapter 02 · The Log

One data structure: the log

Kafka's heart isn't a queue or a table — it's an append-only log. Producers write records to the tail. That's the entire write path: no index updates, no B-tree rebalancing, no in-place edits.

That UPDATE you wanted? Rejected at the door.

02 · step 2 of 4

Offsets: a number is an index

Every record gets a monotonically increasing offset. Reading is therefore fetch(offset = N) — an O(1) jump, not a search.

Immutability buys more than speed: it makes caching trivial and turns replication into a simple game of catch-up.

02 · step 3 of 4

Disks are fast — if you don't make them think

Kafka's designers measured a six-drive SATA array: ~600 MB/s linear writes against ~100 KB/s random ones — a 6,000× gap on the same hardware. Kafka only ever appends, so its "slow disk" streams.

The log converts the expensive kind of I/O into the cheap kind. That's the whole trick of the write path.

02 · step 4 of 4

Old data is the feature

Consuming a record doesn't delete it. Records are retained for days, weeks, years. A new consumer starts at the tail. A buggy consumer rewinds to Tuesday. A new model trains on three months of history.

Storage is cheap. Re-processing is the payoff — and one log still has a ceiling.

Chapter 03 · Shard the Log

One queue has a ceiling

A topic with a single partition is one ordered stream: one writer contending at the head, one reader advancing. Perfect order — and a hard cap on throughput.

At 81 million messages a second, one queue is a traffic jam. So Kafka shards.

03 · step 2 of 4

hash(key) % P

A topic is split into P partitions. The producer hashes each record's key and lands it deterministically: same key → same partition, forever.

That preserves per-key order — every event for user:8123 stays in sequence — while spreading the load.

03 · step 3 of 4

Parallelism is a dial

Partitions live on different brokers, each accepting writes and serving reads independently — so throughput scales roughly linearly with P and with machines.

3 partitions, 6, 12, hundreds… the dial turns up. (The cost: file handles, memory, rebalance time. Dials have detents.)

03 · step 4 of 4

Order is a local promise

Within a partition: total order — offsets 17 → 18 → 19, no ambiguity. Across partitions: none. Records interleave, and no broker promises otherwise.

That's not a bug to fix; it's the price of parallelism. Design your keys like ordering depends on them — because it does.

Chapter 04 · Survive the Crash

Every partition lives in three places

Each partition is replicated — typically three times. One replica is the leader: all writes and reads go through it. The others are followers, quietly tailing the leader's log.

Replication here isn't a backup strategy. It is the availability strategy.

04 · step 2 of 4

acks=all — durability with a receipt

With acks=all, a producer's batch isn't "sent" until every in-sync replica holds it. Set min.insync.replicas=2 and durability survives a broker vanishing mid-write.

Latency cost: a few milliseconds. Insurance you stop noticing.

04 · step 3 of 4

The ISR club

Followers that lag get benched out of the ISR — the in-sync replica set. The quorum keeps acknowledging writes without them, and a laggard can never be elected leader.

Membership is earned continuously, by staying caught up.

04 · step 4 of 4

When the leader dies

No drama: the remaining ISR members elect a new leader in seconds. Committed records are intact, producers reconnect, consumers barely notice.

Kafka's goal isn't to prevent failure. It's to make failure boring.

Chapter 05 · Never Copy Bytes

The CPU is the real toll

At this rate, moving bytes is the expense. Kafka's radical answer: cache nothing itself. Records land in the OS page cache and stay there — hot data is served from RAM without ever entering the JVM heap.

No cache code, no eviction policy, no GC pressure. The operating system is Kafka's cache.

05 · step 2 of 3

The naive path: four copies

Classic read() + write(): disk → kernel buffer → user-space buffer → socket buffer → NIC. Four copies, four context switches — and the JVM did nothing but shuttle bytes it never looked at.

A very expensive ferry service.

05 · step 3 of 3

sendfile(): fire the middleman

One syscall, and the kernel DMAs bytes straight from page cache to the NIC: two copies, both DMA — zero through the CPU. The JVM orchestrates but never touches payload bytes.

It's the same trick web servers use for static files. Kafka serves logs like static files — and spends its CPU on traffic, not ferrying.

Chapter 06 · Readers Keep Pace

The broker never pushes

Consumers pull. The broker never dictates pace, so a slow consumer can't drown the cluster — it just accrues lag, a number you can chart and alert on.

Push systems hide backpressure. Kafka exposes it as data.

06 · step 2 of 4

The consumer group

Consumers band into a group, and Kafka hands each partition to exactly one member: six partitions, three consumers, two each. The group is one logical subscriber that scales horizontally.

Different groups read the same stream independently — like radio stations off one tower.

06 · step 3 of 4

Scale reads the way you scaled writes

Add consumers until there's one per partition; rebalance redistributes automatically. Add a seventh? It idles — parallelism is capped at P.

Another reason partition count is the big dial in the room.

06 · step 4 of 4

Your bookmark is just a number

Each consumer's position is an offset — stored in an internal topic, __consumer_offsets. Crashed? Resume from the number. Shipped a bug? Rewind and reprocess. Deleted a downstream store? Rebuild it from history.

Your place in the stream is state you own — not state the broker hoards.

Chapter 07 · The Arithmetic

The ledger

Nothing here is exotic in isolation. The scale comes from compounding: sequential I/O, the page cache, zero-copy — already on the books from earlier chapters.

Now the multipliers stack on top.

07 · step 2 of 4

Batch: ship trucks, not parcels

Producers linger a few milliseconds (linger.ms) and pack thousands of records into one request. Brokers store the batch as a single unit; consumers fetch in chunks.

Amortize one network round-trip over 10,000 records — that's where the order of magnitude lives.

07 · step 3 of 4

Compress end-to-end

Batches are compressed — lz4, zstd — before leaving the producer. The broker writes and forwards them without ever decompressing; consumers inflate at the edge.

Quarter the bytes on wire and disk, for almost nothing.

07 · step 4 of 4

Then multiply by N

Kreps' benchmark: three cheap machines, default settings, 2,024,032 writes per second — about 675,000 per machine. The curve stays linear because partitions barely share state — they don't coordinate, so they don't contend.

Run the arithmetic: LinkedIn's 81-million-a-second average is about 120 such brokers. Trillions is just multiplication.

The receipt.

what trillions actually cost — itemized
Write path
append-only · sequential · batched · cached by the OS, not the JVM
Read path
sendfile() · zero CPU copies · batched fetches
Durability
acks=all · replication factor 3 · min.insync.replicas=2
Unit of scale
the partition — shard the log, then add machines
Failure model
ISR election in seconds · no committed record lost
State the broker keeps
just the log — your position is yours (offsets)
0

messages that flowed through a LinkedIn-scale stream while you read this page

N.01
Keys decide order. Hash routing is deterministic — pick keys that match how you'll consume.
N.02
Partitions decide parallelism. Write and read throughput are both capped at P. Budget headroom; repartitioning is painful.
N.03
linger.ms trades latency for throughput. A few milliseconds of patience buys 10,000-record batches.
N.04
Replay is a feature — pay for it in disk. Retention is how you debug the past and train on it.
Sources & further reading Kreps, Narkhede, Rao — “Kafka: a Distributed Messaging System for Log Processing”, NetDB 2011
Jay Kreps — “The Log: What every software engineer should know about real-time data's unifying abstraction”, 2013
LinkedIn Engineering — Kafka at LinkedIn scale (2019): 7 trillion messages/day
Jay Kreps — Benchmarking Apache Kafka: 2 Million Writes Per Second (On Three Cheap Machines), LinkedIn Engineering, 2014
Apache Kafka design documentation — kafka.apache.org (linear vs random I/O: ~600 MB/s vs ~100 KB/s)