Shekharrajak commented on code in PR #22357:
URL: https://github.com/apache/kafka/pull/22357#discussion_r3469553421
##########
clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/consumer/ShareConsumerTransactionTest.java:
##########
@@ -129,50 +129,66 @@ public void testTransactionalShareAckCommitAccept()
throws Exception {
}
@ClusterTest
- public void testTransactionalShareAckRejectsStaleMemberEpoch() throws
Exception {
- String groupId = "txn-share-stale-member-epoch";
+ public void testTransactionalShareAckRejectsFencedMemberEpochs() throws
Exception {
+
assertTransactionalShareAckRejectedForInvalidMemberEpoch("future-member-epoch",
1);
+
assertTransactionalShareAckRejectedForInvalidMemberEpoch("stale-member-epoch",
-2);
+ }
+
+ private void assertTransactionalShareAckRejectedForInvalidMemberEpoch(
+ String scenario,
+ int memberEpochDelta
+ ) throws Exception {
+ String groupId = "txn-share-" + scenario;
+ String topic = "txn-share-" + scenario;
+ TopicPartition topicPartition = new TopicPartition(topic, 0);
+ createTopic(topic);
alterShareAutoOffsetReset(groupId, "earliest");
try (Producer<byte[], byte[]> producer = createProducer();
- Producer<byte[], byte[]> transactionalProducer =
createTransactionalProducer("txn-share-stale-member-epoch-producer");
+ Producer<byte[], byte[]> transactionalProducer =
createTransactionalProducer(groupId + "-producer");
ShareConsumer<byte[], byte[]> shareConsumer = createShareConsumer(
groupId,
Map.of(ConsumerConfig.SHARE_ACKNOWLEDGEMENT_MODE_CONFIG,
EXPLICIT));
Admin admin = createAdminClient()) {
- producer.send(record("first")).get();
+ producer.send(record(topicPartition, scenario)).get();
producer.flush();
- shareConsumer.subscribe(Set.of(tp.topic()));
+ shareConsumer.subscribe(Set.of(topic));
ConsumerRecords<byte[], byte[]> records =
waitedPoll(shareConsumer, 2500L, 1);
ConsumerRecord<byte[], byte[]> record = records.iterator().next();
assertEquals(0L, record.offset());
- assertEquals("first", new String(record.value(),
StandardCharsets.UTF_8));
+ assertEquals(scenario, new String(record.value(),
StandardCharsets.UTF_8));
ShareGroupMetadata groupMetadata =
shareConsumer.shareGroupMetadata();
- ShareGroupMetadata staleGroupMetadata = new ShareGroupMetadata(
+ ShareGroupMetadata fencedGroupMetadata = new ShareGroupMetadata(
groupMetadata.groupId(),
groupMetadata.memberId(),
- groupMetadata.memberEpoch() + 1
+ invalidMemberEpoch(groupMetadata.memberEpoch(),
memberEpochDelta)
);
shareConsumer.acknowledge(record, AcknowledgeType.ACCEPT);
ShareAcknowledgements acknowledgements =
shareConsumer.acknowledgementsForTransaction();
assertFalse(acknowledgements.isEmpty());
transactionalProducer.initTransactions();
- transactionalProducer.partitionsFor(tp.topic());
+ transactionalProducer.partitionsFor(topic);
transactionalProducer.beginTransaction();
try {
assertThrows(CommitFailedException.class,
- () ->
transactionalProducer.sendShareAcknowledgementsToTransaction(acknowledgements,
staleGroupMetadata));
+ () ->
transactionalProducer.sendShareAcknowledgementsToTransaction(acknowledgements,
fencedGroupMetadata));
Review Comment:
the request is rejected because its epoch does not match coordinator-owned
membership state.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]