Lesson 9 / 25

Consumer Groups and Committed Offsets

Explain how a group shares partitions and how committed offsets and lag are tracked.

One partition, one consumer per group

Consumers that share the same group.id form a consumer group. Kafka assigns each partition to exactly one consumer in the group, so the group processes the topic in parallel, while a different group reads the full topic independently. As a consumer works it commits the offset of the next record to read, stored in an internal topic (__consumer_offsets). Lag is the log-end offset minus the committed offset per partition: a growing lag means consumers are falling behind. A group cannot use more consumers than partitions; extra consumers sit idle.

Share the work, track the position

Consumers in a group split a topic's partitions; Kafka remembers each group's committed offsets.

Four ideas: group, assignment, offset, rebalance.
Figure 3.1 — Group, assignment, offset and rebalance.

Group lag from a real broker, run

I read 4 of the 7 records with group billing, then ran kafka-consumer-groups.sh --describe. Partitions 1 and 2 are fully read (lag 0); partitions 4 and 5 still have 1 and 2 unread records (lag 3 in total). The group has no active members because the console consumer exited.

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic orders \
  --group billing --from-beginning --max-messages 4
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing

Output:

Consumer group 'billing' has no active members.

GROUP   TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
billing orders  0          0               0               0
billing orders  1          1               1               0
billing orders  2          3               3               0
billing orders  3          0               0               0
billing orders  4          0               1               1
billing orders  5          0               2               2

Lag arithmetic, run

Total lag is the sum over partitions; here 0 + 80 + 420 = 500 records behind.

log_end   = {0: 1500, 1: 1480, 2: 1520}
committed = {0: 1500, 1: 1400, 2: 1100}
lag = {p: log_end[p] - committed[p] for p in log_end}
print(lag, sum(lag.values()))

Output:

{0: 0, 1: 80, 2: 420} 500

Quick check: What is consumer lag?

  • Network latency to the broker
  • The age of the topic
  • Log-end offset minus the committed offset
  • The number of partitions
Answer

Log-end offset minus the committed offset — It counts how many records are written but not yet processed by the group.