This is an automated email from the ASF dual-hosted git repository.
AndrewJSchofield pushed a commit to branch 4.3
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.3 by this push:
new e61da6429c1 KAFKA-20694: Fix share read leader epoch sentinel (#22581)
e61da6429c1 is described below
commit e61da6429c1509df00de86f5990d240d31143101
Author: Shekhar Prasad Rajak <[email protected]>
AuthorDate: Tue Jun 16 14:39:56 2026 +0530
KAFKA-20694: Fix share read leader epoch sentinel (#22581)
https://issues.apache.org/jira/browse/KAFKA-20694
Fixes ReadShareGroupStateRequest handling for leaderEpoch = -1. This
value is used as a sentinel for “do not update leader epoch”, so share
coordinator read validation should not fence it as stale. The read path
now matches write-state validation behavior.
Reviewers: Sushant Mahajan <[email protected]>, Andrew Schofield
<[email protected]>
---
.../org/apache/kafka/coordinator/share/ShareCoordinatorShard.java | 2 +-
.../org/apache/kafka/coordinator/share/ShareCoordinatorShardTest.java | 4 ++++
2 files changed, 5 insertions(+), 1 deletion(-)
diff --git
a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java
b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java
index 9eb4ea3f0da..07f7ce9d825 100644
---
a/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java
+++
b/share-coordinator/src/main/java/org/apache/kafka/coordinator/share/ShareCoordinatorShard.java
@@ -812,7 +812,7 @@ public class ShareCoordinatorShard implements
CoordinatorShard<CoordinatorRecord
topicId, partitionId, Errors.INVALID_REQUEST,
READ_UNINITIALIZED_SHARE_PARTITION.getMessage()));
}
- if (leaderEpochMap.containsKey(mapKey) && leaderEpochMap.get(mapKey) >
partitionData.leaderEpoch()) {
+ if (partitionData.leaderEpoch() != -1 &&
leaderEpochMap.containsKey(mapKey) && leaderEpochMap.get(mapKey) >
partitionData.leaderEpoch()) {
log.error("Read request leader epoch is smaller than last recorded
current: {}, requested: {}.", leaderEpochMap.get(mapKey),
partitionData.leaderEpoch());
return
Optional.of(ReadShareGroupStateResponse.toErrorResponseData(topicId,
partitionId, Errors.FENCED_LEADER_EPOCH, Errors.FENCED_LEADER_EPOCH.message()));
}
diff --git
a/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorShardTest.java
b/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorShardTest.java
index ecdbb991554..4787e34eb0e 100644
---
a/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorShardTest.java
+++
b/share-coordinator/src/test/java/org/apache/kafka/coordinator/share/ShareCoordinatorShardTest.java
@@ -1006,6 +1006,8 @@ class ShareCoordinatorShardTest {
CoordinatorResult<ReadShareGroupStateResponseData, CoordinatorRecord>
result2 = shard.readStateAndMaybeUpdateLeaderEpoch(request2);
assertTrue(result2.records().isEmpty()); // Leader epoch -1 - no
update.
+ assertEquals(Errors.NONE.code(),
result2.response().results().get(0).partitions().get(0).errorCode());
+ assertEquals(2, shard.getLeaderMapValue(SHARE_PARTITION_KEY));
ReadShareGroupStateRequestData request3 = new
ReadShareGroupStateRequestData()
.setGroupId(GROUP_ID)
@@ -1019,6 +1021,8 @@ class ShareCoordinatorShardTest {
CoordinatorResult<ReadShareGroupStateResponseData, CoordinatorRecord>
result3 = shard.readStateAndMaybeUpdateLeaderEpoch(request3);
assertTrue(result3.records().isEmpty()); // Same leader epoch - no
update.
+ assertEquals(Errors.NONE.code(),
result3.response().results().get(0).partitions().get(0).errorCode());
+ assertEquals(2, shard.getLeaderMapValue(SHARE_PARTITION_KEY));
verify(shard.getMetricsShard()).record(ShareCoordinatorMetrics.SHARE_COORDINATOR_WRITE_SENSOR_NAME);
}