Hello! Welcome to your first lesson in the course on designing high-load distributed systems.
Given your extensive experience in building high-performance payment and trading systems, you're already familiar with the critical importance of durability and reliability. This course aims to deepen that practical knowledge by examining the architectural choices and configuration trade-offs of specific technologies that underpin these systems.
We'll begin with Apache Kafka, a cornerstone of modern data infrastructure. Today's lesson focuses on its core design.
Lesson Goal: By the end of this session, you will be able to explain Kafka's log-based architecture and the mechanisms that provide its durability guarantees. This foundation is essential for making informed decisions when configuring and operating Kafka in production environments.
1. The Foundation: The Write-Ahead Log (WAL)
Before we dive into Kafka specifically, let's start with a fundamental pattern that is at the heart of many reliable systems, including PostgreSQL, which we'll cover later, and Kafka itself: the Write-Ahead Log (WAL).
The core principle is simple but powerful: log first, apply later. Any change is first written to a durable, append-only log before the actual data is modified in place. This ensures that even if the system crashes mid-operation, it can recover to a consistent state by replaying the log.
To understand this concept in more detail, please read the following sections from the article "The Write-Ahead Log: A Foundation for Reliability".
The Write-Ahead Log: A Foundation for Reliability in ...
This article provides an excellent overview of the WAL pattern, explaining its role in ensuring durability and consistency. It uses a database as an initial example, which provides a familiar context before we apply the concept to Kafka.
Please read the sections 'What is a Write-Ahead Log (WAL)?' and 'Why Write-Ahead Logs Are Everywhere'. Focus on how the 'log first, apply later' sequence provides both durability and a clear path for crash recovery.
As you can see, the WAL transforms complex, risky random-access writes into simple, fast, and reliable sequential appends. This is the key to providing strong guarantees without sacrificing performance.
2. Kafka's Architecture: The Log as the System
Now, let's see how Kafka takes this principle to its logical conclusion. For many systems, the WAL is an internal implementation detail for ensuring atomicity and durability. For Kafka, the log is the system. The primary abstraction that Kafka exposes to its clients is a distributed, replicated log.
Let's break down the core components of this architecture.
Topics, Partitions, and Offsets
- Topic: A topic is a logical name for a stream of records. For example, you might have a
paymentstopic or anorderstopic. - Partition: A topic is split into one or more partitions. Each partition is an independent, ordered, immutable sequence of records—a physical log file. By distributing partitions across different servers (called brokers), Kafka achieves horizontal scalability.
- Offset: Each record within a partition is assigned a unique, sequential ID number called an offset. This offset immutably defines the position of a record within that partition.
The diagram below illustrates this fundamental structure. Producers append records to the end of the partitions, and consumers read records sequentially, tracking their position using offsets.
This model provides a crucial guarantee: ordering is preserved within a partition, but not across partitions of the same topic.
Physical Log Structure: Segments
To manage these potentially massive log files efficiently, Kafka further divides each partition log into segments. A partition is a logical sequence, but on disk, it's a collection of segment files.
Please read the following sections to understand this physical structure.
Kafka Logs: Concept & How It Works & Format
The article 'Kafka Logs: Concept & How It Works & Format' provides a clear, detailed look at the physical storage layer of Kafka. This is key to understanding performance and retention.
Read the sections 'Understanding Kafka Logs', 'Kafka Log Structure and Components', and 'How Kafka Logs Work'. Focus on the roles of log segments, the .log, .index, and .timeindex files, and the concept of the 'active segment'.
To summarize, when a producer sends a message:
- It is appended to the active segment of the leader broker for that partition.
- The message is written to the
.logfile. - Its offset and position are recorded in the
.indexfile to enable fast lookups.
This segmentation is what allows Kafka to efficiently manage data retention. Instead of deleting individual messages, Kafka can simply delete entire segment files once they are older than the configured retention period.
3. Kafka's Durability Guarantees
Now we can connect this architecture to the durability guarantees it provides. Durability in Kafka is not a single feature but a result of several mechanisms working together, from disk flushing to distributed replication.
Durability on a Single Broker
At the most basic level, durability means the data is persisted to disk and will survive a process restart. Kafka leverages the operating system's page cache for high performance. Writes are made to the in-memory page cache first, which is fast, and the OS handles flushing this data to the physical disk in the background.
While you can force a flush after every message, this severely impacts performance. In practice, durability is achieved through replication.
Durability Through Replication
To protect against broker failure, Kafka replicates partitions across multiple brokers in the cluster. For each partition, one broker is elected as the leader, and the others act as followers.
- All producer writes and consumer reads for a partition go through the leader.
- Followers fetch data from the leader, attempting to keep their own copy of the log fully in sync.
A write is only considered committed and durable once it has been successfully replicated to a set of followers. This set is called the In-Sync Replicas (ISR). The ISR is a dynamic set of the leader and all followers that are caught up with the leader's log.

The producer's acks configuration controls the durability guarantee for each write:
acks=0: The producer does not wait for any acknowledgment. This offers the lowest latency but no durability guarantee.acks=1: The producer waits for an acknowledgment from the partition leader only. The write is durable if the leader remains alive but can be lost if the leader crashes before followers replicate it.acks=all(or-1): The producer waits for an acknowledgment from the leader after the write has been replicated to all brokers in the ISR. This provides the strongest durability guarantee. If a leader fails, a new leader can be elected from the up-to-date followers in the ISR without data loss.
Durability Through Retention
Finally, durability in Kafka also refers to how long the data is stored. Unlike traditional message queues that delete messages after they are consumed, Kafka retains messages based on configurable policies. This allows multiple consumers to read the same data and enables "rewinding" to re-process historical data.
There are two main cleanup policies:
- Delete (default): Old log segments are deleted based on age (
log.retention.ms) or total log size (log.retention.bytes). - Compact: Kafka retains only the most recent record for each unique key in a partition. This is useful for maintaining the latest state of an object, similar to a database table, and ensures that state is durable indefinitely (or until a new value for that key arrives).
We will explore the configuration of these policies in a later lesson. For now, the key is to understand that Kafka's log is a durable storage system, not just a transient message bus.
Key Takeaways
- Kafka's architecture is a direct and powerful implementation of the Write-Ahead Log (WAL) pattern, where the log itself is the core product.
- Topics are logical streams, but the unit of parallelism and storage is the partition, which is an ordered, append-only log.
- On disk, partitions are broken into segments, which allows for efficient storage management and data retention.
- Durability is achieved through two main mechanisms:
- Replication: The leader-follower model and the In-Sync Replica (ISR) list ensure that writes are copied to multiple brokers before being acknowledged as committed.
- Retention: Configurable policies (
deleteorcompact) determine how long data persists in the log, making Kafka a durable storage system.
Next Up
In our next lesson, we will move from theory to practice and learn how to configure Kafka topics with partitions and a replication factor. This will allow you to apply the architectural concepts we've discussed today to set up a topic with the desired level of parallelism and fault tolerance.