AWS Big Data Blog

Optimize consumer rebalancing on Amazon MSK with next generation protocol

If you run large consumer groups on Apache Kafka and Amazon Managed Streaming for Apache Kafka (Amazon MSK), you’ve likely experienced the pain of slow rebalances: processing stalls across all consumers, “rebalance storms” triggered by routine scaling events, and prolonged recovery times that impact downstream applications. With the classic rebalance protocol, even a single consumer joining or leaving the group forces a global synchronization barrier, pausing every consumer regardless of whether its partition assignments changed.

The KIP-848 consumer protocol, introduced in Apache Kafka 4.0, fundamentally redesigns how consumer group rebalancing works. Also referred to as “the Next Generation Consumer Rebalance Protocol”, KIP-848 shifts coordination logic from the client to the broker-side group coordinator. This supports fully incremental, server-driven rebalancing that significantly improves performance for large consumer groups. You can use the consumer protocol on Amazon MSK on all 4.x Apache Kafka versions on both MSK Standard and Express brokers.

In this post, we explain how the consumer protocol works, how to enable it on Amazon MSK, and how to diagnose and resolve slow rebalancing issues to help improve performance.

The classic protocol compared to the consumer protocol

The classic protocol relied on client-side rebalance logic with a global synchronization barrier. Every rebalance caused all consumers in the group to pause processing simultaneously regardless of whether their partition assignments were changing. This led to “rebalance storms” in large consumer groups where cascading rebalances could take minutes to resolve. The CooperativeStickyAssignor is a client-side partition assignment strategy that supports incremental, cooperative rebalancing. It significantly improves rebalance performance and minimizes disruption to groups during rebalance events. However, it still suffers from bottlenecks as consumer group size and partition count increase. For large workloads, client-side rebalancing behavior can result in longer rebalancing times and require significant client tuning and monitoring during rebalances.

The consumer protocol addresses these limitations by moving all rebalancing logic to the server. The broker handles coordination using a continuous heartbeat mechanism and server-driven reconciliation process. Only affected partitions move during a rebalance, and consumers with unchanged assignments continue processing uninterrupted. This results in faster recovery compared to the classic protocol, and improved scalability as workloads grow.

The following table compares the classic and next generation protocols across key dimensions.

Aspect Classic Protocol Next Generation Protocol (KIP-848)
Rebalance logic Client-side Fully server-driven
Consumer impact Depends on assignor, all consumers pause, or rebalance is limited by group size Only affected consumers impacted, scales effectively as groups grow
Mechanism Client-side algorithm and cross-group coordination Incremental, async reconciliation
Commit processing Paused during rebalance Able to progress during rebalance
Scalability Complex, fragile at scale Resilient, broker-driven
Rebalance storms Common in large groups Eliminated

Server-side configuration

With the consumer protocol, key parameters are now configured on the server rather than the client:

  • group.consumer.heartbeat.interval.ms – Controls the consumer heartbeat interval (server-side).
  • group.consumer.session.timeout.ms – Controls the session timeout (server-side).
  • group.consumer.assignors – Specifies available assignors (uniform and range by default).

In Amazon MSK Express brokers, these configurations are read-only and cannot be modified. In Amazon MSK Standard brokers, these configurations are managed with broker configurations. To update these configurations in Amazon MSK Standard brokers, refer to Update the configuration of an Amazon MSK cluster.

When to use the consumer protocol

The consumer protocol provides the most benefit to workloads with the following requirements:

  • Large consumer groups: Groups with many consumers and partitions see the most significant improvements because of the elimination of global synchronization barriers.
  • High-availability applications: Applications that cannot afford processing interruptions benefit from continuous message processing during rebalances. Financial services, real-time analytics, and fraud detection systems are ideal candidates.
  • Frequently rebalancing environments: Automatic scaling deployments, Kubernetes with frequent pod restarts, or continuous integration and continuous delivery (CI/CD) environments experience significantly less disruption.
  • Dynamic partition scaling: Workloads that regularly add partitions or topics benefit from the incremental, server-driven approach.

Prerequisites

Before you begin, make sure that you have the following:

  • An Amazon MSK cluster running Apache Kafka version 4.0 or later (both MSK Standard and Express brokers are supported).
  • A Kafka client library that supports the KIP-848 consumer protocol (see Step 4 for supported versions).
  • Basic familiarity with Apache Kafka consumer groups and partition assignment.
  • An AWS account with appropriate permissions to manage your MSK cluster.

Enabling the consumer protocol on Amazon MSK

The following steps walk you through verifying your cluster version, configuring your consumer client, removing deprecated configurations, and confirming client library support.

Step 1: Verify cluster version

The consumer protocol requires Apache Kafka 4.0 or later. To use the consumer protocol on Amazon MSK, verify that your cluster is running Apache Kafka version 4.0.x or later. You can verify your cluster’s Apache Kafka version using the AWS Management Console, AWS Command Line Interface (AWS CLI), or AWS SDKs:

aws kafka describe-cluster-v2 --cluster-arn <your-cluster-arn> \
    --query "ClusterInfo.Provisioned.CurrentBrokerSoftwareInfo.KafkaVersion"

If your cluster is running Apache Kafka 4.0.x or later, the consumer protocol is automatically enabled on the server and ready to use. No additional server-side feature flag verification is needed.

Step 2: Configure consumer client

Set group.protocol=consumer in your consumer configuration. The protocol is not enabled by default:

# confluent-kafka-python example
config = {
    'bootstrap.servers': bootstrap_servers,
    'group.id': group_id,
    'group.protocol': 'consumer',  # Required — defaults to 'classic' if omitted
    'auto.offset.reset': 'earliest'
}

The consumer protocol can be changed in-place for existing consumer groups. When you update the group.protocol, perform a rolling restart of your consumers. The broker-side group coordinator automatically handles the upgrade to the consumer protocol and handles classic protocol requests from old clients alongside the upgraded clients.

Step 3: Remove deprecated client configurations

When the consumer protocol is enabled, the following client-side configurations are no longer supported because they are controlled by the brokers:

  • heartbeat.interval.ms.
  • session.timeout.ms.
  • partition.assignment.strategy.

Step 4: Verify client library support

Verify that your Kafka client version supports the consumer protocol:

  • Java clients: Generally available (GA) in Apache Kafka 4.0+.
  • confluent-kafka-python: Version 2.12.0+ (GA support for KIP-848). See the release notes.
  • librdkafka-based clients (Go, .NET, C/C++): Based on librdkafka 2.12.0+.

Note: For other Kafka client libraries, verify your client library’s documentation for group.protocol=consumer support before enabling the next generation protocol. If your client doesn’t support KIP-848, it will continue to use the classic protocol.

Diagnosing slow consumer group rebalancing with the consumer protocol

Even after enabling the consumer protocol, you may encounter situations where consumer group rebalancing takes longer than expected. The following sections help you diagnose and resolve these issues.

Common symptoms

  • Consumer group rebalancing takes longer than expected despite setting group.protocol=consumer.
  • Consumers pause processing during rebalances.
  • Broker logs show “member session expired” or “fenced” messages.
  • Frequent rebalances triggered during rolling deployments or pod restarts.

Step 1: Confirm the consumer protocol is active using broker logs

Before troubleshooting performance, verify which protocol your consumers are actually using. Check broker logs in Amazon CloudWatch Logs Insights. The log patterns differ significantly between protocols.

Consumer protocol expected logs:

Key indicators: “consumer protocol”, “epoch” terminology, “target assignment” with server-side assignor, “fenced” for member removal.

[GroupCoordinator id=X] [GroupId <group-id>] Member <member-id> joins the consumer group using the consumer protocol.
[GroupCoordinator id=X] [GroupId <group-id>] Bumped group epoch to 309 with metadata hash 4064309670987706693.
[GroupCoordinator id=X] [GroupId <group-id>] Computed a new target assignment for epoch 309 with 'uniform' assignor in 0ms.
[GroupCoordinator id=X] [GroupId <group-id>] Member <member-id> fenced from the group because the member session expired.

Classic protocol expected logs:

Key indicators: “PreparingRebalance” state, “old generation” terminology, “Assignment received from leader”.

If you see classic protocol logs, the consumer protocol is not active. Proceed to Step 2 to troubleshoot why.

[GroupCoordinator id=X] Preparing to rebalance group <group-id> in state PreparingRebalance with old generation X
[GroupCoordinator id=X] Stabilized group <group-id> with X members
[GroupCoordinator id=X] Assignment received from leader for group <group-id>

Step 2: Troubleshoot why the consumer protocol is not active

Verify that your client configuration, client library versions, and cluster versions support the consumer protocol, as described in the preceding Step 1 through Step 4.

Step 3: Resolve slow rebalancing when KIP-848 is active

After you verify the consumer protocol is active but rebalancing is still slow, investigate the following causes:

A. Consumer session timeout causing premature member removal

With the consumer protocol, session timeout is server-controlled through group.consumer.session.timeout.ms (default: 45 seconds). The diagnostic path depends on whether you are using static group membership. The following table outlines the diagnostic path and recommended actions for each scenario.

Scenario Symptom Root cause Recommended action
With static group membership (group.instance.id configured) Slow rebalancing when a static member terminates without calling consumer.close() The coordinator waits for the full session timeout before reassigning partitions. This is the most common cause of slow rebalancing in containerized environments. MSK Standard: Implement graceful shutdown to trigger an immediate leave-group request, or increase the session timeout: group.consumer.session.timeout.ms=60000 (default is 45000). MSK Express: This configuration is not editable in Amazon MSK Express clusters. For Amazon MSK Express, optimize your client’s cold starts to allow members to restart within the 45 second consumer session timeout.
Without static group membership Session timeouts expiring during normal operations Your consumer is freezing or becoming unresponsive, which prevents heartbeats from reaching the coordinator.

Investigate long-running message processing, garbage collection pauses, network connectivity issues, or resource exhaustion on the consumer host. Look for this in broker logs:

[GroupCoordinator id=X] [GroupId <group-id>] Member <member-id> has timed out

B. Missing graceful shutdown handling

When consumers terminate without calling consumer.close(), the coordinator waits for the full session timeout before removing the member. This is the most common cause of slow rebalancing in containerized environments.

Resolution: Implement proper SIGTERM handling to trigger an immediate leave-group:

import signal
import sys
from confluent_kafka import Consumer

class GracefulKafkaConsumer:
    def __init__(self, config):
        self.running = True
        self.consumer = Consumer(config)
        signal.signal(signal.SIGTERM, self.shutdown_handler)
        signal.signal(signal.SIGINT, self.shutdown_handler)

    def shutdown_handler(self, signum, frame):
        print(f"Received signal {signum}, initiating graceful shutdown...")
        self.running = False

    def consume_loop(self):
        self.consumer.subscribe(['your-topic'])
        while self.running:
            msg = self.consumer.poll(timeout=1.0)
            if msg is None:
                continue
            # Process message

        print("Closing consumer gracefully...")
        self.consumer.close()  # Sends LeaveGroup — triggers immediate rebalance
        sys.exit(0)

For Kubernetes, verify that terminationGracePeriodSeconds allows time for consumer.close() to complete:

spec:
  terminationGracePeriodSeconds: 60
  containers:
    - name: kafka-consumer

C. Frequent rebalances from unstable consumers

If consumers repeatedly join and leave (out-of-memory (OOM) kills, CrashLoopBackOff, or short-lived tasks), each event triggers a new rebalance epoch.

Resolution: Use static membership by assigning a unique group.instance.id:

config = {
    'bootstrap.servers': bootstrap_servers,
    'group.id': 'my-group',
    'group.protocol': 'consumer',
    'group.instance.id': f'consumer-{unique_identifier}'  # Unique per consumer
}

With static membership:

  • Short restarts within the session timeout don’t trigger rebalances.
  • The consumer rejoins with the same partition assignment.
  • Scaling up (adding new consumers) still works. New group.instance.id values trigger assignment of unassigned partitions only.

Monitoring and validation

After applying changes, confirm the improvement:

  • Check broker logs: Confirm that “member session expired” messages no longer appear during normal operations or deployments.
  • Monitor consumer lag: Use the SumOffsetLag and EstimatedMaxTimeLag Amazon CloudWatch metrics to verify that lag returns to zero quickly after a rebalance.
  • Describe consumer group: Use kafka-consumer-groups.sh --describe to verify that all members are active and stable.

Conclusion

After implementing the consumer protocol, you should observe the following behavior for consumer group rebalances:

  • Consistently faster rebalance times compared to the classic protocol.
  • Fewer session timeout-related rebalances.
  • More stable consumer group membership.
  • Smooth scaling operations without disrupting existing consumers.
  • Fewer unnecessary rebalances during consumer restarts when using static membership.
  • Clean consumer departures without waiting for timeout expiration when using graceful shutdown.

To get started, try the consumer protocol in your non-production workloads and observe the rebalance improvements as you scale your workload up and down.

To learn more about Amazon MSK and the consumer rebalance protocol, see the following resources:

 


About the authors

Yashika Jain

Yashika Jain

Yashika is a Senior Cloud Analytics Engineer at AWS, specializing in real-time analytics and event-driven architectures. She is committed to helping customers by providing deep technical guidance, driving best practices across real-time data platforms and solving complex issues related to their streaming data architectures.

Vinayaka Gangadhar

Vinayaka Gangadhar

Vinayaka is an Analytics Specialist at Amazon Web Services (AWS), where he helps customers build and troubleshoot scalable data platforms and derive meaningful insights through AWS analytics services, with deep expertise in Amazon Redshift and Amazon OpenSearch. When not solving complex analytics challenges, he enjoys exploring new technologies and spending quality time with his family.

Kalyan Janaki

Kalyan Janaki

Kalyan is Senior Big Data & Analytics Specialist with Amazon Web Services. He helps customers architect and build highly scalable, performant, and secure cloud-based solutions on AWS.