Learn Labs
6. Kafka Internals

6.7 Compaction

The swap is what makes compaction crash-safe: the original segment is intact until the replacement is complete.

Log compaction

state
12
records
0
reclaimed
user:1=v1user:2=v1user:3=v1user:4=v1user:1=v2user:2=v2user:3=v2user:4=v2user:1=v3user:2=v3user:3=v3user:4=null
Problem

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.

With cleanup.policy=compact the log keeps at least the most recent value for every key, so the topic becomes a replayable snapshot of current state rather than a history of changes. A null value is a tombstone and deletes the key.

7.1 Why it exists — two use cases

  1. “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.
  2. “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

PolicyBehavior
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 compact only 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:

CLEANcompacted beforeDIRTYwritten after the last compactionone partition, oldest offsets on the left
  • 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”.
Figure 6.7.2Clean and dirty

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 partValue part
16-byte hash of the message key8-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:

each message in the clean segmentsoldest firstkey NOT in the offset mapthis value is STILL THE LATESTkey IS in the offset mapa newer value exists laterCOPY to a replacementsegmentOMIT the messagereplacement SWAPPED for the originalthen on to the next segmentEnd state: ONE MESSAGE PER KEY — the one with the latest value.

A message is omitted because “there is a message with an identical key but newer value later in the partition”.

Figure 6.7.4The copy phase

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?"

after that timeproduce that key with a NULL VALUE= a TOMBSTONEcleaner thread finds itnormal compaction, retaining ONLY the null-value messagekept around for a configurable amount of timeconsumers CAN SEE it and know the value is deletedcleaner removes the tombstonethe key is GONE

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.”

Figure 6.7.57.4 Deleted events — tombstones

⚠️ 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."

consumer offline while the tombstoneis created AND removedon restart, the key simply DOESN’T EXISTin the topicno event telling the consumer to deleteTHE DOWNSTREAM DATABASE STILL HAS THE USERGDPR deletion silently failedin the system that matters
Figure 6.7.6⚠️ Why the tombstone retention window matters so much

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 deleteRecords method. 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."

TombstonedeleteRecords
per-key deletionper-offset (time-range) deletion
propagates to consumers as an event they can act onrecords become invisible; consumers get no notification
works only on compacted topicsworks on both policies
mechanism: compactionmechanism: 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:

ConfigGuarantees
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.ms is 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.


On this page