Key Technology: Kafka and Distributed Logs
Kafka is a distributed, replicated, append-only log. It is the standard answer when a design needs to move a large volume of events between systems durably, let several consumers read the same events independently, and replay history. This chapter explains how it works well enough to reason about ordering, throughput, delivery guarantees and failure, which is what interviewers actually ask.
Product features and defaults change between versions, so treat specific settings as things to confirm against the documentation of the version you would run. The concepts here are stable.
1. The core model
A topic is a named stream of events. It is divided into partitions, and each partition is an ordered, append-only log stored on disk. Every event in a partition has an increasing number called its offset.
- Producers append events to a topic. The partition is chosen from the event's key by hashing it, so all events with the same key go to the same partition. With no key, events are spread across partitions.
- Brokers are the servers that store partitions.
- Consumers read events in order from a partition, starting at an offset they choose. Reading does not delete the event. Events stay until a retention limit by time or size, or forever for a compacted topic.
The consequence is the model's main strength: the log is a durable, shared record, and each consumer group keeps its own position in it.
<!--fig:log-->2. Why Kafka is fast
Kafka reaches very high throughput, often hundreds of thousands to millions of messages per second across a cluster, using a few simple ideas:
- Sequential disk access. Appending to the end of a log and reading in order is much faster than random access, and the operating system's page cache serves recent data from memory.
- Batching. Producers and consumers move many messages per request, and brokers store them as batches, which amortises overhead and enables compression.
- Zero-copy transfer. The broker can send data from the file cache to the network without copying it through application memory.
- Partitioning. Work is spread across partitions and brokers, so the cluster scales out.
Say that throughput depends on message size, batch settings, replication and hardware, and that you would benchmark for your workload.
3. Partitions, ordering and parallelism
This is the part of Kafka that interviews return to most.
Ordering is only guaranteed within a partition. There is no global order across partitions. If events for one order must be processed in sequence, give them the same key, such as the order identifier, so they land in the same partition.
Parallelism is limited by partitions. Within a consumer group, each partition is read by exactly one consumer. With 12 partitions you can usefully run up to 12 consumers in a group, and a thirteenth sits idle. Choose the partition count with future growth in mind, because increasing it later changes the key-to-partition mapping and can break ordering for existing keys.
Key choice matters. A skewed key, such as a customer that generates half of all events, creates a hot partition. Choose a key with an even spread, or add a salt for hot keys at the cost of per-key ordering across the salted pieces.
Consumer groups. Consumers sharing a group id divide the partitions among themselves. Different groups each receive every event. So the same topic can feed billing, search indexing and analytics, each with its own group and its own pace. When a consumer joins or leaves, the group rebalances and partitions are reassigned, causing a brief pause.
4. Replication and durability
Each partition has a replication factor, typically 3. One replica is the leader, which handles all reads and writes for that partition. The others are followers, which copy the leader's log. The set of followers that are caught up forms the in-sync replicas.
<!--fig:cluster-->Producer acknowledgements set how much safety you buy per write:
| Setting | Meaning | Trade-off |
|---|---|---|
| Fire and forget | Do not wait | Fastest, can lose messages |
| Leader only | Wait for the leader to write | Lost if the leader dies before replication |
| All in-sync replicas | Wait for every in-sync replica | Safest, higher latency |
Combine "all in-sync replicas" with a minimum in-sync replica count, for example 2 of 3. If fewer replicas are available, the broker rejects writes instead of accepting data that is not safely replicated. This is the standard "no acknowledged write is lost" configuration, and the trade-off is availability: when too many replicas are down, writes fail.
If a leader fails, one of the in-sync followers is elected the new leader, so no acknowledged data is lost. Allowing an out-of-sync replica to become leader restores availability sooner and can lose data, and it is normally disabled for important topics.
5. Delivery guarantees
Say this precisely, because it is a classic question.
- At-most-once: the consumer commits its offset before processing. A crash after the commit and before the work loses the message.
- At-least-once: the consumer processes, then commits the offset. A crash between the two causes the message to be processed again. This is the common default.
- Exactly-once effects: Kafka offers idempotent producers, which prevent duplicates caused by producer retries, and transactions, which let a consumer-process-produce loop commit its output and its input offset atomically, within Kafka. For effects outside Kafka, such as writing to a database, you still need idempotent consumers: store a processed identifier or make the write naturally idempotent.
The accurate summary: the system gives at-least-once delivery, and you build exactly-once effects with idempotent writes, using Kafka's transactions where the whole flow stays inside Kafka.
6. Retention, replay and compaction
Retention. Events are kept for a configured period or size, for example seven days. Within that window any consumer can replay: start a new group at the beginning to build a new index, or rewind an existing group after a bug fix.
Log compaction. A compacted topic keeps at least the latest value for each key and removes older ones. It turns a topic into a changelog of the current state, which is useful for rebuilding a cache or a table by reading the topic from the start.
Storage cost. Volume equals event rate times size times retention times replication. At 100,000 events a second of 1 KB with seven days of retention and three replicas, that is about bytes, close to 180 TB. Doing this arithmetic out loud is a strong signal.
7. Where Kafka fits in designs
Decoupling and buffering. Put a topic between a bursty producer and a slower consumer, such as clickstream events into an analytics pipeline. The log absorbs the spike.
Fan-out. One event, many consumer groups: an order placed event feeds billing, inventory, email and analytics, each independent.
Change data capture. Stream a database's change log into a topic, so search indexes, caches and warehouses stay in sync without dual writes.
Event sourcing and replay. Keep events as the record of what happened and rebuild state by replay.
Stream processing. Consume, transform, aggregate over windows and write results to another topic.
8. Where Kafka is the wrong choice
- A simple task queue with per-message acknowledgement, delays and priorities. A queue system built for that is simpler. Kafka's unit of progress is an offset per partition, so one slow message blocks the messages behind it in that partition, and there is no built-in delay or priority.
- Low volume, where operating a cluster is overkill. A managed queue or a lightweight option may be enough.
- Request-response. Kafka is not designed for low-latency synchronous calls.
- Random access to individual events. It is a log, not a database with indexes.
9. Failure modes and operations
- Consumer lag. The gap between the newest offset and a group's offset. Rising lag means consumers cannot keep up. Monitor it per group, scale consumers up to the partition count, and optimise the handler.
- Poison message. A message that always fails blocks its partition. Retry a limited number of times, then send it to a dead-letter topic and continue.
- Rebalance storms. Frequent consumer restarts cause repeated rebalances and stall processing. Use stable group membership and avoid long processing that exceeds the heartbeat timeout.
- Hot partition. Fix the key.
- Broker failure. A new leader is elected from the in-sync replicas. Re-replication of the lost copies uses network and disk, so throttle it.
- Unclean leader election. Know that enabling it trades data safety for availability.
- Schema evolution. Producers and consumers change independently, so use a schema with compatibility rules, so that old consumers can read new events and the reverse.
- Cluster metadata. Recent versions manage cluster metadata with a built-in consensus protocol, and older ones used a separate coordination service. Mention which you assume, and that you would confirm for the version in use.
10. Interview questions and model answers
Q: How do you guarantee ordering in Kafka? Only within a partition. I choose a key such as the order identifier so that all events for one order go to the same partition, and a single consumer in the group reads that partition in order.
Q: How many consumers can read a topic in parallel? Up to the number of partitions in a consumer group. Extra consumers sit idle. I would pick the partition count from the expected peak throughput and leave room to grow.
Q: How do you avoid losing messages? Replication factor 3, producer acknowledgements from all in-sync replicas, a minimum in-sync replica count of 2, and consumers that commit offsets after processing.
Q: How do you get exactly-once processing? At-least-once delivery plus idempotent processing. Kafka's idempotent producer and transactions help when the flow stays inside Kafka. For an external database, I store the processed message identifier with the result in one transaction.
Q: A consumer group falls far behind. What do you do? Check lag and the handler's speed, scale consumers up to the partition count, increase partitions if needed, fix slow downstream calls, and consider shedding or deferring low-priority work.
Q: Kafka or a message queue? Kafka when I need high throughput, durable replay, multiple independent consumers and ordering by key. A queue when I need per-message acknowledgement, delays, priorities and simple competing consumers at modest scale.
11. Common mistakes
- Assuming global ordering across partitions.
- Choosing a skewed key and creating a hot partition.
- More consumers than partitions, expecting more throughput.
- No dead-letter strategy, so a poison message blocks a partition.
- Treating exactly-once delivery as free, with non-idempotent consumers.
- Forgetting retention and replication in the storage estimate.
- Using Kafka as a task queue with per-message delays.