Current state: Under Discussion
Discussion thread: here
JIRA: here
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Since Kafka transactions protocol has been improved throught KIP-447: Producer scalability for exactly once semantics and Kafka Streams has already adopted eos-v2 by KIP-732: Deprecate eos-alpha and replace eos-beta with eos-v2 we should now drop support for the legacy transactions protocol (which doesn't use Consumer Group generation id and member id) in Kafka Producer.
Present API informs developers that they should avoid legacy protocol by marking void sendOffsetsToTransaction(Map<TopicPartition,OffsetAndMetadata> offsets, String consumerGroupId) deprecated, but still developers can easily use a non depreacated method void sendOffsetsToTransaction(Map<TopicPartition,OffsetAndMetadata> offsets, ConsumerGroupMetadata groupMetadata) and pass there ConsumerGroupMetadata instance created by hand using non deprecated constructor.
This situation creates confustion when migrating the legacy applications, as developers might simply wrap a string consumer group id in a ConsumerGroupMetadata object and pass it to the void sendOffsetsToTransaction(Map<TopicPartition,OffsetAndMetadata> offsets, ConsumerGroupMetadata groupMetadata) to avoid using deprecated APIs - but this will still lacks the necessary generation id and member id.
Users should not create instances of ConsumerGroupMetadata, and this class should not have any public constructor.
ConsumerGroupMetadata class should be an interface, we should provide a new internal class for interface implementation and use java default scope for class visibility.
public interface ConsumerGroupMetadata {
String groupId();
int generationId();
String memberId();
Optional<String> groupInstanceId();
} |
class DefaultConsumerGroupMetadata implements ConsumerGroupMetadata {
private final String groupId;
private final int generationId;
private final String memberId;
private final Optional<String> groupInstanceId;
DefaultConsumerGroupMetadata(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");
}
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 DefaultConsumerGroupMetadata that = (DefaultConsumerGroupMetadata) 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);
}
}
|
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.
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);
}
}
|