A topic is a log: segments and retention
What you'll learn
- to read a table of log segments: where writes go, what is closed, where the log start offset is
- to explain why “six hours” means “at least six”
- to predict what a lagging consumer will see: a silent gap or an error
- to enter the log by time and know why that is almost as cheap as reading from the end
- to choose the segment size and so that the log start offset lands where you need it
17:40. One topic erased everything older than six hours, the other nothing
This morning, after bringing the level back, you switched the twelve antennas to intake at 04:40 and gave the topic signal_raw a six-hour : there is little space on the comms level. You also told it to cut the log often. On disk the log is not one file but segments in a row, and you set a new segment every fifteen minutes or every 900 megabytes, whichever comes first — so that old data would be removed more neatly. You brought the topic up in a hurry, on three partitions.
It is 17:40 now. The reading room asks to reread the feed from 11:00: their reports have a gap between eleven and half past twelve. You point a consumer at the right offset in partition 0 — with the setting “never jump silently” — and get an error instead of data: offset 456,000 no longer exists, the earliest available one is 504,000, which is 11:40. Eleven in the morning was six hours forty minutes ago: those frames crossed the six-hour line forty minutes ago and have already been erased.
You go to the neighboring topic signal_ack — the dispatch room writes its replies to requests there, and the retention is the same, six hours. And you find messages thirteen hours old.
One topic deleted everything older than six hours — and the 11:00 request can no longer be served. The neighboring one deleted nothing. Neither of them lied.
QUERY: Six hours means “not before”. Nobody said “exactly”.

Key numberThe feed rate
48 frames a second — that is the feed's normal rate: twelve antennas, 4 frames a second each.
A frame weighs 450 KB on average, so every second adds up to 21.6 MB, every minute to 2,880 frames, every hour to 172,800, and every day to 4,147,200. These numbers will come up again in the course.
In this lesson we look at partition 0. Every partition has its own message offsets, and the key source_id puts five of the twelve antennas into partition 0 — so it gets 20 frames a second, 9 MB/s. Why the key spreads the antennas unevenly is the next lesson.
TermsThree more words
| word | what it is |
|---|---|
| message and frame | the same thing: in the mechanics we say “message”, while in shift scenes a feed message is called a frame |
| broker | a machine of a Kafka cluster that keeps partitions on its disk; RabbitMQ and its kin are a “queue broker” |
| downstream system | an external system that receives the result: the reading room and its neighbors. A consumer is a program, a downstream system is whoever uses the result |
The English names are in the interview callout at the end of the lesson.
signal_raw and what cleanup did to them. Second: what a consumer asking for the feed from 11:00 will see — under two different settings. Third: the sparse signal_ack under the same retention. Watch the “time” column of the first segment.Why “six hours” means “at least six”
A partition's log is not one file but log segments in a row: each segment is a separate file on disk. Writes always go only to the last one, the active segment. When it grows to its size or time limit, it is closed and a new one is opened.
You cannot delete a single message from the middle of a file. So works on whole segments, by three rules:
- only a whole segment is deleted — never an individual message;
- the decision is made by the newest message in the segment: a segment goes when even that message is older than the retention;
- the active segment is not touched while it holds a fresh message.
This explains both symptoms. In signal_raw a segment closed not after fifteen minutes but after 100 seconds: at 9 MB/s, 900 MB fills up faster. On a dense stream the segment is cut by size. The segments are small, so the cut lands right up against six hours: 252 segments deleted, and the oldest surviving message is exactly six hours old.
signal_ack gets only four replies an hour, and the default segment closes once every seven days — in thirteen hours it has not closed once. It is still active and holds a fresh message, so it is left alone: hence thirteen hours instead of six. And if writes to it had stopped and it had expired entirely, the broker would open a new empty segment and delete the old one — that is the bottom row of the diagram below.
The offset of the oldest message in a partition that can still be read is the log start offset. Each partition has its own, just like the offsets themselves: in partition 0 cleanup moved it forward from 0 to 504,000.
What a lagging consumer sees
The consumer asked for 11:00 — and those messages are gone. What it will see is decided by one consumer setting, not by the broker:
- with
auto_offset_reset='earliest'the consumer silently jumps to the log start offset and gets a gap 11:00–11:40 without a single error: data is lost, and nobody knows; - with
'latest'— which is the default in both kafka-python and the Java client — it jumps to the end and skips everything in between. Silently, too; - with
'none'it gets an error on the very first read: the offset does not exist, and you are told about it instead of being kept in the dark.
How to reread a window from the past, and when this setting kicks in at all, is covered in chapter 3.
How to enter the log by time
The broker keeps each segment's time bounds in memory, so it finds the right segment without opening the log. And next to each segment lie two small index files: one finds the offset by time, the other finds the position in the file by offset. So asking for “from 11:00” is almost as cheap as reading from the end: neither one scans the log, and rereading yesterday is a routine operation, not an emergency.
SEGMENT_MS (and RETENTION_MS, if you like) for signal_ack so that after cleanup at 17:40 the oldest available message is no older than seven hours and no newer than six. Replies are written every fifteen minutes, from 04:40 to 17:25, into a single partition.
The function oldest_age(segment_ms, retention_ms) is already plugged into the editor, and its code is hidden: it creates the topic, writes 52 replies, runs cleanup at 17:40 and returns the age of the oldest surviving reply in milliseconds — or None if none survived. The bottom lines of the starter code call it and print the age — you see their output when you press “Check”. The constants MIN and HOUR (a minute and an hour in milliseconds) are already defined.What in this lesson is real and what is simulated
In the training Kafka a log segment is a counter and a pair of boundaries in the tab's memory, not a file on disk, and its size is computed from a training frame weight of 450 KB. The cell does not store the frames themselves: it counts them and generates only the ones you asked to read. On a real cluster the feed would have written about a terabyte in these thirteen hours.
In the sandbox you trigger cleanup yourself, by calling cluster.apply_retention(...). A real broker does it on its own, in the background and on its own schedule. So in the boundary moves not at the moment you look but when cleanup gets round to it — one more reason why “six hours” means “at least six”.
What is real are the rules, and they carry over to a production cluster without caveats: an offset lives inside a partition, deletion goes by whole segments, the decision is made by the newest message, the active segment is not touched, the log start offset moves, and what a lagging consumer sees — a gap or an error — is decided by its setting.
Interview question
How this comes up in interviews
“We set .ms to 6 hours, but the data stays for 13. Why?” They expect three things: the log segment, the active segment, and the fact that the decision is made by the newest message in a segment. Bonus points for saying that the cleaner runs on a schedule, not at the moment retention expires.
“What happens to a consumer when its group's committed offset points at something already deleted?” The right answer starts with “it depends on auto.offset.reset”: earliest — a silent gap, latest — a silent skip, none — an error. And separately: the client default is latest, and that is worth knowing by heart.
“How does Kafka find a message by time?” Through two indexes next to each segment: the time index gives the offset, and the offset index gives the position in the file. That is why offsetsForTimes does not scan the log. Names that will come up: log, retention.ms, segment.ms, segment.bytes, log start offset.
Key takeaways
| topic | segments | deleted | oldest message age | |
|---|---|---|---|---|
signal_raw, partition 0 | by size, 100 s | 6 h | 252 segments, log start 0 → 504,000 | 6 h 00 min |
signal_ack | default, one active | 6 h | nothing | 13 h 00 min |
Retention works on whole log segments and decides by the newest message in a segment; the active segment is not touched while it holds a fresh message. That is why “six hours” means “at least six”. What a lagging consumer sees — a gap or an error — is decided by its setting. And entering the log by time is cheap: next to each segment lie indexes.
QUERY: You cut the log every fifteen minutes. The feed cut it every hundred seconds — on its own.