Versions Compared

Key

  • This line was added.
  • This line was removed.
  • Formatting was changed.

Table of Contents

Status

Current state: "Vote in progress"  Accepted

Discussion thread: here

Vote thread:  here

JIRA: KAFKA-19426

Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).

Motivation

TopicBasedRemoteLogMetadataManager(TBRLMM) is the key built-in implementation for tiered-storage feature. It uses the Kafka topic to maintain the metadata for remote storage.

Thus. when we begin to support tiered-storage with it. We found that the TBRLMM's initialization failed in some Kafka clusters sometimes.

The root cause is the initialization(TopicBasedRemoteLogMetadataManager#initializeResources) will use the __remote_log_metadata  but the server is not ready for handle the request.

We can check with broker's startup sequence refer to follow snapshot:

Image Removed

So the initialization will have to relay on the retry to get success until the broker is ready. The retry many happen with multiple times:

FYI:  according to one of kamalcph 's test result: 
Image Removed

Thus. the critical bad case is that the default retry time (DEFAULT_REMOTE_LOG_METADATA_INITIALIZATION_RETRY_MAX_TIMEOUT_MS: 2 Minutes) is not enough for some Kafka cluster which take > 2 minutes to complete the startup. 

Base on this. the fail will happen. As one result. The fail will cause the feature broken and local disk never get deleted and some other issues for different cases.

Image Removed

As one workaround solution. we can check every kafka cluster's  startup consume time and set a very big value for the DEFAULT_REMOTE_LOG_METADATA_INITIALIZATION_RETRY_MAX_TIMEOUT_MS due to the startup time may increased in future.

But It is not reasonable for one configure need to change again and again or set to a very large value . What's more, as we mentioned above you will find there are lots of warn log which hint the connection not available. This can be avoided by this change.

...

[2025-07-19 18:00:17,923] WARN [AdminClient clientId=adminclient-1] Connection to node -1 (10.20.4.98:9559) could not be established. Node may not be available. (org.apache.kafka.clients.NetworkClient)

19784

PR: https://github.com/apache/kafka/pull/20691

Motivation

Apache Kafka supports rack-aware partition assignment to improve fault tolerance and reduce cross-rack network traffic/cost. However, the current Admin API does not expose rack information for all type of consumer group members, despite this information being available at the protocol level.

Current State


The ConsumerGroupDescribeResponse (API Key 69) protocol includes a rackId field for each group member, as defined in the protocol specification:

{ "name": "RackId", "type": "string", "versions": "0+", 
  "nullableVersions": "0+", "default": "null",
  "about": "The member rack ID." }


However, when users call AdminClient.describeConsumerGroups(), the returned MemberDescription objects do not include this rack information. The rack ID is available in the wire protocol but is discarded during response processing in DescribeConsumerGroupsHandler.

The Similar case is for AdminClient.describeShareGroups().
The current status can be summarized as follows:


Class

Type of Group

contain RackId? 

protocoal response contain rackId?

MemberDescription

Consumer Group

✅ 

ShareMemberDescription

Share Group

❌ 

✅ 

StreamsGroupMemberDescription

Streams Group

✅ 

✅ 

Problem Statement


This limitation creates several issues:

  1. Monitoring and Observability: Operators cannot determine the rack distribution of consumer group members through the Admin API, making it difficult to verify that rack-aware assignment is working correctly.
  2. Inconsistency: Other group types expose rack information:

    • StreamsGroupMemberDescription includes rackId() method
    • The underlying protocol supports it for consumer groups
    • Only the public Admin API omits this information
  3. Diagnostics: When troubleshooting rack-aware assignment issues or network problems, operators need to use lower-level tools or custom code to access rack information that should be readily available.
  4. Third-party Tools: Monitoring and management tools built on the Admin API cannot display rack information, limiting their usefulness for rack-aware deployments.

Take one example:  currently we have to implemen our AZ/Rack analysis using a workaround — passing the rack information into the clientId field and parsing it afterward.

kafkaConsumerConfig.customConfig(ConsumerConfig.CLIENT_ID_CONFIG, generateClientIdWithRack(ip, rack));

Image Added


Public Interfaces

So create this KIP to solve this issue.

Public Interfaces

...

We need to do tiny change for the existed public class:
1. Add rack ID support to
org.apache.kafka.clients.admin.MemberDescription

:

Code Block
languagejava
titleRemoteLogMetadataManagerMemberDescription
public class MemberDescription {
    private final String memberId;
    private final Optional<String> groupInstanceId;
    private final Optional<String> rackId;  // new field
    private final String clientId;/**
 * Simple callback interface for broker ready notification.
 */
public interface BrokerReadyCallback {
    /**
     * This method will be called during broker startup for the implementation
     * which needs delayed initialization until the broker can process requests.
     */
    void onBrokerReady();

}

add the interface into RemoteLogMetadataManager's implements with default implement.

/omit other codes        
    public Optional<String> rackId() {  // new method
       return rackId;
    }
    //omit other codes}

2.  Add rack ID support to org.apache.kafka.clients.admin.ShareMemberDescription:

Code Block
languagejava
title ShareMemberDescription
public class ShareMemberDescription{     
    private final String memberId;
    private final Optional<String> rackId; // new field
    private final String clientId;
    //omit other codes        
    public Optional<String> rackId() {  // new method
       return rackId;
    }
    //omit other codes
}
Code Block
public interface RemoteLogMetadataManager extends BrokerReadyCallback, Configurable, Closeable 

Proposed Changes

 You can refer to https://github.com/apache/kafka/pull/20203/files

 We postpone the TopicBasedRemoteLogMetadataManager's initialization part (quering metedata from remote topic) after the server is ready for the request. 

...

20691:

  •  add rackId into MemberDescription and ShareMemberDescription
  •  pass through rack ID from protocol response

Compatibility, Deprecation, and Migration Plan 

...


Backward Compatibility

This change is fully backward compatible:

  1. Binary Compatibility:

    • The old constructor remains available (marked as @Deprecated)
    • Existing code will continue to compile and run
    • The new field is an Optional, defaulting to empty()
  2. Source Compatibility:

    • Existing code using the old constructor continues to work
    • No changes required to existing applications
  3. Behavioral Compatibility:

    • Existing behavior is unchanged
    • Only adds new information when available

Migration Path

For users upgrading:

  • No action required
  • Rack information automatically available when using Admin API
  • Access via memberDescription.rackId() when needed

For developers using MemberDescription:

  • Old constructor still works
  • Recommended to migrate to new constructor over time
  • Old constructor may be removed in a future major version (e.g., Kafka 5.0)

Deprecation Plan

  • Mark old constructor as @Deprecated immediately
  • Keep it available for at least 2 major releases and then remove them.

Test Plan

We can use follow test to cover the change:

Unit Tests

  • Test equality with and without rack ID
  • Test rack ID is correctly extracted from DescribeResponse

Integration/System Tests

  • Behavior with mixed rack configurations
  • Behavior with classic and new protocol groups

Rejected Alternatives

Alternative 1: Create a New API
Proposal: Create describeConsumerGroupsWithRack() method
Rejection Reason:

  • Unnecessary API proliferation
  • The information already exists in the protocol
  • Extending existing API is more intuitive
  • Similar precedent: KIP-345 added groupInstanceId to existing MemberDescription

Test Plan

Describe in few sentences how the KIP will be tested. We are mostly interested in system tests (since unit-tests are specific to implementation details). How will we know that the implementation works as expected? How will we know nothing broke?

  • Will cover the patch with deploy test and check if the startup can success without any connection retry error for remote storage. You can refer to test case

Rejected Alternatives

If there are alternative ways of accomplishing the same thing, what were they? The purpose of this section is to motivate why the design is the way it is and not some other way.

  • Another discussed approach is not to call RLMM#configure() method while instantiating the RemoteLogManager#L422 and define a new method in RemoteLogManager#configureRLMM and this can be called from the BrokerServer. You can refer to the code.
    But the change breaks that contract:
    Accroding to KIP-877: "If a plugin implements this interface, the withPluginMetrics() method will be called when the plugin is instantiated (after configure() if the plugin also implements Configurable). "