Consumers with the same group.id divide the subscribed partitions: each partition goes to exactly one member, which therefore sees its records in order. The group's coordinator, the broker leading the group's partition of __consumer_offsets, tracks members and offsets. An assignor decides who gets what: range, roundrobin, sticky and cooperative-sticky run in a client under the classic protocol; uniform and range run in the broker under the new one (The New Protocol (KIP-848)). Here four members share the three-partition topic.
import subprocess, threading, time
from confluent_kafka import Consumer
def member(stop):
c = Consumer({"bootstrap.servers": "localhost:33092", "group.id": "booknest-four",
"group.protocol": "consumer", "auto.offset.reset": "earliest"})
c.subscribe(["booknest.order-events"])
while not stop.is_set():
c.poll(0.5)
c.close()
stop = threading.Event()
threads = [threading.Thread(target=member, args=(stop,)) for _ in range(4)]
for t in threads: t.start()
time.sleep(15) # let all four join
subprocess.run("docker exec l3-kafka /opt/kafka/bin/kafka-consumer-groups.sh --describe "
"--bootstrap-server l3-kafka:9092 --group booknest-four --members --verbose "
"| awk 'NF {print $5, $6, $7, $8}' | column -t", shell=True)
stop.set()
for t in threads: t.join()Output
#PARTITIONS CURRENT-EPOCH CURRENT-ASSIGNMENT TARGET-EPOCH 0 5 - 5 1 5 booknest.order-events:0 5 1 5 booknest.order-events:1 5 1 5 booknest.order-events:2 5
The fourth member is idle: a group cannot use more consumers than partitions (Topic Configuration). The epochs come from the new protocol, which numbers each assignment.