Kafka Multiple Consumers: Consumer Groups and Partitions

Kafka consumers with the same group.id split the partitions of a topic, while each consumer group gets every message. Real runs with 1 to 4 consumers, rebalance logs, assignment strategies, the new KIP-848 protocol and Spring Boot @KafkaListener concurrency.

kafka-logo

In Kafka, multiple consumers share the work through consumer groups, and every consumer joins a group with its group.id. Consumers with the same group.id split the partitions of a topic, so each message goes to only one of them. Consumers with different group.id values do not share anything, so each group gets every message.

We put several consumers in one group to spread the load of one application over several instances, and we add a second group when another application, such as a shipping service, needs the same messages.

Topic orders with partitions 0, 1 and 2; in group billing, billing-1 reads partition 0, billing-2 reads partition 1 and billing-3 reads partition 2; in group shipping, shipping-1 reads all three partitions
The billing consumers split the 3 partitions. The shipping group still gets every order.

The following example is a normal Kafka consumer that joins the group billing, and its group.id is the only setting that decides how the consumers share the topic.

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "billing");          // same id = share the partitions
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Set.of("orders"));
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));   // only records from this consumer's partitions

// 3 partitions, 3 consumers in "billing"   -> 1 partition each
// 3 partitions, 4 consumers in "billing"   -> billing-4 gets no partition
// one more consumer in group "shipping"     -> reads all 3 partitions again

Notice the last three comments, which show the limits of a group. A fourth consumer in a group with 3 partitions gets no partition, whereas a consumer in a new group reads all partitions again.

Next, we start 1, 2, 3 and 4 consumers in one group on a Kafka 4.3.1 broker, add a second group, and stop a consumer to see the rebalance logs. Then we compare the partition assignment strategies, try the new consumer group protocol from KIP-848, and build the same setup with Spring Boot @KafkaListener.

1. What Is a Kafka Consumer Group?

A Kafka topic is split into partitions, and each partition is one ordered list of messages.

When we run several copies of the same application, we want each copy to do part of the work. For example, an online shop runs three instances of its billing service, and each instance should bill a different part of the orders. A consumer group splits the work, because Kafka gives each partition to only one consumer of the group.

TermMeaningIn our example
ConsumerOne KafkaConsumer object that reads messagesbilling-1, billing-2, …
group.idThe name of the group the consumer joinsbilling and shipping
PartitionOne ordered part of a topicTopic orders has partitions 0, 1 and 2
AssignmentThe list of partitions one consumer ownsbilling-1 owns partition 0
RebalanceKafka moves partitions when a consumer joins or leavesbilling-3 stops, its partition moves
Group coordinatorThe broker that tracks the members of a groupOur single broker

Two rules follow from the way Kafka gives out the partitions.

  • Inside one group, a partition is never read by two consumers at the same time. So the number of partitions is the upper limit for useful consumers in a group.
  • Kafka does not delete a message after a group reads it. Each group stores its own position (offset) in every partition, so a second group reads the same messages starting from its own position.

2. Setting Up Kafka and the Example Project

We need only one Kafka broker, so the example uses the official apache/kafka:4.3.1 image in KRaft mode, with one node, based on our Kafka cluster setup with Docker Compose. Applications on the host connect to localhost:9092.

docker compose up -d

The plain Java examples need only the Kafka client library and a logger.

<dependency>
  <groupId>org.apache.kafka</groupId>
  <artifactId>kafka-clients</artifactId>
  <version>4.3.1</version>
</dependency>
<dependency>
  <groupId>org.slf4j</groupId>
  <artifactId>slf4j-simple</artifactId>
  <version>2.0.20</version>
</dependency>

The kafka-clients 4.3.1 library pulls in slf4j-api 1.7.36, and with slf4j-simple 2.x Kafka then prints no log lines at all. To fix this, we add slf4j-api 2.0.20 as a direct dependency, as the example project does.

The demo program creates the topic orders with 3 partitions and starts each consumer on its own thread. Each consumer prints the partitions it gets and every order it reads.

A producer sends 6 orders with the keys order-1 to order-6, and the key decides the partition. In every run, the orders go to the same partitions, because the same key always gives the same partition.

PartitionOrders
0order-2, order-3, order-6
1order-1
2order-4, order-5

3. Multiple Consumers in the Same Group

We start N consumers with group.id=billing and wait until Kafka has assigned all partitions, and then send the 6 orders.

mvn -q compile exec:exec -Ddemo="group 1"
mvn -q compile exec:exec -Ddemo="group 2"
mvn -q compile exec:exec -Ddemo="group 3"
mvn -q compile exec:exec -Ddemo="group 4"

3.1. One Consumer Reads All Partitions

With one consumer in the group, that consumer owns all 3 partitions.

billing-1 assigned [0, 1, 2]
--- summary ---
billing-1  partitions [0, 1, 2] orders [order-1, order-2, order-3, order-4, order-5, order-6]

3.2. Two and Three Consumers Split the Partitions

With two consumers, the 3 partitions cannot be split evenly, so one consumer gets 2 partitions.

billing-1  partitions [0, 1]    orders [order-1, order-2, order-3, order-6]
billing-2  partitions [2]       orders [order-4, order-5]

With three consumers, each consumer gets one partition, and each order still arrives only once in the group.

billing-1  partitions [0]       orders [order-2, order-3, order-6]
billing-2  partitions [1]       orders [order-1]
billing-3  partitions [2]       orders [order-4, order-5]

Kafka splits partitions, not messages. The consumer billing-2 got only one order because partition 1 had only one order, so for even work we need keys that spread evenly over the partitions.

3.3. The Fourth Consumer Stays Idle

A fourth consumer has no partition left, so billing-4 stays in the group and waits.

billing-4 assigned []
--- summary ---
billing-1  partitions [0]       orders [order-2, order-3, order-6]
billing-2  partitions [1]       orders [order-1]
billing-3  partitions [2]       orders [order-4, order-5]
billing-4  partitions []        orders []
Four panels for topic orders with 3 partitions: 1 consumer owns partitions 0, 1 and 2; 2 consumers own 0 and 1, and 2; 3 consumers own one partition each; with 4 consumers billing-4 is idle
Adding consumers helps only up to the number of partitions.

The idle consumer is not useless, because when another consumer stops, Kafka gives the stopped consumer’s partition to the idle one. For more parallel work, we add partitions to the topic first.

3.4. Checking the Assignment With kafka-consumer-groups.sh

The kafka-consumer-groups.sh tool shows which consumer owns which partition, so we keep 4 consumers running with the watch 4 demo and ask the broker.

docker exec kafka /opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server kafka:19092 \
  --describe --group billing --members
GROUP           CONSUMER-ID                                             HOST            CLIENT-ID          #PARTITIONS
billing         consumer-billing-4-0ae4e02c-d898-401d-ac47-7d623eeacfbe /172.19.0.1     consumer-billing-4 0
billing         consumer-billing-2-fa1969cb-da6e-4b2b-a9ce-cf99b5535670 /172.19.0.1     consumer-billing-2 1
billing         consumer-billing-3-74c3b57e-a05a-4d95-9272-725419bb8f25 /172.19.0.1     consumer-billing-3 1
billing         consumer-billing-1-11f2d6dd-110b-4142-bd49-f5ea89e82055 /172.19.0.1     consumer-billing-1 1

Without –members, the tool prints one line per partition with the owner, the committed offset and the lag, which is the number of messages the group has not read yet.

4. Multiple Consumer Groups on One Topic

Next, we start two groups at the same time. The group billing has 2 consumers, and the group shipping has 1 consumer.

startGroup(bootstrap, topic, "billing", 2, "classic");    // billing-1, billing-2
startGroup(bootstrap, topic, "shipping", 1, "classic");   // shipping-1

Both groups receive all 6 orders. Inside billing, the orders are split, whereas shipping-1 is alone in its group, so it owns all 3 partitions.

billing-1  partitions [0, 1]    orders [order-1, order-2, order-3, order-6]
billing-2  partitions [2]       orders [order-4, order-5]
shipping-1 partitions [0, 1, 2] orders [order-1, order-2, order-3, order-4, order-5, order-6]

Two groups on one topic give us the publish-subscribe pattern. We use one group per application or per purpose, for example one group for billing and one for shipping, and we use several consumers inside a group to share the load of one application.

GoalWhat to do
Share the load of one applicationSame group.id in every instance
Every application gets every messageA different group.id per application
Read everything again with a new serviceA new group.id with auto.offset.reset=earliest

5. Rebalancing When a Consumer Leaves

A rebalance happens when a consumer joins or leaves the group. A consumer leaves when it calls close() or when it stops sending heartbeats. A heartbeat is a small “I am alive” request that the consumer sends to the broker every few seconds.

For example, during a rolling deployment, each instance of a billing service stops and starts again, and every stop causes a rebalance. In the rebalance demo, 3 consumers start in billing, and then billing-3 closes.

mvn -q compile exec:exec -Ddemo=rebalance

The example enables INFO logs for the class ConsumerCoordinator, so we see what the Kafka client does, and the log lines are trimmed to the important part.

[billing-3] Member consumer-billing-3-515d...7d4 sending LeaveGroup request to coordinator localhost:9092 due to the consumer is being closed
[billing-1] Request joining group due to: group is already rebalancing
[billing-2] Request joining group due to: group is already rebalancing
billing-1 revoked  [0]
billing-2 revoked  [1]
[billing-1] Successfully joined group with generation Generation{generationId=9, ..., protocol='range'}
[billing-2] Successfully joined group with generation Generation{generationId=9, ..., protocol='range'}
[billing-2] Finished assignment for group at generation 9: {consumer-billing-1-...=Assignment(partitions=[orders-0, orders-1]), consumer-billing-2-...=Assignment(partitions=[orders-2])}
billing-1 assigned [0, 1]
billing-2 assigned [2]

The log shows five steps.

  1. The consumer billing-3 sends a LeaveGroup request when we call close().
  2. The broker starts a rebalance, and billing-1 and billing-2 learn about it with their next heartbeat.
  3. Both remaining consumers release all their partitions, even the ones that will not move, and our ConsumerRebalanceListener prints the revoked lines.
  4. The consumers join again. The generationId grows from 8 to 9, because every rebalance creates a new generation.
  5. One member, the group leader, computes the new assignment with the range strategy, so billing-1 owns partitions 0 and 1.

All 6 orders sent after the rebalance arrive, so no order is lost.

billing-1  partitions [0, 1]    orders [order-1, order-2, order-3, order-6]
billing-2  partitions [2]       orders [order-4, order-5]

A consumer that crashes cannot send a LeaveGroup request, so the broker waits for session.timeout.ms (45 seconds by default) without heartbeats before it starts the rebalance. Until then, nobody reads the partitions of the crashed consumer.

A ConsumerRebalanceListener runs our code during a rebalance. For example, the listener can commit offsets before a partition moves to another consumer.

consumer.subscribe(Set.of(topic), new ConsumerRebalanceListener() {
  @Override
  public void onPartitionsRevoked(Collection<TopicPartition> revoked) {
    System.out.printf("%s revoked  %s%n", name, numbers(revoked));
  }

  @Override
  public void onPartitionsAssigned(Collection<TopicPartition> assigned) {
    System.out.printf("%s assigned %s%n", name, numbers(assigned));
  }
});

6. Partition Assignment Strategies

The assignment strategy decides which consumer gets which partition. With the classic protocol, the consumers choose it with partition.assignment.strategy. Kafka 4.3 includes four strategies, and only one of them keeps the other consumers reading during a rebalance.

StrategyHow it assignsStops all consumers during a rebalance?
RangeAssignorSplits each topic into ranges; consumer 1 gets the first rangeYes
RoundRobinAssignorGives partitions to the consumers one at a time, in turnYes
StickyAssignorBalanced, and keeps partitions where they wereYes
CooperativeStickyAssignorLike sticky, but moves only the partitions that must moveNo

The default value in the 4.3.1 client is the list [RangeAssignor, CooperativeStickyAssignor]. The consumers use the first strategy in the list, so the default is RangeAssignor.

The rebalance log in section 5 shows protocol=’range’, and CooperativeStickyAssignor is in the list only so that a group can switch to it with a rolling restart.

To use the cooperative strategy, we set it as the only entry.

props.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, "classic");
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
    CooperativeStickyAssignor.class.getName());

We repeat the rebalance demo with this setting, and this time billing-1 and billing-2 print no revoked line. Both keep their partitions and keep reading, and only partition 2 of billing-3 moves.

--- stopping billing-3 ---
billing-3 revoked  [2]
[billing-1] Successfully joined group with generation Generation{generationId=3, ..., protocol='cooperative-sticky'}
[billing-1] Notifying assignor about the new Assignment(partitions=[orders-0, orders-2])
billing-1 assigned [2]
[billing-2] Notifying assignor about the new Assignment(partitions=[orders-1])
billing-2 assigned []

With a cooperative strategy, onPartitionsAssigned() receives only the newly added partitions. So billing-1 assigned [2] means “partition 2 was added to partition 0”. To get the full list, we call consumer.assignment().

7. The New Consumer Group Protocol (KIP-848)

Kafka 4.0 made the next generation consumer rebalance protocol, KIP-848, generally available.

With this protocol, the broker computes the assignment instead of a consumer, and it moves partitions one by one through the heartbeats. No consumer has to stop for a rebalance.

The broker side is on by default since Kafka 4.0, but the client side is not. In kafka-clients 4.3.1, group.protocol still defaults to classic.

We turn the new protocol on with one setting.

props.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, "consumer");   // default: "classic"

7.1. What Changes With group.protocol=consumer

Several classic settings move to the broker, so the client rejects them with the new protocol.

org.apache.kafka.common.config.ConfigException: partition.assignment.strategy cannot be set when group.protocol=CONSUMER
Classic client settingWith group.protocol=consumer
partition.assignment.strategyNot allowed; the broker uses group.consumer.assignors (uniform,range), and uniform is the default
Choose a strategy per groupClient setting group.remote.assignor, for example range
heartbeat.interval.ms (3 s)Broker setting group.consumer.heartbeat.interval.ms (5 s)
session.timeout.ms (45 s)Broker setting group.consumer.session.timeout.ms (45 s)

The broker values in the table come from kafka-configs.sh –describe –all on our 4.3.1 broker.

7.2. Rebalancing With the New Protocol

We run the same rebalance demo with the new protocol, and this time the log comes from ConsumerMembershipManager.

--- stopping billing-3 ---
billing-3 revoked  [1]
[billing-3] Member 2LuUsVIjSoK_svcP0oAXvg with epoch 18 transitioned from STABLE to PREPARE_LEAVING.
[billing-3] Member 2LuUsVIjSoK_svcP0oAXvg with epoch -1 transitioned from LEAVING to UNSUBSCRIBED.
[billing-2] Member ABb-almRQXG0kXfp3mjPOg with epoch 19 transitioned from STABLE to RECONCILING.
[billing-2] Reconciling assignment with local epoch 2
billing-2 assigned [1]
--- summary ---
billing-1  partitions [0]       orders []
billing-2  partitions [1, 2]    orders []

The consumer billing-1 does not appear in the log at all, because partition 0 stayed with billing-1 and its poll loop never stopped.

The consumer billing-2 keeps partition 2 and only adds partition 1. The member epoch counts the assignment versions, like the generationId of the classic protocol.

Three rows for billing-3 leaving a group of 3: with the default RangeAssignor all consumers release their partitions and billing-1 ends with 0 and 1; with CooperativeStickyAssignor and with group.protocol=consumer the remaining consumers keep reading and only the partition of billing-3 moves
With the default classic setup, every consumer stops during a rebalance. The cooperative strategy and the new protocol move only what must move.

In a small demo, the new protocol needs more time to finish the assignment. When 3 consumers start at the same moment, the first one gets all 3 partitions, and the broker moves 2 of them with the next heartbeats. In our run, the move took 5 to 10 seconds.

billing-1 assigned [0, 1, 2]
billing-2 assigned []
billing-3 assigned []
billing-1 revoked  [1, 2]
billing-1 assigned []
billing-3 assigned [1]
billing-2 assigned [2]

The tool kafka-consumer-groups.sh shows the group type and the server-side strategy.

GROUP                     TYPE
shipping                  Classic
billing                   Consumer

GROUP           COORDINATOR (ID)          ASSIGNMENT-STRATEGY  STATE                #MEMBERS
billing         kafka:19092  (1)          uniform              Stable               3

A group switches between Classic and Consumer when it is empty, so our billing group switched between the demo runs without any extra step.

8. One KafkaConsumer per Thread

A KafkaConsumer is not thread-safe. Two threads must never call poll() on the same consumer, and when we try it, the second thread fails at once.

java.util.ConcurrentModificationException: KafkaConsumer is not safe for multi-threaded access. currentThread(name: worker-2, id: 24) otherThread(id: 23)

The fix is one consumer per thread. Each OrderConsumer in the example creates its own KafkaConsumer and runs its poll loop on its own thread.

// OrderConsumer: the constructor creates its own KafkaConsumer
thread = Thread.ofPlatform().name(name).start(this);

// run(): the poll loop, only on this thread
while (true) {
  for (ConsumerRecord<String, String> rec : consumer.poll(Duration.ofMillis(100))) {
    received.add(rec);
  }
}

// stop(): called from another thread
consumer.wakeup();   // the only thread-safe method; poll() throws WakeupException
thread.join();       // run() catches it and calls consumer.close()

To start 3 consumers, we create 3 OrderConsumer objects, and each one joins the group as a separate member, as the logs in section 5 show.

When processing a record is slow, a common alternative is one consumer thread that passes records to a thread pool. In that case, the consumer thread must commit offsets only for records the pool has finished.

9. Multiple Consumers in Spring Boot

Spring for Apache Kafka creates the consumers for us. The concurrency attribute of @KafkaListener starts that many KafkaConsumer objects, each on its own thread, all in the same group.

A second listener with another groupId forms a second group, and the rest is the basic Spring Boot with Kafka setup.

@KafkaListener(id = "billing", topics = "orders", groupId = "billing", concurrency = "3")
void billing(ConsumerRecord<String, String> rec) {
  log.info("billing  read {} from partition {}", rec.key(), rec.partition());
}

@KafkaListener(id = "shipping", topics = "orders", groupId = "shipping")
void shipping(ConsumerRecord<String, String> rec) {
  log.info("shipping read {} from partition {}", rec.key(), rec.partition());
}

The application sends the same 6 orders at startup. The first column is the thread name, so billing-0-C-1, billing-1-C-1 and billing-2-C-1 are the 3 billing consumers, and each one reads its own partition.

billing-1-C-1            billing  read order-1 from partition 1
shipping-0-C-1           shipping read order-1 from partition 1
billing-0-C-1            billing  read order-2 from partition 0
shipping-0-C-1           shipping read order-2 from partition 0
billing-0-C-1            billing  read order-3 from partition 0
shipping-0-C-1           shipping read order-3 from partition 0
shipping-0-C-1           shipping read order-4 from partition 2
billing-2-C-1            billing  read order-4 from partition 2
shipping-0-C-1           shipping read order-5 from partition 2
billing-2-C-1            billing  read order-5 from partition 2
billing-0-C-1            billing  read order-6 from partition 0
shipping-0-C-1           shipping read order-6 from partition 0

The partition limit also applies here, so a concurrency higher than the number of partitions creates idle consumers. Spring Boot passes any Kafka client setting through spring.kafka.consumer.properties, including the new protocol.

spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.consumer.auto-offset-reset=earliest
spring.kafka.consumer.properties[group.protocol]=consumer

A test of the example project runs both listeners with this property, and the 3 billing consumers end with one partition each, as with the classic protocol. Spring also supports class-level listeners with @KafkaListener and @KafkaHandler.

10. Testing Consumer Groups With Testcontainers

The example project tests every consumer group rule against a real broker. Testcontainers starts apache/kafka:4.3.1 in Docker for the test class.

@Container
static final KafkaContainer KAFKA = new KafkaContainer("apache/kafka:4.3.1");

@ParameterizedTest(name = "{0} consumers, protocol {1}")
@CsvSource({
    "1, classic,  3",
    "2, classic,  2;1",
    "3, classic,  1;1;1",
    "4, classic,  1;1;1;0",
    "3, consumer, 1;1;1",
    "4, consumer, 1;1;1;0"})
void consumersInOneGroupSplitThePartitions(int count, String protocol, String expected)

The class has 13 tests, and each test creates its own topic so the tests do not affect each other. The tests check five rules.

  • 1 to 4 consumers get 3, 2+1, 1+1+1 and 1+1+1+0 partitions, with both protocols.
  • Two groups each receive all 6 orders.
  • After one consumer leaves, the remaining consumers own all 3 partitions.
  • With the cooperative strategy and the new protocol, the remaining consumers keep their partitions.
  • Setting partition.assignment.strategy with group.protocol=consumer throws a ConfigException.

11. Kafka Multiple Consumers FAQs

11.1. Can Two Consumers in the Same Group Read the Same Partition?

No. Inside one group, Kafka gives each partition to only one consumer, so two consumers can read the same partition only when they are in different groups.

Since Kafka 4.2, share groups (KIP-932, “Queues for Kafka”) are production-ready. In a share group, several consumers can read from one partition, and share groups use a different class, KafkaShareConsumer.

11.2. How Many Consumers Should a Group Have?

At most as many as the topic has partitions, because extra consumers stay idle, as billing-4 did. A few idle consumers are fine as spares, but for more parallel work we increase the number of partitions first.

11.3. Is group.id Required?

Yes, for subscribe(). Without a group.id, the call fails with an InvalidGroupIdException.

org.apache.kafka.common.errors.InvalidGroupIdException: To use the group management or offset commit APIs, you must provide a valid group.id in the consumer configuration.

A consumer without a group can still read with assign(), where we choose the partitions ourselves. In that case, Kafka does no rebalancing, and the consumer cannot commit offsets.

11.4. Do Multiple Consumers Keep the Message Order?

Within one partition, yes, because each partition has one reader in the group and that reader gets the messages in order. Across partitions, there is no order, but messages with the same key go to the same partition, so all messages of order-1 stay in order.

11.5. What Happens to Messages When a Consumer Crashes?

No message is lost, because the messages stay in the partition. After the session timeout, another consumer of the group gets the partition and starts from the last committed offset. Messages after that offset can be processed twice, so the processing code should handle duplicates.

12. Conclusion

Consumers with the same group.id split the partitions of a topic, and each group gets every message. The number of partitions limits how many consumers in a group can work at the same time.

With default settings, a rebalance stops every consumer in the group. CooperativeStickyAssignor or group.protocol=consumer moves only the partitions that must move.

For new applications on Kafka 4.x, we should test the new protocol, and in every setup each KafkaConsumer is used by only one thread.

13. References

Happy Learning !!

Source Code on Github

Leave a Comment

  1. ListenableFuture is deprecated in Boot 3.1.5 with JDK 17
    kafkaTemplate.send() doesnt return the ListenableFuture

Comments are closed.

About Us

HowToDoInJava provides tutorials and how-to guides on Java and related technologies.

It also shares the best practices, algorithms & solutions and frequently asked interview questions.