Offsets, retention, and where bookmarks live
A consumer's offset is just another piece of data stored in the __consumer_offsets topic; reading and committing are separate ack channels, and retention — not the consumer — decides when records age out.
Reads don't delete and consumers carry their own bookmark. So who DOES delete records, and where does that bookmark physically live? Both answers are themselves Kafka topics.
Scene 02a
Offsets, retention, and where bookmarks live
- Watch
- Try it
- Predict
- Capture
Watch the read and committed offsets — they start in lockstep. They don't have to. Retention crawls in from the left; messages stay until retention ages them out.
Highlighted lines are the ones running in the diagram right now.
def poll():resp = broker.fetch(topic, partition,fromOffset = self.read,)for record in resp.records:self.read += 1 # in-memory onlyyield recordif enable.auto.commit:# every auto.commit.interval.mscommitOffset(self.read)
def commitOffset(offset):# __consumer_offsets is a normal Kafka topic;# the key is (group, topic, partition).record = OffsetCommit(group = self.group,topic = topic,part = partition,offset = offset,)coordinator.append(record)self.committed = offset
# runs continuously per partitiondef maybeRoll():for seg in segments:tooOld = age(seg) > log.retention.mstooLarge = totalBytes() > log.retention.bytesif tooOld or tooLarge:delete(seg) # baseOffset goneif active.size >= log.segment.bytes:roll() # start a new segment
Where this sits in Build Kafka
Scene 02a of 13, in the Why a log? act — Orientation — the log is the database, not a queue.. Read and commit are separate ack channels; retention, not consumers, ages records out.
Up next. Partitions — splitting the log. One partition is one ordering lane; multiple partitions trade strict global order for parallelism, and the key picks the lane.
Designs that use this
- Twitter / X TimelinePush or pull? Both. The canonical fanout problem.
- Uber / Lyft — Match Drivers and RidersMatch a rider to the closest acceptable driver in under 3 s. Geohash, S2, surge.
- Slack / DiscordChannels and history. Push or pull — and how a hot-channel fanout doesn't melt the gateway.
- WhatsApp / MessengerHundreds of millions of long-lived sockets, sub-second 1:1 + group delivery, E2E-encrypted, multi-device, multi-region active-active.