skip to content

Explain the __consumer_offsets topic: why it's log-compacted, how offsets are keyed, and how the group coordinator uses it.

level: seniorimportance: should knowfreq 45%

answer

  1. internal topic, 50 partitions default, RF 3
  2. key = (group, topic, partition)
  3. cleanup.policy=compact, keep latest per key
  4. coordinator = leader of hash(groupId) % numPartitions
  5. tombstone + offsets.retention.minutes expire offsets

basics

~20 s

__consumer_offsets is an internal Kafka topic (50 partitions by default) where consumer commits are stored as messages keyed by (group, topic, partition). It's log-compacted so only the latest committed offset per key is retained. The group coordinator broker reads/writes it to track group progress.

solid answer

~40 s

Kafka stores committed offsets as records in the internal topic __consumer_offsets, created automatically with offsets.topic.num.partitions (default 50) partitions and replication offsets.topic.replication.factor (default 3). Each commit is a message keyed by (groupId, topic, partition) with a value containing the offset, metadata, and timestamp. The topic uses cleanup.policy=compact (log compaction): compaction keeps only the most recent value per key, so the latest commit for each (group,topic,partition) survives while old commits are garbage-collected — bounding the topic's size despite continuous commits. A group's coordinator is the broker that leads the __consumer_offsets partition determined by hash(groupId) % numPartitions; that broker handles OffsetCommit/OffsetFetch requests and group membership. On startup the coordinator replays its partition to rebuild the in-memory offset map. Offset retention (offsets.retention.minutes) and tombstones handle expiry.

go deeper

for a junior

Know offsets live in an internal Kafka topic, not an external DB.

for a middle

Know it's compacted, keyed by group/topic/partition, with 50 partitions by default.

for a senior

Explain coordinator selection via hash(groupId), compaction's role, and recovery by replaying the partition.

for a principal

Reason about RF and partition-count sizing for many groups, offset retention/tombstone expiry, and coordinator failover impact on availability.

## What it is `__consumer_offsets` is an **internal, auto-created** Kafka topic that durably persists consumer group state — primarily committed offsets, plus group metadata. It's a normal Kafka topic used as a key-value store via the log. ## Why a topic (not a separate database) Kafka reuses its own replicated, durable log instead of an external store. Commits become **append-only messages**; durability and replication come for free from Kafka's partition replication. ## Partitions and replication - Default **50 partitions** (`offsets.topic.num.partitions`). - Default replication factor **3** (`offsets.topic.replication.factor`), so group state survives broker loss. A consumer group maps to exactly one partition of this topic via `hash(groupId) % offsets.topic.num.partitions`. The **leader** of that partition is the group's **coordinator** — the broker responsible for that group's offset commits/fetches and rebalance orchestration. ## Keying Each committed offset is a message whose **key is (groupId, topic, partition)** and whose **value** encodes the committed offset (next-to-read), optional metadata, leader epoch, and commit timestamp. Because the key uniquely identifies the slot being updated, repeated commits to the same slot share a key. ## Why log compaction The topic is configured `cleanup.policy=compact`. **Log compaction** retains, for each distinct key, **only the latest message**, discarding older ones in the background. Without compaction, a busy group committing every few seconds would write millions of obsolete offset records forever. With compaction, the topic's size is bounded to roughly one record per active (group,topic,partition). A **tombstone** (null value) for a key deletes that key entirely — used to expire offsets for groups/partitions no longer in use, governed by `offsets.retention.minutes`. ## The coordinator's use of it 1. **OffsetCommit**: a consumer's commit becomes an append to the coordinator's __consumer_offsets partition; once replicated, the commit is durable. 2. **OffsetFetch**: on (re)join, the consumer asks the coordinator for its committed offsets; the coordinator serves them from its in-memory map. 3. **Recovery**: if a coordinator broker restarts or leadership moves, the new leader **replays** the __consumer_offsets partition from the log to rebuild the in-memory offset/group map before serving requests. ## Practical notes - You can inspect it with `kafka-console-consumer --topic __consumer_offsets --formatter kafka.coordinator.group.GroupMetadataManager\$OffsetsMessageFormatter` (or the consumer-group tooling, which is friendlier). - `kafka-consumer-groups.sh --describe --group X` reads from here to show current offset, log-end offset, and lag. - Don't manually produce to it; it's managed by the coordinator.

  • How is a consumer group's coordinator broker determined?
    By hashing the groupId modulo the number of __consumer_offsets partitions to pick a partition; the broker that leads that partition is the group coordinator and owns that group's commits and rebalances.
  • What would happen if __consumer_offsets were NOT compacted?
    It would grow without bound as every commit appends a new record, eventually consuming huge disk and slowing coordinator recovery (which replays the partition). Compaction caps it to roughly one latest record per (group,topic,partition).

saying these in an interview costs you the question

  • Saying offsets are stored in ZooKeeper (legacy pre-0.9 only)
  • Calling __consumer_offsets a deletion (time-based) topic — it's compacted, not delete-policy
  • Thinking each consumer instance has its own coordinator (the coordinator is per-group)
  • Manually producing to or editing __consumer_offsets

context