Apache Kafka
Apache Kafka
Definition: Apache Kafka is a distributed event streaming platform that lets many producers publish events and many consumers read them, independently and at very high throughput, durably persisting every record to disk rather than holding it only in memory. It was originally built at LinkedIn to handle the company’s activity-tracking and log-aggregation data at a scale existing message queues couldn’t sustain, then open-sourced and donated to the Apache Software Foundation, where it remains one of the most widely deployed pieces of open-source data infrastructure. Structurally, Kafka behaves less like a traditional Message Queue and more like a distributed, partitioned, replicated commit log — a distinction that shapes almost everything else about how it’s used.
Architecture
- Topics — named, append-only categories of events that producers publish to and consumers read from, roughly analogous to a table or a named channel
- Partitions — each topic is split into one or more partitions, the unit of parallelism; records within a single partition are strictly ordered, but ordering across an entire topic is not guaranteed
- Brokers — the servers that make up a Kafka cluster, each holding a subset of partitions; a production cluster commonly runs anywhere from a handful to hundreds of brokers
- Producers — client applications that publish records, optionally supplying a key that’s hashed to deterministically choose which partition a record lands in
- Consumer groups — one or more consumers cooperating to read a topic, where each partition is assigned to exactly one consumer within the group at a time, spreading out the read load
- Offsets — a monotonically increasing position within a partition; each consumer group tracks its own committed offset, so different groups can read the same topic at completely different points
- Replication — each partition is copied across multiple brokers (commonly a replication factor of 3), with one replica acting as leader and the rest as passive followers for fault tolerance
- ZooKeeper vs KRaft — Kafka historically relied on Apache ZooKeeper to store cluster metadata and elect controllers; newer versions default to KRaft, Kafka’s own Raft-based consensus mode, removing the external ZooKeeper dependency entirely
Why It Matters
- Became the backbone of real-time data infrastructure at most large tech companies, decoupling data producers from consumers at massive scale
- Lets the same stream of events be consumed independently by many different downstream systems — fraud detection, analytics, notifications — without producers needing to know who’s listening or how many
- Durable, replayable storage means a downstream bug or a brand-new consumer can reprocess historical data from any point, something a traditional delete-on-read queue can’t offer
- Its throughput and horizontal scalability made previously batch-only workloads achievable in near real time, commonly feeding stream processors like Apache Spark
- Standardized how services publish and subscribe to domain events at scale, becoming the de facto backbone underneath much of Event-Driven Architecture and Microservices Architecture
Under the Hood: Partition Leader Election and In-Sync Replicas (ISR)
Kafka’s durability guarantees hinge on one specific mechanism, not just “replication” as a vague idea:
- Every partition has one leader replica and zero or more follower replicas; only the leader ever serves reads or writes from clients — followers exist purely for durability
- Followers continuously fetch new records from the leader and append them to their own copy of the log, in the same order the leader wrote them
- The In-Sync Replica (ISR) set is the leader plus every follower that has fully caught up within a configurable lag threshold; a follower that falls too far behind is temporarily dropped from the ISR
- With
acks=all, a producer’s write is only considered committed once every replica currently in the ISR has it, so an acknowledged record survives the loss of any single broker - If the broker hosting a leader crashes, the cluster’s controller promotes one of the remaining in-sync followers to leader — safe specifically because that follower was already fully caught up
- Historically the controller depended on ZooKeeper’s ephemeral nodes to detect broker failure and decide who becomes controller; KRaft replaces this with a Raft-based quorum of controller nodes instead
- Unclean leader election allows a replica outside the ISR to become leader as a last resort, explicitly trading potential data loss for availability during an extended outage — disabled by default in modern versions
Topic “orders” here has three partitions spread across three brokers, and each broker is the leader for one partition while holding follower copies of the other two — parallelism (three partitions doing work simultaneously) and fault tolerance (any single broker can vanish without losing data or availability) come from the same layout, not two separate mechanisms.
Comparison: Kafka vs Traditional Message Queue vs Amazon Kinesis
| Dimension | Kafka | Traditional Message Queue | Amazon Kinesis |
|---|---|---|---|
| Durability/retention | Retained on disk for a configurable period (or indefinitely via log compaction), regardless of whether it’s been read | Message typically deleted once acknowledged/consumed | Retained 24 hours by default, extendable up to 365 days |
| Ordering | Guaranteed within a partition, not across the whole topic | Varies — some (SQS FIFO) guarantee strict order, most don’t by default | Guaranteed within a shard, conceptually the same model as a partition |
| Replay | Native — rewind a consumer group’s offset and reread history | Rare or unsupported — a consumed message is generally gone | Native, within the retention window |
| Throughput | Very high, scales near-linearly by adding partitions and brokers | Moderate, tuned for reliable task delivery over raw volume | High, scales by provisioning additional shards |
| Operational model | Self-managed cluster, or a managed offering like Confluent Cloud/MSK | Often fully managed (SQS) or simple to self-host (RabbitMQ) | Fully managed AWS service, no cluster to operate |
| Typical use case | Event streaming, log aggregation, feeding stream processors | Task distribution, background job processing | Managed streaming ingestion within the AWS ecosystem |
Common Pitfalls
- Treating Kafka like a simple task queue when a lighter traditional message queue would be simpler and sufficient
- Underestimating the operational complexity of running and tuning Kafka clusters reliably
- Choosing a poor (or no) partition key, creating hot partitions where one partition absorbs disproportionate traffic while others sit nearly idle
- Ignoring consumer lag until it’s already an incident, when a slowly climbing lag is usually the earliest available signal that a group can’t keep up
- Assuming ordering holds across an entire topic, when Kafka only guarantees it within a single partition
- Setting retention too short for real reprocessing needs, losing the ability to recover from a downstream bug by replaying history that’s already gone
- Running production topics with a replication factor of 1, turning an ordinary single-broker failure into permanent data loss
Worked Example
Tracing a single event from a checkout service through Kafka to two independent downstream consumers:
| Step | Component | What Happens |
|---|---|---|
| 1 | Producer | Application calls producer.send("orders", key="order-123", value=payload) |
| 2 | Producer | The key is hashed to deterministically pick a partition, e.g. partition 2 of 6 |
| 3 | Broker (leader) | Partition 2’s leader broker appends the record to its log at the next offset, e.g. 5001 |
| 4 | Followers | In-sync replica followers fetch and append the identical record to their own copies of the log |
| 5 | Broker (leader) | Once every ISR member has replicated it (acks=all), the broker sends the producer an acknowledgment |
| 6 | Consumer Group A | Polls partition 2 from its last committed offset, receives record 5001, processes it, commits offset 5002 |
| 7 | Consumer Group B | Independently polls the same partition from its own committed offset, perhaps still catching up around offset 4990 |
| 8 | Retention | Record 5001 remains on disk per the topic’s retention policy, available for replay long after both groups have read it |
Real-World Use
- Log aggregation — centralizing application and infrastructure logs from thousands of servers into one durable, queryable stream instead of scattered local files
- Activity tracking — capturing user clicks, page views, and interactions at web scale, LinkedIn’s original motivating use case
- Event sourcing — using Kafka’s durable, ordered log as the system of record for state changes, with downstream services rebuilding state by replaying events
- Microservice decoupling — services publish domain events like
order.placedorpayment.completedwithout knowing or caring which downstream services react, see Event-Driven Architecture - Real-time analytics pipelines — feeding Apache Spark or Flink streaming jobs that compute rolling aggregates, fraud scores, or recommendations within seconds of an event occurring
Best Practices
- Size partition counts for your target parallelism upfront — increasing them later doesn’t rebalance existing historical data and can break key-based ordering assumptions
- Choose a partition key that spreads load evenly while keeping every record that must stay ordered relative to each other on the same key
- Design consumer groups around independent scaling needs — separate groups for fraud detection, analytics, and notifications so a slow consumer in one group never blocks another
- Monitor consumer lag continuously as the primary early-warning signal that a consumer group is falling behind producer throughput
- Set replication factor to at least 3 in production and pair
acks=allwith a sensiblemin.insync.replicasfor real durability guarantees - Set retention deliberately per topic — long enough to support replay and recovery, short enough to control storage cost, rather than leaving every topic on the cluster default
FAQ
Is Kafka a message queue? Not really — it’s better described as a distributed commit log; unlike a traditional queue, records aren’t deleted on consumption, and multiple consumer groups can independently replay the same data.
Does Kafka guarantee message ordering? Only within a single partition, not across an entire topic — records sharing a key always land on the same partition and stay ordered relative to each other there.
What replaced ZooKeeper in modern Kafka? KRaft (Kafka Raft) mode, which uses a self-managed Raft consensus quorum of controller nodes instead of an external ZooKeeper ensemble for cluster metadata.
Can Kafka lose acknowledged messages?
Only under misconfiguration — with replication factor 3, acks=all, and min.insync.replicas set sensibly, a committed record survives the loss of any single broker.
How is Kafka different from Amazon Kinesis? They’re conceptually similar (partitions versus shards, offsets versus sequence numbers), but Kinesis is fully managed by AWS with simpler operations, while Kafka offers more configuration control and typically higher achievable throughput.
Do consumers remove messages from a topic by reading them? No — reading is non-destructive; a record is only deleted when its retention period or compaction policy expires, never simply because a consumer read it.
History
- Built at LinkedIn starting around 2010 to solve the company’s growing pains around activity-stream and log-aggregation data, led by Jay Kreps, Neha Narkhede, and Jun Rao
- Named after author Franz Kafka — Jay Kreps has said he picked the name because the system was optimized for writing, and he’d taken plenty of literature classes in college
- Open-sourced by LinkedIn in 2011, then donated to the Apache Software Foundation, graduating to a top-level Apache project around 2012
- Its three original creators left LinkedIn in 2014 to found Confluent, a company built around commercializing Kafka with enterprise features and a managed cloud offering
- KRaft mode reached production-ready status in Kafka 3.3 (2022), with ZooKeeper support fully removed a couple of major releases later
- Now one of the most widely adopted pieces of open-source data infrastructure, underpinning real-time systems across finance, retail, ride-sharing, and most large tech companies
Common Interview Questions
- “How does Kafka achieve both high throughput and fault tolerance at once?” — expect partitions (parallelism) and replication plus ISR (fault tolerance) named as two separate, complementary mechanisms
- “What happens if a consumer in a group crashes mid-processing?” — expect an explanation of consumer group rebalancing, reassigning the crashed consumer’s partitions and resuming from the last committed offset
- “Why can’t Kafka guarantee ordering across an entire topic?” — expect recognition that ordering is a per-partition property, in direct tension with spreading records across partitions for parallelism
- “When would you choose Kafka over SQS or RabbitMQ?” — expect reasoning centered on replay, multiple independent consumers of the same data, and sustained high throughput
- “What’s the difference between at-least-once, at-most-once, and exactly-once delivery in Kafka?” — expect each tied to concrete configuration choices:
acks, offset commit timing, idempotent producers, and transactions
Related Terms
- Message Queue
- Event-Driven Architecture
- Batch vs Stream Processing
- Change Data Capture (CDC)
- Microservices Architecture
- Apache Spark
- Data Pipeline
Example
LinkedIn, Kafka’s creator, uses it to stream billions of events per day, from activity tracking to log aggregation across its entire infrastructure. A single page-view topic there can have dozens of independent consumer groups — one materializing a real-time analytics dashboard, another feeding a fraud-detection model, another archiving to a data lake — all reading the exact same partitions at their own pace, with none of them aware the others exist.
Referenced by