Apache Kafka for Data & AI Engineers

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 fourOne 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 / Note

Within 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

python
if 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

python
consumer_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:

python
def 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:

python
consumer.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_assign runs when this instance is handed partitions.
  • on_revoke runs 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

bash
docker 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:

bash
docker exec -it kafka-1 bash export PATH=$PATH:/opt/kafka/bin
bash
kafka-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:

bash
python 5.2_consumer_group_member.py instance-1
bash
python 5.2_consumer_group_member.py instance-2
bash
python 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:

bash
kafka-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:

bash
kafka-consumer-groups.sh --describe --group recommendation-service --members --bootstrap-server kafka-1:29092

This lists the members and how many partitions each one has.

bash
kafka-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:

bash
python 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:

bash
docker 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 coordinatesHow 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:

  1. Take the group.id and hash it.
  2. Map the hash to one partition of the internal topic __consumer_offsets (50 partitions by default).
  3. 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:

TermWhat it is
Partition leaderThe broker that holds the main copy of a partition
Group leaderOne 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 / Note

Good to know: Newer Kafka brokers (the apache/kafka:4.3.1 in our cluster is one) also offer a second protocol, where the broker computes the assignment and there is no group leader. confluent-kafka uses 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 consumer group 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 speaks classic unless you explicitly set group.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.id means a separate, full copy of the data
  • Ran three instances with the same group.id and watched partitions get assigned
  • Confirmed each instance received messages only from its own partition
  • Used kafka-consumer-groups.sh --describe, with --members and --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_assign callback

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.