DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
...
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 | ||||
|---|---|---|---|---|
| ||||
public class ConsumerGroupMetadata {
private final String groupId;
private final int generationId;
private final String memberId;
private final Optional<String> groupInstanceId;
public ConsumerGroupMetadata newInstance(KafkaConsumer<?,?> kafkaConsumer) {
return kafkaConsumer.groupMetadata();
}
private ConsumerGroupMetadata(String groupId,
int generationId,
String memberId,
Optional<String> groupInstanceId) {
this.groupId = Objects.requireNonNull(groupId, "group.id can't be null");
this.generationId = generationId;
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 generationId() {
return generationId;
}
public String memberId() {
return memberId;
}
public Optional<String> groupInstanceId() {
return groupInstanceId;
}
@Override
public String toString() {
return String.format("GroupMetadata(groupId = %s, generationId = %d, memberId = %s, groupInstanceId = %s)",
groupId,
generationId,
memberId,
groupInstanceId.orElse(""));
}
@Override
public boolean equals(final Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
final ConsumerGroupMetadata that = (ConsumerGroupMetadata) o;
return generationId == that.generationId &&
Objects.equals(groupId, that.groupId) &&
Objects.equals(memberId, that.memberId) &&
Objects.equals(groupInstanceId, that.groupInstanceId);
}
@Override
public int hashCode() {
return Objects.hash(groupId, generationId, memberId, groupInstanceId);
}
}
|
...