Consumer Groups & Partition Assignment — Building a Multi-Consumer Group
The Requirement
Clickstream volume just tripled.
Your single RecommendationConsumer from the last lecture reads all three partitions of clickstream-events, alone. It is falling behind, and recommendations are stale by the time they update.
The fix product wants is simple: add more consumer capacity, so the work is shared.
So can we just start two more copies of our script?
Careful. It depends on how you start them. If the three copies are unrelated, each one reads everything. All three process every event, and your engine reacts to each click three times. That is not sharing the work. That is tripling it.
Why a Group, Not Just More Consumers
One consumer, a group of three, and a group of four
What makes consumers "related"?
One thing: the group.id. You set it in the last lecture ("recommendation-service") without a full explanation. Here is the explanation.
Output / NoteWithin one group, each partition is read by exactly one consumer.
Kafka splits the partitions between the members of the group. Three partitions and three consumers: each consumer gets one partition. Nothing is read twice, and you get the three-way parallelism you wanted.
Let's answer some questions you may already have.
What if I add a fourth consumer to this group?
A topic with 3 partitions can keep only 3 consumers of one group busy. The fourth one gets no partition. It sits idle as a standby, ready to take over if another member fails. So more consumers than partitions does not add speed.
So how do I scale beyond three consumers?
Add partitions to the topic. The number of partitions is the ceiling for parallelism inside one group.
What if the same order of events matters?
Good thinking. Kafka keeps order inside one partition. Since one partition goes to only one consumer, that consumer sees its events in order. If you use a key (for example user_id), all events of one user land in one partition, and so with one consumer.
What about a second group, like an analytics service?
Groups are independent. A group named analytics-service reads the same topic on its own, and gets every message again, with its own position. So:
- Same
group.id= share the work. - Different
group.id= each group gets a full copy of the data.
Let's Build It
We reuse almost everything from the last lecture: same config shape, same poll loop. There are only two additions. First, a way to tell running copies apart. Second, a way to see which partitions each copy receives.
Step 1: Accept an Instance Name
pythonif len(sys.argv) != 2: print("Usage: python 5.2_consumer_group_member.py <instance-name>") sys.exit(1) instance_name = sys.argv[1]
This is not a Kafka idea. It only helps you tell your terminal windows apart when three of them print at the same time.
Step 2: Same group.id, Different client.id
pythonconsumer_config = { "bootstrap.servers": BOOTSTRAP_SERVERS, # The same group.id in every instance is what makes them share work. "group.id": GROUP_ID, # Only a label, so you can tell instances apart in the terminal # and in kafka-consumer-groups.sh output. "client.id": instance_name, "auto.offset.reset": "earliest", }
Which line makes the magic happen?
group.id. Every instance uses the exact same value. That is all Kafka needs. There is no other signal, and the instances never talk to each other directly.
Then what is client.id for?
It is a label. It has no effect on how work is shared. But it shows up in Kafka's own tools, so you can see which instance owns which partition.
Step 3: Watch Assignment Happen, Live
In the last lecture, assignment happened silently. This time, let's watch it:
pythondef on_assign(consumer, partitions): """Runs (inside poll) when this instance is handed partitions.""" assigned = [p.partition for p in partitions] if assigned: print(f"\n>>> [{instance_name}] ASSIGNED partitions: {assigned}\n") else: print(f"\n>>> [{instance_name}] ASSIGNED no partitions (idle standby)\n") def on_revoke(consumer, partitions): """Runs (inside poll) when this instance is about to lose partitions.""" revoked = [p.partition for p in partitions] print(f"\n>>> [{instance_name}] REVOKED partitions: {revoked}\n")
Then we pass them when subscribing:
pythonconsumer.subscribe([TOPIC], on_assign=on_assign, on_revoke=on_revoke)
What are these functions?
They are callbacks. You give them to Kafka, and Kafka calls them at the right moment:
on_assignruns when this instance is handed partitions.on_revokeruns when this instance is about to lose partitions.
Who calls them?
Your own poll() does. Remember from the last lecture: poll() also runs callbacks. So these messages appear only while your loop is running.
Everything else in the script is the same as before: the poll loop, the JSON decoding, and close().
Running It
Step 1: Start Your Cluster
bashdocker compose start docker ps
Step 2: Check the Partition Count
Our whole demo assumes the topic has three partitions. Let's confirm it. Open a broker terminal:
bashdocker exec -it kafka-1 bash export PATH=$PATH:/opt/kafka/bin
bashkafka-topics.sh --describe --topic clickstream-events --bootstrap-server kafka-1:29092
You should see PartitionCount: 3. Keep this terminal open. We will use it again.
Step 3: Open Three Terminals
Download 5.2_consumer_group_member.py into your chapter-5 folder. Open three separate terminal windows on your host machine, and run one instance in each:
bashpython 5.2_consumer_group_member.py instance-1
bashpython 5.2_consumer_group_member.py instance-2
bashpython 5.2_consumer_group_member.py instance-3
Now watch the ASSIGNED partitions lines.
What should you see?
With three partitions and three instances, each instance ends up with exactly one partition. No overlaps, nothing left out.
Did instance-1 first get all three partitions, and then lose two?
It can happen, depending on how fast you start the instances. The first one to join may briefly own everything. Then, when the next one joins, you see REVOKED followed by a new ASSIGNED. This is normal. Kafka rearranged the split because the group changed. That is called a rebalance, and it is the topic of the next lecture. For now, just watch the final result.
Will instance-1 always get partition 0?
Not guaranteed, but very often yes with these names. The default strategy sorts members by their member ID, and that ID starts with the client.id. So instance-1, instance-2, instance-3 usually get partitions 0, 1, 2 in this order.
Step 4: Watch the Load Split
The group already read all the old events in the last lecture. So the three terminals will stay quiet. That is expected.
In a fourth terminal, run one of your clickstream producer scripts from the producer chapter to send new events. Watch the three consumer terminals. Each one prints only the events from its own partition, and never from the others.
Step 5: Confirm It From the Cluster's Side
Go back to your broker terminal:
bashkafka-consumer-groups.sh --describe --group recommendation-service --bootstrap-server kafka-1:29092
This is the cluster's own view of the group. Look at three columns: PARTITION, CLIENT-ID, and LAG. You can see which instance owns which partition, and how far behind each one is. This is the exact tool you will use in production to check the health of a consumer group.
Two more views of the same group:
bashkafka-consumer-groups.sh --describe --group recommendation-service --members --bootstrap-server kafka-1:29092
This lists the members and how many partitions each one has.
bashkafka-consumer-groups.sh --describe --group recommendation-service --state --bootstrap-server kafka-1:29092
This shows the group's state (it should be Stable), its assignment strategy (range), and its coordinator: a specific broker. Remember this line. We explain it in the next section.
Step 6: Try a Fourth Instance (Optional)
Open one more terminal and start a fourth member:
bashpython 5.2_consumer_group_member.py instance-4
You should see the group rebalance, and the new instance report that it has no partitions (idle standby). Re-run the --members command: the fourth instance appears with 0 partitions. This proves the rule from earlier: more consumers than partitions adds no speed.
Step 7: Stop Everything
Press Ctrl+C in each consumer terminal, then:
bashdocker compose stop
How Group Coordination Actually Works
You just watched three separate scripts reach a clean, non-overlapping split, and they never talked to each other. How?
How a consumer group coordinates
Step 1: Find the Group Coordinator
Every group has a Group Coordinator. Who is it?
It is one specific broker. Kafka picks it with a simple rule:
- Take the
group.idand hash it. - Map the hash to one partition of the internal topic
__consumer_offsets(50 partitions by default). - The broker that is the leader of that partition is the coordinator.
This is the same leader and partition idea you learned in the chapter on brokers. It is just applied to an internal topic. And every member gets the same answer, because the rule uses only the group.id. That is why the coordinator column in your --state output showed one broker.
Step 2: Every Instance Sends JoinGroup
Each instance sends a JoinGroup request to the coordinator. It says: "I want to join this group, and here are the topics I subscribe to."
Step 3: One Member Becomes the Group Leader
Does the coordinator decide who gets which partition?
No, and this surprises many people. The coordinator picks one member of the group and names it the group leader. Usually this is the first member that joined. The coordinator sends the leader the full member list and the subscriptions. Then the leader itself runs the assignment strategy, on the client side.
Do not mix up two ideas:
| Term | What it is |
|---|---|
| Partition leader | The broker that holds the main copy of a partition |
| Group leader | One consumer in the group that computes the assignment |
Which strategy does it run?
In confluent-kafka, the default is range (the --state output showed it). For one topic, range sorts the members and splits the partitions into equal chunks. With 3 partitions and 3 members, each gets one. With 3 partitions and 4 members, the extra member gets none.
range is not the only strategy — it's just the one this group happened to run, so it's the one we can point at on screen right now. The next lecture is entirely about why the strategy you pick matters: we'll put a second, very different strategy called cooperative-sticky next to range and watch them behave completely differently during a rebalance.
Step 4: SyncGroup Delivers the Result
The group leader sends the finished assignment back to the coordinator, in a SyncGroup request. The other members send SyncGroup too, and the coordinator answers each member with its own slice of the plan.
When that answer reaches your consumer, your poll() runs your on_assign callback. That is the exact moment you saw >>> ASSIGNED partitions in your terminal. It was not a coincidence. It is the visible result of this handshake.
Output / NoteGood to know: Newer Kafka brokers (the
apache/kafka:4.3.1in our cluster is one) also offer a second protocol, where the broker computes the assignment and there is no group leader.confluent-kafkauses the classic protocol described above unless you switch it on. Everything in this course uses the classic one.Out of scope in this course. That second protocol has a name: the
consumergroup protocol (Kafka calls its introduction KIP-848, the "Next Generation Consumer Rebalance Protocol"). We stop at naming it and go no further, on purpose:
- It is not the client default.
confluent-kafka(like every mainstream client) still speaksclassicunless you explicitly setgroup.protocol: consumer. Teaching it would mean explaining machinery you are not actually exercising.- It changes who computes the assignment, not what a developer does day to day. Your code —
subscribe(),on_assign,on_revoke, the poll loop — looks the same either way. The payoff for a working developer is small next to what it would cost to explain properly.- It deserves its own treatment once you're comfortable with the classic handshake in this lecture and the rebalancing behavior in the next one. Learning the new protocol before you understand what it's replacing would be memorizing a diagram, not understanding a trade-off.
And when does all this run again?
Not just at startup. Whenever the group changes, it runs again. That is the subject of the next lecture.
Common Mistakes at This Stage
Starting "more consumers" with different group.id values. Each one then gets every event, and the work is duplicated instead of shared.
Expecting more consumers than partitions to go faster. The extra ones just wait.
Giving every instance the same client.id. It works, but the tool output becomes hard to read. Use a unique label per instance.
Thinking the coordinator assigns the partitions. The group leader, one of your own consumers, computes the assignment.
Checklist
- Understood why unrelated consumers duplicate work instead of sharing it
- Understood the rule: within one group, each partition is read by exactly one consumer
- Understood that extra consumers beyond the partition count stay idle
- Understood that a different
group.idmeans a separate, full copy of the data - Ran three instances with the same
group.idand watched partitions get assigned - Confirmed each instance received messages only from its own partition
- Used
kafka-consumer-groups.sh --describe, with--membersand--state, to see the cluster's own view - Understood what the Group Coordinator is and how a broker gets picked for that role
- Understood JoinGroup, group leader, SyncGroup, and how they connect to your
on_assigncallback
You just scaled a consumer horizontally and watched Kafka coordinate it. Next: we make one of these instances disappear in the middle of the stream, and watch a real rebalance happen.