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.

Previously

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

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Producersend(record)appends onlyPARTITION 0 (OF 3 — ZOOMED IN)events3 partitionsshowing P0 — pedagogy zoomREAD VS COMMITrecord = consumer.poll()// read advancesprocess(record)consumer.commit()// writes to __consumer_offsets// on restart: resume from committed__consumer_offsetsinternal Kafka topicwhere the bookmark physically livesConsumeroffset = 0Consumerread=0, committed=0
What to watch for

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.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Consumer.poll
advances the in-memory cursor; auto-commit piggybacks here
def poll():
resp = broker.fetch(
topic, partition,
fromOffset = self.read,
)
for record in resp.records:
self.read += 1 # in-memory only
yield record
if enable.auto.commit:
# every auto.commit.interval.ms
commitOffset(self.read)
Consumer.commitOffset
persists the bookmark — a write into __consumer_offsets
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
LogCleaner.maybeRoll
broker-side retention — ages segments out by time or size
# runs continuously per partition
def maybeRoll():
for seg in segments:
tooOld = age(seg) > log.retention.ms
tooLarge = totalBytes() > log.retention.bytes
if tooOld or tooLarge:
delete(seg) # baseOffset gone
if 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.

All 13 scenes in Build Kafka · Every curriculum

Built with Arqly
Every scene in Build Kafka builds on the one before it.All 13 Build Kafka scenes