Skip to content

Commit 33ea66b

Browse files
KAFKA-20736: Improve share session handling on leader change (apache#22766)
In some situations, the share consumer does not tidy up share sessions sufficiently on leadership change. In particular, if a broker loses leadership of a partition, the share consumer should send a further ShareFetch/ShareAcknowledge to remove the partition from its share session. In most cases, this removal happens naturally, but it depends upon whether the share consumer has other requests to send to that broker. This PR tidies up the maintenance of share sessions when leadership changes occur. Reviewers: Lianet Magrans <lmagrans@confluent.io>, Shivsundar R <shr@confluent.io>, Apoorv Mittal <apoorvmittal10@gmail.com>
1 parent 260a397 commit 33ea66b

5 files changed

Lines changed: 710 additions & 166 deletions

File tree

clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/consumer/ShareConsumerTest.java

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -33,10 +33,9 @@
3333
import org.apache.kafka.common.errors.GroupMaxSizeReachedException;
3434
import org.apache.kafka.common.errors.InterruptException;
3535
import org.apache.kafka.common.errors.InvalidTopicException;
36-
import org.apache.kafka.common.errors.NotLeaderOrFollowerException;
36+
import org.apache.kafka.common.errors.NetworkException;
3737
import org.apache.kafka.common.errors.RecordDeserializationException;
3838
import org.apache.kafka.common.errors.SerializationException;
39-
import org.apache.kafka.common.errors.ShareSessionNotFoundException;
4039
import org.apache.kafka.common.errors.UnknownTopicIdException;
4140
import org.apache.kafka.common.errors.WakeupException;
4241
import org.apache.kafka.common.header.Header;
@@ -1054,7 +1053,7 @@ public void testLeaderRestartWithoutLeadershipChangeExplicitAcknowledgementSync(
10541053

10551054
AtomicBoolean callbackCalled = new AtomicBoolean(false);
10561055
shareConsumer.setAcknowledgementCommitCallback((offsetsByTopicPartition, exception) -> {
1057-
assertInstanceOf(NotLeaderOrFollowerException.class, exception);
1056+
assertInstanceOf(NetworkException.class, exception);
10581057
callbackCalled.set(true);
10591058
});
10601059

@@ -1083,7 +1082,7 @@ public void testLeaderRestartWithoutLeadershipChangeExplicitAcknowledgementSync(
10831082
assertEquals(1, commitResult.size());
10841083
TopicIdPartition tidp = commitResult.keySet().iterator().next();
10851084
assertTrue(commitResult.get(tidp).isPresent());
1086-
assertInstanceOf(NotLeaderOrFollowerException.class, commitResult.get(tidp).get());
1085+
assertInstanceOf(NetworkException.class, commitResult.get(tidp).get());
10871086

10881087
assertTrue(callbackCalled.get());
10891088
}
@@ -1098,7 +1097,7 @@ public void testLeaderRestartWithoutLeadershipChangeExplicitAcknowledgementAsync
10981097

10991098
AtomicBoolean callbackCalled = new AtomicBoolean(false);
11001099
shareConsumer.setAcknowledgementCommitCallback((offsetsByTopicPartition, exception) -> {
1101-
assertInstanceOf(NotLeaderOrFollowerException.class, exception);
1100+
assertInstanceOf(NetworkException.class, exception);
11021101
callbackCalled.set(true);
11031102
});
11041103

@@ -1147,7 +1146,7 @@ public void testLeaderRestartWithoutLeadershipChangeImplicitAcknowledgement() {
11471146

11481147
AtomicBoolean callbackCalled = new AtomicBoolean(false);
11491148
shareConsumer.setAcknowledgementCommitCallback((offsetsByTopicPartition, exception) -> {
1150-
assertInstanceOf(ShareSessionNotFoundException.class, exception);
1149+
assertInstanceOf(NetworkException.class, exception);
11511150
callbackCalled.set(true);
11521151
});
11531152

0 commit comments

Comments
 (0)