6.7 Compaction
The swap is what makes compaction crash-safe: the original segment is intact until the replacement is complete.
Log compaction
- 12
- records
- 0
- reclaimed
12 records for 4 distinct keys — 9 of them are superseded. Compaction never touches the active segment, so the most recent writes are always still there in full.
7.1 Why it exists — two use cases
- “you use Kafka to store shipping addresses for your customers. It makes more sense to store the last address for each customer rather than data for just the last week or year. This way, you don't have to worry about old addresses, and you still retain the address for customers who haven't moved in a while.” — delete-retention would lose inactive customers entirely.
- “an application that uses Kafka to store its current state. Every time the state changes, the application writes the new state into Kafka. When recovering from a crash, the application reads those messages from Kafka to recover its latest state. It only cares about the latest state before the crash, not all the changes.”
7.2 The three retention policies
| Policy | Behavior |
|---|---|
delete | "deletes events older than retention time" |
compact | "only stores the most recent value for each key in the topic" |
delete.and.compact | "combines compaction with a retention period. Messages older than the retention period will be removed EVEN IF THEY ARE THE MOST RECENT VALUE FOR A KEY." Purpose: "prevents compacted topics from growing overly large and is also used when the business requires removing records after a certain time period." |
⚠️ "setting the policy to
compactonly makes sense on topics for which applications produce events that contain BOTH a key and a value. If the topic contains NULL KEYS, compaction will FAIL."
7.3 How compaction works
Clean and dirty:
- Clean: “Messages that have been compacted before. Contains only one value for each key, the latest value at the time of the previous compaction.”
- Dirty: “Messages that were written after the last compaction”.
Thread setup: if compaction is enabled at startup ("using the awkwardly named log.cleaner.enabled configuration"), "each broker will start a compaction manager thread and a number of compaction threads."
Partition selection: "Each thread chooses the partition with the highest ratio of dirty messages to total partition size."
The offset map — and its beautiful efficiency:
The cleaner thread reads the dirty section and builds an in-memory map:
| Key part | Value part |
|---|---|
| 16-byte hash of the message key | 8-byte offset of the previous message with this same key |
= 24 bytes per entry.
Worked example from the book: a 1 GB segment at ~1 KB per message → 1,000,000 messages → only a 24 MB map is needed to compact the segment. And “we may need a lot less — if the keys repeat themselves, we will reuse the same hash entries often and use less memory.” ⇒ “This is quite efficient!”
Memory configuration and the failure mode:
"the administrator configures how much memory compaction threads can use for this offset map. Even though each thread has its own map, the configuration is for TOTAL memory across all threads. If you configured 1 GB and you have 5 cleaner threads, each thread will get 200 MB."
⚠️ "Kafka doesn't require the entire dirty section to fit into the size allocated for this map, but AT LEAST ONE FULL SEGMENT HAS TO FIT. If it doesn't, Kafka will log an error, and the administrator will need to either allocate more memory for the offset maps OR USE FEWER CLEANER THREADS. If only a few segments fit, Kafka will start by compacting the oldest segments that fit into the map. The rest will remain dirty and wait for the next compaction."
The counterintuitive fix: if compaction fails for lack of memory, reducing the thread count may help — because total memory is divided among threads. More threads = less memory each.
The copy phase:
A message is omitted because “there is a message with an identical key but newer value later in the partition”.
The swap is what makes compaction crash-safe: the original segment is intact until the replacement is complete.
7.4 Deleted events — tombstones
The question: "If we always keep the latest message for each key, what do we do when we really want to delete ALL messages for a specific key, such as if a user left our service and we are legally obligated to remove all traces of that user?"
During the window, a consumer copying Kafka into a relational database “will see the tombstone message and know to delete the user from the database.”
⚠️ Why the tombstone retention window matters so much
"It is important to give consumers enough time to see the tombstone message, because if our consumer was down for a few hours and MISSED the tombstone message, it will simply NOT SEE THE KEY when consuming and therefore NOT KNOW that it was deleted from Kafka or that it needs to be deleted from the database."
This is a compliance bug disguised as a config default. Tombstone retention must exceed your worst-case consumer downtime.
7.5 deleteRecords — a completely different mechanism
"Kafka's admin client also includes a
deleteRecordsmethod. This method deletes all records before a specified offset, and it uses a COMPLETELY DIFFERENT MECHANISM. When called, Kafka will move the LOW-WATER MARK — its record of the first offset of a partition — to the specified offset. This will prevent consumers from consuming the records below the new low-water mark and effectively makes these records inaccessible until they get deleted by a cleaner thread. This method can be used on topics with a retention policy AND on compacted topics."
| Tombstone | deleteRecords |
|---|---|
| per-key deletion | per-offset (time-range) deletion |
| propagates to consumers as an event they can act on | records become invisible; consumers get no notification |
| works only on compacted topics | works on both policies |
| mechanism: compaction | mechanism: move the low-water mark |
Practical guidance: "delete this user everywhere" → tombstone (downstream systems learn about it). "Delete everything older than 30 days" → deleteRecords (Ch. 5 §7.3).
7.6 When are topics compacted?
“In the same way that the Delete policy never deletes the current active segments, the Compact policy Never compacts the current segment. Messages are eligible for compaction Only on inactive segments.”
The default trigger and its rationale:
"By default, Kafka will start compacting when 50% of the topic contains dirty records. The goal is not to compact too often (since compaction can impact the read/write performance on a topic) but also not to leave too many dirty records around (since they consume disk space). Wasting 50% of the disk space used by a topic on dirty records and then compacting them in one go seems like a reasonable trade-off, and it can be tuned."
Two timing controls:
| Config | Guarantees |
|---|---|
min.compaction.lag.ms | "the MINIMUM length of time that must pass after a message is written before it could be compacted" |
max.compaction.lag.ms | "the MAXIMUM delay between the time a message is written and the time the message becomes eligible for compaction" |
The compliance use case, named explicitly: "
max.compaction.lag.msis often used in situations where there is a business reason to guarantee compaction within a certain period; for example, GDPR requires that certain information will be deleted within 30 days after a request to delete has been made."
A GDPR-compliant compacted topic needs both:
max.compaction.lag.ms→ the tombstone is acted on in time- tombstone retention long enough → consumers see it
Default settings satisfy neither guarantee.