Why does the replication factor and configuration of the __consumer_offsets topic matter for the Group Coordinator, and what can go wrong if it's misconfigured?
answer
- coordinator state = __consumer_offsets log
- offsets.topic.replication.factor default 3
- RF=1 + broker loss -> offsets gone -> auto.offset.reset
- auto-created at available brokers if cluster too small
- compacted; num.partitions=50 fixes coordinator spread
basics
~10 sThe coordinator's state and committed offsets live in __consumer_offsets. If its replication factor is too low (e.g. 1), losing that broker loses group/offset state and breaks coordinator failover, causing duplicate processing or stuck groups.
solid answer
~50 sBecause the coordinator is the leader of a `__consumer_offsets` partition and that partition's log *is* the group's durable state (committed offsets + group metadata), the topic's durability settings directly determine coordinator reliability. The key config is **`offsets.topic.replication.factor`** (default 3, but auto-created at cluster's available brokers in tiny dev clusters). With replication factor 1, a single broker loss destroys offset/group state for every group hashed to its partitions — failover can't restore state, so groups reset to `auto.offset.reset` (reprocessing or skipping data) and may be unable to elect a coordinator until the broker returns. Other relevant settings: `offsets.topic.num.partitions` (50, fixes coordinator spread — changing it later is disruptive), and the topic uses **log compaction** (`cleanup.policy=compact`) so the latest offset per key survives. A classic incident is a 3-broker cluster created with the offsets topic at RF=1 because it was bootstrapped while only one broker was up.
go deeper
Know committed offsets are stored in the __consumer_offsets topic.
Explain that coordinator state lives there and replication factor protects it.
Diagnose the RF=1 auto-creation pitfall and its data-loss/duplicate consequences.
Set durability policy (RF, ISR, monitoring) for internal topics and reason about failover recovery semantics.
## The coordinator's state lives in a topic The Group Coordinator is stateful: it holds, per group, the committed offsets and the group metadata (members, generation, assignment). But that state is **not** kept only in broker memory — it is written to the internal **`__consumer_offsets`** topic as the durable source of truth. When a coordinator fails over, the new coordinator **replays its `__consumer_offsets` partition** to rebuild everything. So the durability of this topic *is* the durability of coordination. ## Why replication factor is critical - **`offsets.topic.replication.factor`** (default **3**) controls how many copies of each offsets partition exist. - If it is **1** and the broker holding a partition dies, the offsets/metadata for every group hashed to that partition are **gone**. The partition has no other in-sync replica to take leadership, so: - No new coordinator can be elected for those groups until the broker recovers. - When consumers eventually reconnect with lost offsets, they fall back to **`auto.offset.reset`** (`earliest` → reprocess everything, or `latest` → skip data and lose messages). - **Gotcha:** the offsets topic is auto-created on first use using the *currently available* broker count if it is below the configured RF. Bootstrapping a 3-broker cluster while only one broker is up can silently create the topic at **RF=1**, leaving a latent single point of failure. Fix it by raising RF via a reassignment, or set the config before first consumer use. ## Other relevant configuration - **`offsets.topic.num.partitions`** (default **50**): determines how groups spread across coordinators via `hash(group.id) % num`. Changing it after groups exist re-maps groups to different partitions/coordinators and is operationally disruptive — treat it as fixed. - **`cleanup.policy=compact`**: the offsets topic is **log-compacted**, so only the latest record per key (group+topic+partition for offsets, or group metadata) is retained — keeping it bounded while preserving current state. - **`min.insync.replicas`** + producer acks for internal commits affect how many replicas must ack an offset commit before it's durable. ## What goes wrong in practice - RF=1 offsets topic → broker loss → mass offset loss → duplicate processing or data skips across many groups at once. - Under-replicated offsets partitions → coordinator failover stalls, groups stuck in rebalance, `COORDINATOR_NOT_AVAILABLE` errors. - Treat `__consumer_offsets` like a production-critical topic: RF ≥ 3, healthy ISR, and monitored under-replicated-partition counts.
- How can a production 3-broker cluster end up with an RF=1 __consumer_offsets topic?If the topic is auto-created while fewer brokers are up than offsets.topic.replication.factor, Kafka creates it at the available broker count. Bootstrapping with only one broker live yields RF=1 — a latent single point of failure that survives even after the other brokers join.
- Why is the __consumer_offsets topic log-compacted rather than time-retained?Only the latest committed offset (and latest group metadata) per key matters, so compaction keeps the current state indefinitely while discarding superseded records, bounding the topic size.
saying these in an interview costs you the question
- Assuming the offsets topic always honors RF=3 regardless of cluster size at creation time.
- Saying coordinator state is purely in-memory and lost on any failover (it's recovered from the compacted log).
- Recommending changing offsets.topic.num.partitions on a live cluster as harmless.