NOTE
2.7 Kafka Consumer
Kafka consumers, consumer groups, message-consumption flow, partition assignment and rebalancing, and duplicate/missed consumption.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What Is a Consumer?
- One of Kafka’s clients. It is responsible for consuming messages and pulling messages from Kafka.
1.1. Consumer Group
- A Consumer Group can have multiple Consumers.
- Each Consumer Group can independently consume all messages in a Topic.
- For a Topic, two Consumers in the same Group cannot consume the same Partition of that Topic.
-

-
Producer
.\kafka-console-producer.bat --broker-list localhost:9092 --topic first -
group1 Consumer
.\kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic first -
group2 Consumers
.\kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic first --consumer.config ..\..\config\consumer.properties.\kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic first --consumer.config ..\..\config\consumer.properties -
Result analysis
- Data produced by the producer is consumed by both group1 and group2, but only one of the group2 consumers can consume it.
-
2. Consumer Message-Consumption Flow
Message -> partitioner -> deserializer -> interceptor
2.1. Interceptor
- After a Consumer pulls messages from a Broker, it can perform preprocessing such as filtering messages, counting messages, or modifying message content.
2.2. Deserializer
- Messages pulled by a Consumer from a Broker are byte arrays and need to be converted back into the original messages through a deserializer.
- A custom deserializer must implement
org.apache.kafka.common.serialization.Deserializer.
2.3. Partitioner
- Determines which Partition a Consumer consumes messages from.
- RoundRobin
- Applies to all Topics.
- Consumers in the same group that subscribe to these Topics are sorted lexicographically by name.
- Partitions are then assigned to each consumer one by one in round-robin order.
- Range (default)
- Applies to one Topic.
- Consumers in the same group that subscribe to this Topic are sorted lexicographically by name.
- Let n = number of partitions / number of consumers, and m = number of partitions % number of consumers. Then each of the first m consumers is assigned n + 1 partitions, while each of the remaining (number of consumers - m) consumers is assigned n partitions.
- Problem: consumers may receive unequal amounts of data (the more Topics are subscribed to, the more serious the imbalance becomes).
2.4. Offset Commit
- Each message in a Partition has an offset representing the message’s position.
- There is also an offset concept on the Consumer side, representing the position of a consumed message in the partition. After consuming a message, the Consumer updates the offset on the broker (ZooKeeper in older versions).
3. Partition Rebalancing
3.1. What Is Partition Rebalancing?
- Consumer partition rebalancing occurs in the following situations, reallocating partitions to consumers in the consumer group:
- The number of members in the consumer group changes, for example when a new consumer instance joins or leaves the group.
- The number of Topics subscribed to by the consumer group changes.
- The number of partitions corresponding to the Topics subscribed to by the consumer group changes.
- The GroupCoordinator node corresponding to the consumer group changes.
- During rebalancing, consumers cannot consume messages, and the state of the current consumers is lost (which may cause duplicate consumption).
3.2. Why Is Partition Rebalancing Needed?
Try to distribute Leader Partitions evenly among consumers for consumption.
3.3. Principle of Partition Rebalancing
Handled by coordinators.
3.3.1. Coordinator
- The consumer side has a
Consumer Coordinator. - The Broker side has a
Group Coordinator.
3.3.1.1. Why Have Coordinators?
- If multiple consumers in the same consumer group specify different partition strategies, conflicts occur, so a coordinator is needed to handle them.
- Partition rebalancing also needs a component to assign partitions.
3.3.1.2. Coordinator Workflow
Consumers belonging to the same consumer group have the same group_id. Calculate hash(group_id) % number of partitions of the __consumer_offsets topic to obtain a leader partition, then use the broker hosting that partition as the coordinator. That coordinator then selects a consumer from the consumer group as the coordinator. After both coordinators have been selected, the partition-consumption plan is chosen. Once it is selected, the broker coordinator distributes the partition plan to all consumers.
- FIND_COORDINATOR
- JOIN_GROUP

- Elect the leader of the consumer group. If there is no leader, the first member to join the consumer group is regarded as the leader; if the leader goes down, select one at random.
- Elect the partition-assignment strategy.
- Collect all assignment strategies supported by each consumer to form a candidate set,
candidates. - Each consumer finds the first strategy in
candidatesthat it supports and casts one vote for that strategy. - Count the votes for each strategy in the candidate set. The strategy with the most votes becomes the assignment strategy for the current consumer group.
- Collect all assignment strategies supported by each consumer to form a candidate set,
- SYNC_GROUP
- HEARTBEAT
- Consumers send heartbeats to the GroupCoordinator to maintain their membership in the consumer group and their ownership of partitions.
4. How to Prevent Duplicate Consumption by Consumers
4.1. Why Duplicate Consumption Occurs
- Kafka consumer message semantics are
at least onceorat most once.at least once: the consumer goes down after processing a message but before committing the offset, causing duplicate consumption.
- This is the consumer’s own problem, not Kafka’s problem, and needs to be solved by the consumer itself.
4.2. How to Prevent Duplicate Consumption
4.2.1. Create a Deduplication Table
- Record the message ID after consumption and before committing the offset.
- Before consumption, check whether the ID has already been processed.
5. How to Prevent Consumers from Missing Messages
5.1. Why Messages Are Missed
- Kafka consumer message semantics are
at least onceorat most once.at most once: the offset is committed first, but the message has not finished processing, so the message is missed.
5.2. How to Prevent Missing Messages
5.2.1. Change the Consumer to Manual Offset Commit
- There are two ways for a consumer to commit offsets:
- synchronous commit:
consumer.commitSync(); - asynchronous commit:
consumer.commitAsync()orconsumer.commitAsync(callback).
- synchronous commit:
- Disable automatic offset commit with
enable.auto.commit=false, and manually commit the offset after processing is complete. - Duplicate consumption is preferable to missed consumption. Messages that repeatedly fail to be consumed can be placed into a dead-letter queue.

Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub