Versions Compared

Key

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

Table of Contents

Status

Current state: Accepted

Under Discussion thread: here

Discussion Vote thread: here

JIRA: here

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

...

ConsumerGroupMetadata class should be an sealed interface, we should provide a new internal permitted class for interface implementation and use java default scope for class visibility.

Code Block
languagejava
titleConsumerGroupMetadata
public interface ConsumerGroupMetadata {

    String groupId();

    int generationIdmemberEpoch();

    String memberId();

    Optional<String> groupInstanceId();
}

...

Code Block
languagejava
titleDefaultConsumerGroupMetadata
class DefaultConsumerGroupMetadata implements ConsumerGroupMetadata {
    private final String groupId;
    private final int generationIdmemberEpoch;
    private final String memberId;
    private final Optional<String> groupInstanceId;

    DefaultConsumerGroupMetadata(String groupId,
                                        int generationIdmemberEpoch,
                                        String memberId,
                                        Optional<String> groupInstanceId) {
        this.groupId = Objects.requireNonNull(groupId, "group.id can't be null");
        this.generationId memberEpoch= generationIdmemberEpoch;
        this.memberId = Objects.requireNonNull(memberId, "member.id can't be null");
        this.groupInstanceId = Objects.requireNonNull(groupInstanceId, "group.instance.id can't be null");
    }

    public String groupId() {
        return groupId;
    }

    public int generationIdmemberEpoch() {
        return generationIdmemberEpoch;
    }

    public String memberId() {
        return memberId;
    }

    public Optional<String> groupInstanceId() {
        return groupInstanceId;
    }

    @Override
    public String toString() {
        return String.format("GroupMetadata(groupId = %s, generationId memberEpoch= %d, memberId = %s, groupInstanceId = %s)",
                groupId,
            generationId    memberEpoch,
                memberId,
                groupInstanceId.orElse(""));
    }

    @Override
    public boolean equals(final Object o) {
        if (this == o) return true;
        if (o == null || getClass() != o.getClass()) return false;
        final DefaultConsumerGroupMetadata that = (DefaultConsumerGroupMetadata) o;
        return generationId memberEpoch== that.generationIdmemberEpoch &&
                Objects.equals(groupId, that.groupId) &&
                Objects.equals(memberId, that.memberId) &&
                Objects.equals(groupInstanceId, that.groupInstanceId);
    }

    @Override
    public int hashCode() {
        return Objects.hash(groupId, generationIdmemberEpoch, memberId, groupInstanceId);
    }

}

 

Compatibility, Deprecation, and Migration Plan

  • Change is impacting all users using any available constructor of ConsumerGroupMetadata class.
  • For Kafka 4.2.10, we should keep the the ConsumerGroupMetadata class as-is for backward compatibility, but deprecate all constructors and update the JavaDocs to indicate this class will become an interface in the future.
  • For Kafka 5.0, we do the actual code change from class to interface.
  • Once Kafka clients stop supporting Java 11 we can mark the ConsumerGroupMetadata interface as sealed with only default implementation permitted.

Test Plan

Automated tests which are part of CI will cover necessary tests.

...

We can make all constructors in ConsumerGroupMetadata private and provide new static factory method which will be getting instance of group metadata from Kafka Consumer.

This approach is rejected because internal Kafka Consumer components (ConsumerCoordinat and AsyncKafkaConsumer) are coupled with ConsumerGroupMetadata class and decoupling them requires big code refactoring.

We can highlight two major hotspots:

1) ConsumerCoordinator is evaluating the onAssignment method from ConsumerPartitionAssignor which is a public interface - it means that even when we try to use some alternative internal implementation of ConsumerGroupMetadata in ConsumerCoordinator, it will fail at that point.

2) AsyncKafkaConsumer implements the public org.apache.kafka.clients.consumer.Consumer interface, so for implementing the ConsumerGroupMetadata groupMetadata() method, it has to be able to instantiate a ConsumerGroupMetadata object. Unfortunately, AsyncKafkaConsumer is in org.apache.kafka.clients.consumer.internals, so as you mentioned, package-private scope will not fly.


Code Block
languagejava
titlealternative
publicclass DefaultConsumerGroupMetadata classimplements ConsumerGroupMetadata {
    private final String groupId;
    private final int generationIdmemberEpoch;
    private final String memberId;
    private final Optional<String> groupInstanceId;

    public ConsumerGroupMetadata newInstance(KafkaConsumer<?,?> kafkaConsumer) {
        return kafkaConsumer.groupMetadata();
    }
    
    private ConsumerGroupMetadata(DefaultConsumerGroupMetadata(String groupId,
                                 int generationIdmemberEpoch,
                                 String memberId,
                                 Optional<String> groupInstanceId) {
        this.groupId = Objects.requireNonNull(groupId, "group.id can't be null");
        this.generationId memberEpoch= generationIdmemberEpoch;
        this.memberId = Objects.requireNonNull(memberId, "member.id can't be null");
        this.groupInstanceId = Objects.requireNonNull(groupInstanceId, "group.instance.id can't be null");
    }

    private ConsumerGroupMetadata(String groupId) {
        this(groupId,
            JoinGroupRequest.UNKNOWN_GENERATION_ID,
            JoinGroupRequest.UNKNOWN_MEMBER_ID,
            Optional.empty());
    }

    public String groupId() {
        return groupId;
    }

    public int generationIdmemberEpoch() {
        return generationIdmemberEpoch;
    }

    public String memberId() {
        return memberId;
    }

    public Optional<String> groupInstanceId() {
        return groupInstanceId;
    }

    @Override
    public String toString() {
        return String.format("GroupMetadata(groupId = %s, generationId memberEpoch= %d, memberId = %s, groupInstanceId = %s)",
                groupId,
               generationId memberEpoch,
                memberId,
                groupInstanceId.orElse(""));
    }

    @Override
    public boolean equals(final Object o) {
        if (this == o) return true;
        if (o == null || getClass() != o.getClass()) return false;
        final ConsumerGroupMetadataDefaultConsumerGroupMetadata that = (ConsumerGroupMetadataDefaultConsumerGroupMetadata) o;
        return generationId memberEpoch== that.generationIdmemberEpoch &&
                Objects.equals(groupId, that.groupId) &&
                Objects.equals(memberId, that.memberId) &&
                Objects.equals(groupInstanceId, that.groupInstanceId);
    }

    @Override
    public int hashCode() {
        return Objects.hash(groupId, generationIdmemberEpoch, memberId, groupInstanceId);
    }

}