Kafka Consumer for Confluent Platform
An Apache Kafka® consumer is a client application that reads and processes events from Kafka topics. Consumer groups let many consumers cooperate to read from a topic’s partitions in parallel.
Offset management tracks each consumer’s position in the topic, and rebalancing redistributes partitions when consumers join or leave the group. This section covers the Kafka consumer, including consumer groups, offset management, and rebalancing, and introduces the configuration settings for tuning.
Ready to get started?
Sign up for Confluent Cloud, the fully managed cloud-native service for Apache Kafka® and get started for free using the Cloud quick start.
Download Confluent Platform, the self managed, enterprise-grade distribution of Apache Kafka and get started using the Confluent Platform quick start.
How consumers fetch messages
The Kafka consumer works by issuing “fetch” requests to the brokers leading the partitions it wants to consume. Each request specifies an offset in the log, and the consumer receives back a chunk of log that contains all the event records in that topic beginning from that offset. The consumer has significant control over this position and can rewind it to re-consume data if desired.
Consumer groups
A consumer group is a set of consumers that cooperate to consume topic
data. You add a consumer to a group by setting its group.id configuration property.
If you don’t set group.id, the subscribe and commit
methods of the
KafkaConsumer API
throw an exception when called.
When a consumer is assigned to a group, the topic partitions are distributed among the consumers in that group. New members can join, and existing members can leave a group. When this happens, Kafka reassigns the topic partitions to the consumers within the group, ensuring that each member receives a proportional share. This process is known as rebalancing; for more information, see Groups and rebalance protocols.
In a Kafka cluster, one of the brokers is designated as the group coordinator. This specific broker takes on the responsibility of managing a particular consumer group’s membership, tracking the heartbeats of the consumers within that group, and triggering rebalances when needed.
The broker that becomes the group coordinator is selected based on the consumer
group’s ID. The group’s ID is hashed to a specific partition of Kafka’s internal
offsets topic, __consumer_offsets. The broker which is the leader for that
particular partition is then assigned as the coordinator for that consumer
group. This mechanism ensures that the duties of group coordination are spread
out evenly across the various brokers in the Kafka cluster.
Groups and rebalance protocols
Rebalancing is the process of assigning, or reassigning, Kafka topic partitions to the members of a consumer group. Kafka supports two rebalance protocols that determine how rebalancing happens within a consumer group. These protocols are:
classic, the original protocol and the only generally available protocol before Kafka version 4.0
consumer, a newer protocol that is generally available starting with Kafka 4.0.
Warning
The consumer protocol replaces the classic rebalance protocol in Kafka clients, depending on the client version:
- Kafka clients 4.3 and later
The consumer protocol is recommended. The classic protocol issues a warning that deprecation will occur in Kafka clients 5.0.
- Kafka clients 5.0 and later
The consumer protocol is the default. The classic protocol remains available but is deprecated.
- Kafka clients 6.0 and later
There is no client-level support for the classic protocol.
For more information, see KIP-1274: Deprecate and remove support for Classic rebalance protocol in KafkaConsumer.
Overview of the rebalance protocols
Confluent Platform supports both the classic and the consumer rebalance protocol, so your consumer clients can use either protocol. The consumer rebalance protocol is preferred.
The new consumer rebalance protocol improves consumer group scalability by removing the group-wide synchronization barrier, making rebalancing truly incremental. Removing this barrier also enhances stability by preserving existing partition assignments as much as possible during rebalances. The new consumer rebalance protocol also reduces rebalance times and simplifies consumers. By offloading or simplifying consumer-side responsibilities, the entire rebalance is streamlined and less prone to delays caused by complex client-side operations.
Every rebalance produces a new set of partition assignments for the group, but the two protocols compute them differently. In the consumer protocol, the coordinator computes assignments from the topic subscriptions and assignment preferences that each member sends. In the classic protocol, the group leader, which is one of the consumers, computes the assignments and the coordinator distributes them to the members.
When a consumer starts up, it finds its group’s coordinator and sends a request to join the group. The coordinator then begins a rebalance. Every rebalance, regardless of which rebalance protocol the group uses, results in a new generation of the group.
Each member in a group must send heartbeats to the coordinator to remain a group member. If the coordinator does not receive a heartbeat from a member before the configured session timeout expires, the coordinator kicks the member out of the group and triggers a rebalance to reassign its partitions to other members.
The following table shows the differences between the classic and consumer rebalance protocols:
Group behavior |
Classic protocol |
Consumer protocol |
|---|---|---|
use case |
Good for architectures that need to continue support for the classic protocol. |
Ideal for dynamic groups (cloud functions or scaling). |
group leader |
One member performs assignment on behalf of the group. |
Not used. No leader election occurs. |
coordinator role |
Assignments computed by the group leader (a consumer) and applied by the broker coordinator. |
Broker coordinator computes assignments but does not receive one global plan. |
rebalance type |
Eager or cooperative, depending on the assignor. |
Incremental reassignment with minimal disruption. |
pause behavior during rebalance |
In some circumstances, all consumers can pause and revoke all partitions. |
Consumers keep their existing partition assignments during a rebalance. |
join group trigger |
Triggered on join or leave. |
A consumer joins a group by sending a heartbeat request. |
join group disruption |
With eager assignors, revokes and reassigns all partitions on every join or leave. With |
Incremental assignment with fewer disruptions. |
assignment |
Handled by one group leader among the consumers. |
Broker coordinator receives consumer subscriptions, validates them, computes, and then orchestrates assignment. |
Enable the consumer rebalance protocol on brokers
In Confluent Platform, you are responsible for enabling the consumer rebalance protocol on your own brokers. Two broker-side settings control availability of the protocol:
The
group.coordinator.rebalance.protocolsbroker configuration must includeconsumerin its list of enabled rebalance protocols on Confluent Platform versions earlier than 8.0. Starting with Confluent Platform version 8.0, this configuration is deprecated, and may be removed in a future release.The
group.versionfeature flag must be set to a version that supports the protocol. On new Confluent Platform 8.0 or later clusters,group.version=1is enabled by default. On clusters upgraded from an earlier release, upgrade the feature explicitly. In the following commands,${CONFLUENT_HOME}is your Confluent Platform installation directory:${CONFLUENT_HOME}/bin/kafka-features --bootstrap-server localhost:9092 upgrade --feature group.version=1
To turn off the feature, downgrade it:
${CONFLUENT_HOME}/bin/kafka-features --bootstrap-server localhost:9092 downgrade --feature group.version=0
For more information about broker-side settings that tune consumer groups using
the consumer rebalance protocol, such as group.consumer.session.timeout.ms,
group.consumer.heartbeat.interval.ms, and group.consumer.assignors, see
Kafka Broker Configurations.
How a group’s rebalance protocol is determined
The first consumer to join a group determines the rebalance protocol for that
group. If this initial consumer uses group.protocol=consumer, either set
explicitly or as the default in Kafka clients 5.0 and later, the group uses the
consumer rebalance protocol. Otherwise, the group uses the classic rebalance
protocol.
Kafka implicitly creates a consumer group when the first consumer with a
specific group.id starts and attempts to subscribe to a topic. You don’t
need an explicit “create group” command or administrative action.
You can, for example, change a group’s rebalance protocol by shutting down all its consumers and bringing them back up. The following is a step-by-step explanation of how a consumer group’s protocol can change from classic to the consumer rebalance protocol:
Initially, in this scenario, the consumer group is operating with the
classicprotocol. The broker’s group coordinator stores metadata about this group, including the protocol it’s currently using (classic).You shut down all the consumer instances belonging to this group. At this point, the group becomes empty from the broker’s perspective. No active members are sending heartbeat requests.
Even though the group is empty, the broker retains the metadata associated with that group, including its last known protocol (
classicin this scenario). The group isn’t entirely “forgotten.”The first consumer instance you bring back up is configured with
group.protocol=consumer. This consumer sends a request to join the group to the coordinator. This request explicitly specifies that it wants to use theconsumerprotocol.Because the group is currently empty, the broker accepts this new member and updates the stored metadata for that consumer group. The broker now records that this group should use the consumer rebalance protocol.
As other consumer instances (also configured with group.protocol=consumer)
are brought back online and join using the same group.id, the
broker will see that the group is now configured to use the consumer
protocol and will allow them to join using that protocol.
A consumer that uses the classic protocol can still join a group that uses the
consumer protocol, which is what makes rolling upgrades possible. The join fails
with an InconsistentGroupProtocolException only if the member uses an
incompatible protocol type, for example, connect.
Upgrading or switching consumer protocols
You can upgrade or switch consumer groups between the
classic and consumer protocols in a few different ways. If you want to migrate or upgrade your
clients to the new consumer rebalance protocol, your choice is between a rolling
upgrade or an empty group restart. When choosing between the two approaches, take
into consideration these aspects of your consumers:
Client version compatibility. Check and ensure your consumer libraries support the Kafka 4.0 version and the
group.protocolconfiguration. For the consumer rebalance protocol to be effective, all consumers in the group should use it.Broker
group.version. The broker’sgroup.versionfeature flag must be set to a version that supports the desired protocol.For the broker-side settings that control the consumer rebalance protocol, see Enable the consumer rebalance protocol on brokers.
Rolling upgrades are the preferred method for production environments because they reduce downtime. If your environment can tolerate some downtime, you can do an empty group restart.
You can also switch back to the classic rebalance protocol from the consumer
rebalance protocol. In this case, do a rolling upgrade, but set
group.protocol=classic explicitly on each consumer. If you choose this
path, make sure you have a good understanding of the potential feature loss and
coordinate with your team.
Note
If a consumer attempts to join a group operating with an incompatible protocol,
the consumer receives an InconsistentGroupProtocolException.
Before migrating to the consumer rebalance protocol
A consumer uses the new consumer rebalance protocol when it meets all the
following conditions:
The consumer uses a client version that supports the
consumerprotocol.The consumer configuration sets
group.protocol=consumer. In Kafka clients 5.0 and later,consumeris the default, so a consumer that doesn’t setgroup.protocolalso uses the consumer rebalance protocol.The cluster has the consumer rebalance protocol enabled. This requires the broker’s
group.versionfeature flag to have a value of1or higher.
Check your consumers and adjust them as necessary before migrating or upgrading to the new consumer rebalance protocol.
Ensure that classic legacy configurations are removed from your consumer properties. Several configuration properties and two APIs from the classic rebalance protocol are no longer applicable in the new consumer rebalance protocol:
The
heartbeat.interval.msproperty on the consumer side is replaced by the server-sidegroup.consumer.heartbeat.interval.msin the consumer rebalance protocol.The
session.timeout.msproperty on the consumer side is replaced by the server-sidegroup.consumer.session.timeout.msin the consumer rebalance protocol.The
partition.assignment.strategyproperty on classic consumer is replaced in the new consumer rebalance protocol by the server-sidegroup.consumer.assignorson the broker and thegroup.remote.assignoron the consumer.The
enforceRebalance(String)andenforceRebalance()APIs are no longer supported with consumers using the new rebalance protocol.
If your consumer uses any legacy properties or methods, they are either ignored
or result in errors if you use them with a consumer where
group.protocol=consumer.
The consumer rebalance protocol also includes subscribe(SubscriptionPattern)
and subscribe(SubscriptionPattern,ConsumerRebalanceListener) methods.
These methods allow consumers to subscribe to a regular expression.
With these methods, the regular expression uses the RE2J
format and is evaluated on the server side. When you subscribe using a
pattern, use an inclusive pattern that matches only the topics you intend to
consume. For example, new SubscriptionPattern("payments-.*") subscribes
to every topic whose name starts with payments-, rather than an overly
broad pattern like .*, which matches every topic. RE2J does not support
lookahead syntax, so exclusion patterns like ^(?!(internal-|_)).* aren’t
valid. A broad pattern can unintentionally match topics created by other
applications or systems, causing your consumer to process unexpected data.
Tip
For information on using the command line to manage your consumer groups, see the Kafka consumer group tool section later in this page.
How to do a rolling deployment
To do a rolling deployment, upgrade consumers to the consumer rebalance protocol one at a time while the group temporarily runs a mix of protocols.
Follow this procedure to do a rolling upgrade:
Ensure that all consumers in the group are using a client version that supports the consumer rebalance protocol.
Remove any classic configurations from the consumer properties.
Deploy new versions of your consumer applications configured with
group.protocol=consumerone at a time.Restart each consumer instance after the configuration change.
When the first consumer that uses the
consumerprotocol joins an existingclassicgroup, the broker converts the group to the consumer rebalance protocol and updates the group’s metadata.You can use the
kafka-consumer-groupscommand line tool to check the migration progress. For information on how to use this tool, see Kafka consumer group tool later in this page.Continue rolling out the new configuration to all consumers in the group.
After all consumers are upgraded, the group fully uses the new protocol.
While Kafka attempts to interoperate between the classic and consumer rebalance protocols during this transition, running a group with mixed protocols temporarily is not a recommended long-term operating model and can lead to potential issues and limitations such as:
The broker needs to maintain compatibility with both protocols, which can introduce overhead and prevent the group from fully leveraging the optimizations of the new protocol.
Managing a group with mixed protocols can increase the complexity of the rebalance process and potentially lead to less predictable behavior, especially in edge cases or during rapid membership changes.
The differences in how the protocols handle assignment could lead to temporary imbalances in partition distribution among consumers using different protocols.
Monitoring and troubleshooting the rebalance process and group health is more complex when different consumers operate under different protocol rules.
How to do an empty group restart
This method is simpler than a rolling upgrade but involves more downtime.
Shut down all consumer instances in the group, making the group empty.
Ensure that all consumers in the group are using a client version that supports the consumer rebalance protocol.
Remove any classic configurations from the consumer properties.
Bring all consumer instances back online with the
group.protocol=consumerconfiguration.The first joining consumer with the new protocol sets the group’s protocol for subsequent members.
Offset management
Offset management is how a consumer tracks and commits its read position
within each partition. After the consumer receives its assignment from the
coordinator, it must determine the initial position for each assigned
partition. If the group has no committed offset for a partition, for example
when the group is first created or after its committed offsets expire, the
position is set according to a configurable offset reset policy
(auto.offset.reset). Typically, consumption starts either at the earliest
offset or the latest offset.
Committing offsets and reset policy
As a consumer reads messages from its assigned partitions, it must commit the
offsets of the messages it has read. If a consumer crashes or shuts down, its
partitions are reassigned to another group member. That member resumes reading
from the last committed offset in each partition. If no offset has been
committed for a partition yet, for example because the group is new and its
consumer crashes before its first commit, the next consumer starts from the
position defined by the auto.offset.reset policy.
When a consumer commits offsets, the group coordinator checks that the commit
comes from a current member of the group. This prevents consumers that have been
fenced or whose partitions were reassigned from overwriting offsets for
partitions they no longer own. With the classic protocol, the coordinator
validates the member ID and generation. With the consumer group protocol
(group.protocol=consumer), the coordinator validates the member epoch.
Auto-commit offsets
The offset commit policy is crucial to providing the message delivery guarantees needed by your application. By default, the consumer is configured to use an automatic commit policy, which triggers a commit on a periodic interval. The consumer also supports a commit API which can be used for manual offset management. Correct offset management is crucial because it affects delivery semantics.
By default, the consumer is configured to auto-commit offsets. The
auto.commit.interval.ms property sets the upper time bound of the
commit interval. Using auto-commit offsets can give you “at-least-once”
delivery, but you must consume all data returned from a ConsumerRecords<K, V>
poll(Duration timeout) call before any subsequent poll calls, or before
closing the consumer.
If the consumer crashes or restarts before the next auto-commit interval, your
application can lose track of which messages it has processed. This happens
because every time your application calls the poll method and the consumer
fetches data, the consumer is ready to automatically commit the offsets of the
messages that the poll returned, whether you have finished processing them.
If the consumer restarts, it begins consuming from the last committed offset.
When this happens, the last committed position can be as old as the auto-commit interval.
Any messages that have arrived after the last commit are read again.
Manual commit API
If you want to reduce the window for duplicates, you can
reduce the auto-commit interval, but some users want even finer
control over offsets. The consumer supports a commit API
which gives you full control over offsets.
Note that when you use the commit API directly, you should first
disable auto-commit in the configuration by setting the
enable.auto.commit property to false.
Each call to the commit API results in an offset commit request being sent to the broker. Using the synchronous API, the consumer is blocked until that request returns successfully. This can reduce overall throughput since the consumer might otherwise be able to process records while that commit is pending.
Tip
To improve throughput, you can increase the amount of data returned when polling.
Set the fetch.min.bytes configuration to a higher value. The broker waits
until that much data is available (or until fetch.max.wait.ms expires) before
responding.
Be aware that this can increase the number of duplicate messages you might need to handle in failure scenarios.
Asynchronous commits
The commit API supports both synchronous and asynchronous modes. A synchronous commit blocks until the broker confirms the offset. A second option is to use asynchronous commits. Instead of waiting for the request to complete, the consumer can send the request and return immediately by using asynchronous commits.
If it helps performance, why not always use asynchronous commits? The main
reason is that the consumer does not retry the request if the commit fails. This
is something that committing synchronously provides; it retries until
the timeout provided by the user, or until the API timeout from
configuration is reached,
depending on which commitSync function was called. The problem with
asynchronous commits is dealing with commit ordering. By the time the consumer
finds out that a commit has failed, you might have already processed the next
batch of messages and even sent the next commit. In this case, a retry of the
old commit could cause duplicate consumption.
Instead of complicating the consumer internals to try to handle this problem, the API gives you a callback which is invoked when the commit either succeeds or fails. If you like, you can use this callback to retry the commit, but you must handle the same reordering problem.
Dealing with commit failures and rebalances
Offset commit failures are merely annoying if the following commits succeed because they won’t actually result in duplicate reads. However, if the last commit fails before a rebalance occurs or before the consumer is shut down, then offsets are reset to the last commit and you can see duplicates. A common pattern is to combine asynchronous commits in the poll loop with sync commits on rebalances or shutdown. Committing on close is straightforward, but you need a way to hook into rebalances.
A rebalance can revoke partitions from a consumer and assign partitions to it.
The revocation callback (onPartitionsRevoked) is called before the consumer
gives up ownership of partitions and is the last chance to commit offsets for
them. With eager assignors, every rebalance revokes all partitions. With
CooperativeStickyAssignor or the consumer rebalance protocol, the callback
is called only for the partitions that move to another consumer. The assignment
callback (onPartitionsAssigned) is called after the rebalance and is used to
set the initial position of newly assigned partitions. To use the preceding
pattern, commit the current offsets synchronously in the revocation callback.
Sync vs async: safety and performance tradeoffs
In general, you should consider asynchronous commits less safe than synchronous commits. Consecutive commit failures before a crash can result in increased duplicate processing. You can mitigate this danger by adding logic to handle commit failures in the callback or by mixing occasional synchronous commits. However, don’t add too much complexity unless testing shows it is necessary. In consumer rebalance mode, where rebalances might be partial and incremental, you might need fewer synchronous commits, especially if most assignments are preserved across rebalances.
If you need more reliability, synchronous commits are a good option, and you can still scale up by increasing the number of topic partitions and the number of consumers in the group. If you want to maximize throughput and you’re willing to accept some increase in the number of duplicates, then asynchronous commits might be a good option.
Asynchronous commits only make sense for “at least once” message delivery. To get “at most once,” you need to know if the commit succeeded before consuming the message. This implies a synchronous commit unless you have the ability to “unread” a message after you find that the commit failed.
For examples of the commit API and a discussion of the performance and reliability tradeoffs, see Kafka Client Examples for Confluent Platform.
Coordinating offset commits with external systems
When writing to an external system, coordinate the consumer’s position with the output it writes. To do this, have your application store the consumer’s offsets in the same system as its output, in the same atomic write, and seek to those stored offsets when the consumer receives its partition assignment. For example, the HDFS sink connector writes the offsets of the data it reads to HDFS along with the data, so each write updates both the data and the offsets, or neither.
A similar pattern is followed for many other data systems that require these stronger semantics, and for which the messages do not have a primary key to allow for deduplication.
Exactly-once processing and transactions
Kafka provides exactly-once processing through the transactional APIs available in Kafka Streams and the core producer/consumer client. These guarantees ensure that messages are neither lost nor duplicated during processing or transfer between topics.
Kafka Streams achieves this by committing both offsets and output results as
part of a single atomic transaction. Similarly, in a consume-transform-produce
application, a transactional producer can commit the consumer’s offsets as part
of the same transaction as its output writes, and consumers that set
isolation.level=read_committed read only committed messages.
By default, Kafka guarantees at-least-once delivery. You can implement at-most-once delivery by disabling retries on the producer and committing offsets in the consumer before processing a batch of messages.
Tip
Consumers can fetch/consume from out-of-sync follower replicas if using a fetch-from-follower configuration. To learn more, see Multi-Region Clusters.
Kafka consumer configuration
The full list of configuration settings is available in Kafka Consumer Configurations. Several of the key configuration settings and how they affect the consumer’s behavior are highlighted below.
Configure core consumer properties
The following core properties apply to all consumers, whether they use the consumer or the classic rebalance protocol:
bootstrap.servers: You are required to set this property so that the consumer can find the Kafka cluster.client.id: Optional, but you should set this property so that you can correlate requests on the broker with the application that made them. Typically, all consumers within the same group share the same client ID so that client quotas apply to the application as a whole.
Group configuration
The following properties apply to consumer groups.
group.id: Optional but you should always configure a group ID unless you are using the simple assignment API and you don’t need to store offsets in Kafka.session.timeout.ms: Applies only to the classic rebalance protocol. Withgroup.protocol=consumer, the broker-sidegroup.consumer.session.timeout.mssetting applies instead. For the classic protocol, control the session timeout by overriding this value. The default is 45 seconds in the C/C++ and Java clients, but you can increase the time to avoid excessive rebalancing, for example due to poor network connectivity or long GC pauses. The main drawback to using a larger session timeout is that it takes longer for the coordinator to detect when a consumer instance has crashed, which means it also takes longer for another consumer in the group to take over its partitions. For normal shutdowns, however, a dynamic member sends an explicit request to the coordinator to leave the group, which triggers an immediate rebalance. A static member (group.instance.id) doesn’t leave the group on shutdown, so its partitions aren’t reassigned until the session timeout expires.heartbeat.interval.ms: Applies only to the classic rebalance protocol. Withgroup.protocol=consumer, the broker-sidegroup.consumer.heartbeat.interval.mssetting applies instead. For the classic protocol, this setting controls how often the consumer sends heartbeats to the coordinator. It is also the way that the consumer detects when a rebalance is needed, so a lower heartbeat interval generally means faster rebalancing. The default setting is three seconds. For larger groups, it might be wise to increase this setting.max.poll.interval.ms: This property specifies the maximum allowed time between calls to the consumer’s poll method (Consumemethod in .NET) before the consumer process is assumed to have failed. The default is 300 seconds and can be safely increased if your application requires more time to process messages. If you are using the Java consumer, you can also adjustmax.poll.recordsto tune the number of records that are handled on every loop iteration.
Offset management configuration
Two main settings control offset management: whether auto-commit is enabled and the offset reset policy.
enable.auto.commit: This setting enables auto-commit (the default), which means the consumer automatically commits offsets periodically at the interval set byauto.commit.interval.ms. The default interval is 5 seconds.auto.offset.reset: Defines the behavior of the consumer when there is no committed position (which occurs when the group is first initialized) or when an offset is out of range. You can reset the position to theearliestoffset or the defaultlatestoffset, or setnoneto throw an exception to the consumer. In Kafka Java clients 4.0 and later, you can also useby_duration:<duration>, where<duration>is an ISO 8601 duration such asPT1H, to reset to the offset at that duration before the current time.
Partition assignment configuration
partition.assignment.strategy sets the partition assignment strategy for a consumer, meaning how partition ownership
is distributed between consumer instances when group management is used. All consumers in the same consumer
group must list at least one assignor in common. The group uses an assignor that every member supports.
The partition.assignment.strategy parameter is not supported when the
consumer rebalance protocol, group.protocol=consumer, is enabled. This
configuration applies only when using the classic rebalance protocol,
group.protocol=classic, which is the default in Kafka clients earlier than 5.0.
Note
The classic rebalance protocol is deprecated starting with Kafka clients 5.0 and is no longer supported in 6.0. For the timeline, see Groups and rebalance protocols.
partition.assignment.strategy accepts a comma-separated list of fully qualified class names that implement the ConsumerPartitionAssignor interface. The list enables you to update the strategy for a group, while temporarily keeping the old one for consumers that have not transitioned to the new strategy yet.
When you configure Kafka consumers, the choice of assignment strategy is important and depends on the specific requirements
for partition balancing, consumer group stability, and rebalance behavior.
In most cases, the default works well. For the Java client, the default is
RangeAssignor,CooperativeStickyAssignor, so the group uses eager range
assignment until every member removes RangeAssignor from the list. For specific
use cases, changing the assignment strategy can affect performance and reliability.
Available options are:
Range Assignment (Default)
How it Works: The
org.apache.kafka.clients.consumer.RangeAssignorworks by evenly distributing partitions of each topic across the consumers in a consumer group. It sorts both the partitions and consumers. Partitions are assigned to consumers in chunks (ranges), aiming for an even distribution.Advantages: It works well when partition count is higher than consumer count, providing a simple and efficient means of partition distribution.
Disadvantages: Can result in uneven load distribution if the number of partitions is not a multiple of the number of consumers.
Round Robin Assignment
How it Works: The
org.apache.kafka.clients.consumer.RoundRobinAssignordistributes partitions across all consumers one by one in a round-robin fashion. It ensures a more even distribution of partitions across consumers, regardless of the number of partitions.Advantages: Leads to a more balanced partition allocation across consumers, useful when handling a varying number of partitions or when partitions have significantly different sizes.
Disadvantages: Like the Range assignor, it’s an eager assignor that doesn’t try to preserve previous assignments, so each rebalance can move more partitions than a sticky assignor would.
Sticky Assignor
How it Works: The
org.apache.kafka.clients.consumer.StickyAssignoraims to maintain a stable partition assignment while still balancing the partitions across consumers. It tries to keep the previously assigned partitions to a consumer as unchanged as possible if a rebalance occurs.Advantages: It minimizes the number of partition reassignments across rebalances, so consumers usually keep the same partitions and any local state associated with them.
Disadvantages: While it offers stability, it might not always result in the most balanced partition distribution if the cluster or consumer group changes frequently.
Cooperative Sticky Assignor
How it Works: An evolution of the StickyAssignor, the
org.apache.kafka.clients.consumer.CooperativeStickyAssignorenables more incremental rebalancing, which can reduce the latency and resources required during the rebalance process.Advantages: It supports more granular changes to the consumer group memberships or to the partitions themselves, making rebalances less impactful. Note that this assignor reduces rebalance impact, not the frequency of rebalances. When all members subscribe to the same set of topics,
CooperativeStickyAssignorautomatically uses an optimized assignment path, so you don’t need extra configuration.Disadvantages: Not all consumers or versions of Kafka support this assignor; requires careful management to ensure compatibility across the consumer group.
Upgrade to CooperativeStickyAssignor
To migrate an existing consumer group from an eager rebalance assignor, for
example, RangeAssignor or StickyAssignor, to cooperative incremental
rebalancing within the classic rebalance protocol, use a two-step rolling
restart to minimize disruption during the transition. KIP-429: Kafka Consumer Incremental Rebalance Protocol
introduced cooperative incremental rebalancing.
After completing step 1 (configuring dual assignors), the consumer group operates
in a mixed-mode period where members that still use an eager assignor continue
to perform full-stop rebalances. After completing step 2 (switching to the cooperative
assignor), rebalances become incremental and avoid full-stop rebalances.
If your consumers use group.protocol=consumer, follow
Upgrading or switching consumer protocols instead.
Before you begin the upgrade, ensure that all consumers in your group support
CooperativeStickyAssignor: the Kafka Java client version 2.4.0 or later, or
librdkafka-based clients version 1.6.0 or later.
Prepare: use a dual list, and keep the current assignor first.
Set partition.assignment.strategy to include both your current assignor and
CooperativeStickyAssignor, with the current assignor listed first. Perform a rolling restart so all members pick up the new list. The coordinator continues to select the first commonly supported assignor, which is the existing one, so there is no behavior change yet.Example properties:
partition.assignment.strategy=org.apache.kafka.clients.consumer.StickyAssignor,org.apache.kafka.clients.consumer.CooperativeStickyAssignor
Fully qualified class names shown here refer to the Java consumer. For other client libraries, configure the same assignor using the client’s configuration API.
Example Java configuration:
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, Arrays.asList(StickyAssignor.class.getName(), CooperativeStickyAssignor.class.getName()));
Switch: make the cooperative assignor the only value.
After all members have been upgraded to client versions that support cooperative rebalancing, update the configuration so that
CooperativeStickyAssignoris the only value, and do another rolling restart of the group. A consumer uses cooperative rebalancing only when every assignor in its list supports it, so leaving the eager assignor in the list keeps the group on eager rebalances. After the restart, the group transitions incrementally with minimal disruption.Example properties:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
Fully qualified class names shown here refer to the Java consumer. For other client libraries, configure the same assignor using the client’s configuration API.
Example Java configuration:
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, Collections.singletonList(CooperativeStickyAssignor.class.getName()));
Verify the switch: after you roll out the change, check consumer application logs for messages that indicate the selected partition assignor, for example, “Using partition assignor class …CooperativeStickyAssignor”.
CooperativeStickyAssignor upgrade notes
All members in the group must advertise a compatible, ordered list of assignors. The coordinator considers only the assignors that every member lists, and selects the one that the most members rank highest among them.
Do not switch to listing only
CooperativeStickyAssignoruntil all members run client versions that support it, which is 2.4.0 or later. Otherwise, the group fails to find a common assignor.Static membership (
group.instance.id) further reduces movement during restarts and pairs well with cooperative rebalancing. In Kubernetes environments, use StatefulSets to provide stable pod identities that can serve as static member IDs.If using manual offset commits (
enable.auto.commit=false), handleRebalanceInProgressExceptionby callingpoll()in the next loop iteration to complete the rebalancing process.For background and the detailed upgrade rationale, see KIP-429: Kafka Consumer Incremental Rebalance Protocol (see the “Compatibility and Upgrade Path” section).
Kafka consumer group tool
Kafka provides the kafka-consumer-groups command-line utility to view and manage consumer groups, which
is also included with Confluent Platform. Find the tool in ${CONFLUENT_HOME}/bin, where
${CONFLUENT_HOME} is your installation directory.
To use this tool with Confluent Cloud, you must install Confluent Platform and pass a client
configuration file that contains your cluster’s security settings and API key
by using the --command-config option.
You can also use the Confluent CLI to complete some of these tasks. For more information, see the Confluent CLI reference.
List consumer groups
You can get a list of the consumer groups in the cluster using the kafka-consumer-groups tool. On
a large cluster, this might take a while because it collects
the list by inspecting each broker in the cluster.
${CONFLUENT_HOME}/bin/kafka-consumer-groups --bootstrap-server host:9092 --list
The output is a list of all consumer groups for the cluster, including groups for internal use. The output might resemble:
_confluent-controlcenter-7-6-0-1
ConfluentTelemetryReporterSampler--4418883999569981189
test-1234
_confluent-controlcenter-7-6-0-lastProduceTimeConsumer
_confluent-controlcenter-7-6-0-1-command
To use the Confluent CLI for this task, see confluent kafka consumer group list.
Describe groups
The kafka-consumer-groups tool can also be used to collect
information on a current group. For example, to see the current
assignments for the test-1234 group, you could use the following command:
${CONFLUENT_HOME}/bin/kafka-consumer-groups --bootstrap-server host:9092 --describe --group test-1234
The output from this command lists the clients, topics, partitions, and more for that group and resembles:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST CLIENT-ID
test-1234 test-metrics 5 258420 258519 99 test-client-1234 /127.0.0.1 test-client
test-1234 test-metrics 10 257002 257097 95 test-client-1234 /127.0.0.1 test-client
test-1234 test-metrics 4 259580 259660 80 test-client-1234 /127.0.0.1 test-client
test-1234 test-metrics 7 254004 254131 127 test-client-1234 /127.0.0.1 test-client
If you run this command while a rebalance is in progress, the command reports an error. Retry the command to see the assignments for all the members in the current generation.
You can use the --members --verbose flags with the command to see if groups have upgraded from classic to the new rebalance protocol:
$ bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group --members --verbose
GROUP CONSUMER-ID HOST CLIENT-ID #PARTITIONS CURRENT-EPOCH CURRENT-ASSIGNMENT TARGET-EPOCH TARGET-ASSIGNMENT UPGRADED
my-group T4tbGKsvT7CtsxVBH5J2QQ /127.0.0.1 consumer-member 3 3 my_topic:0,1;new_topic:0 3 my_topic:0,1;new_topic:0 true
my-group 6mfoIHq4n3BT7n1HCdqihb /127.0.0.1 classic-member 1 2 - 3 new_topic:1 false
The UPGRADED column is showed when a group is in migration. A true value
in this column means a consumer member upgraded to use the new consumer
rebalance protocol. A false value means the member is using the classic
protocol.
To use the Confluent CLI for this task, see confluent kafka consumer group describe.
Reset offsets
You can also use the kafka-consumer-groups tool to reset the consumer offset in scenarios where a consumer is stalled or lagging.
Before you reset offsets, stop all consumers in the group. The tool resets
offsets only when the group has no active members.
You have many options for changing the offset.
For example, you can reset offsets by shifting forward or backward with --shift-by or reset them to the beginning with --to-earliest.
For all the options, see the kafka-consumer-groups tool usage details.
To reset the offsets back by 20 positions, use the following command:
${CONFLUENT_HOME}/bin/kafka-consumer-groups --bootstrap-server host:9092 --group test-1234 --reset-offsets --shift-by -20 --topic test-metrics --execute
The output contains the group, topic, partition, and new offset for each partition:
GROUP TOPIC PARTITION NEW-OFFSET
test-1234 test-metrics 5 258400
test-1234 test-metrics 10 256982
test-1234 test-metrics 4 259560
test-1234 test-metrics 7 253984
Consumer examples
Confluent provides several resources to help you get started with Kafka consumers.
For a tutorial on how to build a Kafka consumer that can read records from Confluent Cloud or Confluent Platform, see How to build your first Apache KafkaConsumer application.
For consumer examples in several different languages, see Get Started, and use the language selector to choose Java, Python, Go, .NET, JavaScript Client, C/C++, REST, Spring Boot, and more. Click Build Consumer in the navigation menu to see example consumer code for the language you chose.
Use Confluent for VS Code to generate a consumer project for you. You can choose from these languages:
Java
Go
Python
For more information, see Confluent for VS Code with Confluent Cloud and Confluent for VS Code with Confluent Platform.