Compaction: the latest value per key
O que você vai aprender
- to tell apart two questions you can ask a topic: “what happened” and “how do things stand now”
- to apply the compaction rule: the latest message survives for each key
- to read a log with gaps in its offsets without tripping over them
- to delete a key with a message that has no value — and to remember that such a message has a lifetime
- to build a directory from a compacted topic in a single consumer pass
19:10. A one-line directory
2 November 2184, the dispatch room of Vault-9. The reading room needs an antenna directory: twelve rows, s01…s12, each with a sky sector and a status. The level already has such a directory — the topic sources. The previous shift kept it, and this morning you brought it back along with the level. Over the years it has piled up 2.1 million messages for these twelve keys: antennas get retuned, statuses flicker, and every change goes out as a message keyed by source_id. The topic's is “forever”.
A new consumer that needs the directory reads the topic from the beginning: over two million messages for the sake of twelve rows, three minutes at start-up alone. And you do exactly what you learned at 17:40: give the topic a thirty-day retention and run retention cleanup.
A few minutes later the reading room sees one antenna in the directory — s07. The other eleven have not gone anywhere: they simply have not changed for half a year, and the latest message about each of them turned out to be older than thirty days. Retention deleted them in whole log segments — exactly as it should.
You did not delete a single message by mistake. You did everything by the rules — and lost eleven rows out of twelve.
QUERY: Retention can ask only one thing: how old a message is. A directory cares about something else.

sources: 2,002 messages for the same twelve keys instead of 2.1 million — otherwise the table would not fit on the screen, and the result does not depend on scale. A log segment is 500 messages. One and the same stream flows into two topics, and the only difference is the policy: sources gets a 30-day retention, while sources_demo is created with a different policy — you can see it in the code. The third printout deletes one key from sources_demo. What the second policy does, we take apart right after the cell.Two questions for a topic
The first printout is the incident from the scene in miniature: two segments deleted, 12 keys before, 1 after. Retention worked correctly. The mistake lies elsewhere: the topic was asked the wrong question.
There are two different questions you can ask a topic:
- “what happened?” — that is history. Every message is valuable in itself, and old ones can be thrown away once nobody needs them. The feed in
signal_rawis like that, and it needs ; - “how do things stand now?” — that is state. Of all the messages about an antenna only one matters, the latest, and its age means nothing: an antenna that has not changed for half a year has not disappeared. The
sourcesdirectory is like that.
A directory does not need the whole history — it needs the latest value of each key. That sentence is, almost word for word, the cure. Log compaction: of all messages with the same key, the latest is kept and the earlier ones are thrown away, and the log turns into a directory.
Compaction is switched on for a topic with a single setting, and such a topic is called a compacted topic. (ClickHouse's ReplacingMergeTree tables keep the same rule — the latest version per key.)
You are not the one who compacts. A background broker process walks the closed log segments — the log cleaner. For each key it keeps one message — the latest across all the closed segments, not the latest in each one. It does not go into the active segment.
What the second printout showed
- Scale. There were 2,002 messages, 14 are left, and still 12 keys: the latest values of eleven antennas in the closed segments, and for s07 its latest value among the closed segments (offset 1999) plus two more messages in the active segment. The two million messages of the real topic also exist for the sake of twelve rows.
- Offsets are not reassigned. Survivors stay at their old offsets — 988, 989, 991… and straight to 1999, and everything in between is gone for good. These are offset gaps. The log start offset is still 0, but the first message a consumer gets is offset 988. Code that assumes offset n is always followed by n + 1 will break on the very first compacted topic.
poll()steps over the gaps by itself: the consumer's position passes through them and only grows. - The active segment is untouched. Offsets 1999, 2000 and 2001 are three messages of key s07, and all of them are in place: the cleaner does not look into the active segment and does not know that offset 1999 is already outdated. So “one message per key in a compacted topic” is wrong — “eventually” is right.
How to delete a key if the latest always survives
The compaction rule keeps a key's latest message for as long as you like. How, then, do you say “this antenna no longer exists”? Send a tombstone — a message with a key and no value: p.send('sources_demo', None, key='s11'). For the cleaner it is the key's latest message, so it throws away all the earlier values of s11. And a consumer building the directory sees the empty value and removes the key on its side.
The third printout was made on the copy — nobody decommissioned the real s11. It shows two steps:
- at 19:10 the tombstone landed in the active segment, and the cleaner does not care about it yet. At 20:10 s07 sent a new value, and the segment holding the tombstone closed. Only then did compaction see it and remove three messages — the previous value of s11 and two old values of s07;
- s11 is no longer in the directory, 11 keys. But the tombstone itself stays in the log, and it has a lifetime: 24 hours.
That lifetime is the tombstone window: how long a tombstone stays in the log so that lagging consumers get a chance to see it. The clock starts at the compaction that first found it in a closed segment. When the window closes, the next compaction throws it out too — and the log will not have a single word left about the key.
Hence the price. A consumer that already has s11 in its assembled directory and has been idle for longer than the window will never learn about the deletion: it comes back, the tombstone is gone, and s11 stays in its directory forever.
build_directory(consumer): the consumer walks the compacted topic sources_demo from the very beginning to the end, and the function returns the directory — a dictionary antenna → value. The topic is the same as after the cell's third printout: the offsets have gaps, and key s11 has been deleted with a tombstone.
The deleted key must not be in the directory — with any value, not even None. The consumer comes from make_consumer(): it is in a new group and starts from the log start offset — the default in the training Kafka; a real consumer needs auto_offset_reset='earliest' for that. The check builds the directory twice, with two new consumers, and expects the same answer both times.sources_demo: p.send('sources_demo', {'sector': 2, 'status': 'в работе'}). What happens?What in this lesson is real and what is simulated
Compaction in the sandbox is an explicit call to cluster.compact(...). The real cleaner comes by itself and whenever it sees fit: it takes on a log only once its uncompacted part crosses a threshold, and it may stay away for hours. Hence the “eventually” caveat: at any given moment the topic may hold an uncompacted tail.
The cell's training numbers. The threshold for the uncompacted part is lowered — 0.1 instead of the default 0.5, otherwise the second compaction would not happen on the small copy. A log segment is exactly 500 messages because they all weigh the same. But the 24-hour tombstone window is not a training number: it is real Kafka's default.
The sandbox closes a time-based segment by the cluster clock, while a real broker goes by message timestamps. That is why in the cell the 10:00 and 18:40 messages sit in one active segment, while on a cluster the history with old timestamps would land differently; the compaction rule does not change because of that.
The cell does not print the rule from the question about a message without a key, but the sandbox enforces it the same way the broker does: answer the question, then check your answer yourself with one line.
The real rules in the lesson: the latest message per key survives, the active segment is not compacted, offsets are not reassigned, a tombstone lives for a limited time — and the rule from the question about a message without a key.
Interview question
How this comes up in interviews
The first question almost always sounds the same: “How does log compaction differ from time-based , and when do you use which?” Answer not with settings but with a question about the data. If the topic answers “what happened”, it is an event log and needs retention (retention.ms). If it answers “how do things stand now”, it is state and needs compaction (cleanup.policy=compact).
There is a trap built into the question: many people answer “it removes duplicates”. No: it keeps the latest message per key. Two identical messages with different keys are not duplicates, and they are not going anywhere.
The second question is “how do you delete a key from a compacted topic?” They expect not just the word tombstone — a message with a key and an empty value — but also the caveat: it lives for a limited time (delete.retention.ms, one day by default), and a consumer that lags behind for longer than that will miss the deletion.
The third, and the most mature: “can a compacted topic replace a lookup table in a database?” Yes — and a good candidate names the price: you have lost the history of changes. If you need it, you need two topics. More names that will come up: log cleaner, min.cleanable.dirty.ratio, active segment and non-contiguous offsets — offsets with gaps.
Principais pontos
| topic | policy | messages | keys |
|---|---|---|---|
sources | 30-day | 2,002 → 1,002 | 12 → 1 |
sources_demo | compaction | 2,002 → 14 | 12 → 12 |
sources_demo, s11 deleted | compaction | 13, including the tombstone | 11 |
Retention answers the question “what happened”, log compaction answers “how do things stand now”: the latest message survives for each key. The cleaner only walks closed segments and does not reassign offsets — hence offset gaps. A key is deleted with a tombstone that has no value, and it lives for a limited time — the tombstone retention window: a consumer that is late by longer than that will miss the deletion. A directory is built from a compacted topic in a single consumer pass from the beginning.
You switch the real sources to compaction and load the twelve rows again: no policy can bring back what retention has erased.
QUERY: The log remembers everything. A directory only needs each antenna's last word.