This is an automated email from the ASF dual-hosted git repository.

AndrewJSchofield pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 05b3d2adcd1 KAFKA-20736: Improve removal of unassigned partitions from 
share sessions (#22720)
05b3d2adcd1 is described below

commit 05b3d2adcd19e293ea28fd326876a27dad87c816
Author: Andrew Schofield <[email protected]>
AuthorDate: Fri Jul 3 09:29:21 2026 +0100

    KAFKA-20736: Improve removal of unassigned partitions from share sessions 
(#22720)
    
    When all partitions in a share session are unassigned with no pending
    acknowledgements to be sent, no ShareFetch request was being sent. This
    patch sends a ShareFetch to remove the partitions from the share session
    in this situation.
    
    Reviewers: Shivsundar R <[email protected]>, Apoorv Mittal
     <[email protected]>
---
 .../internals/ShareConsumeRequestManager.java      | 21 ++++++
 .../internals/ShareConsumeRequestManagerTest.java  | 61 ++++++++++++++---
 .../server/ShareFetchAcknowledgeRequestTest.scala  | 76 ++++++++++++++++++++--
 3 files changed, 144 insertions(+), 14 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
index 4ff391ecfc1..ac3f891288b 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManager.java
@@ -248,6 +248,27 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
             }
         });
 
+        // Iterate over the session handlers again to remove any partitions 
which are no longer
+        // being fetched and which also have no acknowledgements to send. This 
handles the case
+        // where a node has not otherwise been picked up in this poll, but we 
need to send a
+        // ShareFetch to remove the stale partitions from the share session.
+        sessionHandlers.forEach((nodeId, sessionHandler) -> {
+            Node node = cluster.nodeById(nodeId);
+            if (node != null && !handlerMap.containsKey(node) && 
!sessionHandler.sessionPartitionMap().isEmpty()) {
+                if (nodesWithPendingRequests.contains(node.id())) {
+                    log.trace("Skipping fetch because previous fetch request 
to {} has not been processed", nodeId);
+                } else {
+                    Set<TopicPartition> currentPartitionsToFetch = new 
HashSet<>(partitionsToFetch());
+                    boolean hasPartitionsToRemove = 
sessionHandler.sessionPartitions().stream()
+                        .anyMatch(tip -> 
!currentPartitionsToFetch.contains(tip.topicPartition()));
+                    if (hasPartitionsToRemove) {
+                        handlerMap.put(node, sessionHandler);
+                        log.debug("Added fetch request for previously 
subscribed partitions without acknowledgements to node {}", nodeId);
+                    }
+                }
+            }
+        });
+
         // Iterate over the share session handlers and build a list of 
UnsentRequests.
         List<UnsentRequest> requests = 
handlerMap.entrySet().stream().map(entry -> {
             Node target = entry.getKey();
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
index ca8d0ca0865..4530965a35d 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareConsumeRequestManagerTest.java
@@ -1115,20 +1115,65 @@ public class ShareConsumeRequestManagerTest {
         // Change the subscription.
         subscriptions.assignFromSubscribed(List.of(tp1));
 
-        // Now we will be sending the request to node1 only as leader for tip1 
is node1.
-        // We do not build the request for tip0 as there are no 
acknowledgements to send.
+        // We build a request to node 1 to fetch tip1, and a request to node 0 
to remove tip0
+        // from the share session even though there are no acknowledgements to 
send.
         NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
-        assertEquals(1, pollResult.unsentRequests.size());
-        assertEquals(nodeId1, pollResult.unsentRequests.get(0).node().get());
+        assertEquals(2, pollResult.unsentRequests.size());
 
-        ShareFetchRequest.Builder builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+        ShareFetchRequest.Builder node0Builder, node1Builder;
+        if (pollResult.unsentRequests.get(0).node().get() == nodeId0) {
+            node0Builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+            node1Builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(1).requestBuilder();
+            assertEquals(nodeId1, 
pollResult.unsentRequests.get(1).node().get());
+        } else {
+            node0Builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(1).requestBuilder();
+            node1Builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+            assertEquals(nodeId0, 
pollResult.unsentRequests.get(1).node().get());
+            assertEquals(nodeId1, 
pollResult.unsentRequests.get(0).node().get());
+        }
 
-        assertEquals(1, builder.data().topics().size());
-        ShareFetchRequestData.FetchTopic fetchTopic = 
builder.data().topics().stream().findFirst().get();
+        // node1 fetches the newly assigned partition tip1.
+        assertEquals(1, node1Builder.data().topics().size());
+        ShareFetchRequestData.FetchTopic fetchTopic = 
node1Builder.data().topics().stream().findFirst().get();
         assertEquals(tip1.topicId(), fetchTopic.topicId());
         assertEquals(1, fetchTopic.partitions().size());
         assertEquals(1, 
fetchTopic.partitions().stream().findFirst().get().partitionIndex());
-        assertEquals(0, builder.data().forgottenTopicsData().size());
+        assertEquals(0, node1Builder.data().forgottenTopicsData().size());
+
+        // node0 removes tip0 from the share session and fetches nothing.
+        assertEquals(0, node0Builder.data().topics().size());
+        assertEquals(1, node0Builder.data().forgottenTopicsData().size());
+        assertEquals(tip0.topicId(), 
node0Builder.data().forgottenTopicsData().get(0).topicId());
+        assertEquals(1, 
node0Builder.data().forgottenTopicsData().get(0).partitions().size());
+        assertEquals(0, 
node0Builder.data().forgottenTopicsData().get(0).partitions().get(0));
+    }
+
+    @Test
+    public void testShareFetchRemovesUnassignedPartitionFromSession() {
+        buildRequestManager();
+
+        assignFromSubscribed(Set.of(tp0));
+
+        // Establish the share session by fetching from tp0.
+        sendFetchAndVerifyResponse(records, emptyAcquiredRecords, Errors.NONE);
+        fetchRecords();
+
+        // The partition is no longer assigned and there are no 
acknowledgements to send.
+        subscriptions.assignFromSubscribed(Set.of());
+
+        // We still build a ShareFetch to remove tip0 from the share session 
on the broker.
+        NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
+        assertEquals(1, pollResult.unsentRequests.size());
+
+        ShareFetchRequest.Builder builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+        assertEquals(0, builder.data().topics().size());
+        assertEquals(1, builder.data().forgottenTopicsData().size());
+        assertEquals(tip0.topicId(), 
builder.data().forgottenTopicsData().get(0).topicId());
+        assertEquals(1, 
builder.data().forgottenTopicsData().get(0).partitions().size());
+        assertEquals(0, 
builder.data().forgottenTopicsData().get(0).partitions().get(0));
+
+        // The partition has already been removed from the session, so no 
further ShareFetch is built.
+        assertEquals(0, shareConsumeRequestManager.sendFetches());
     }
 
     @Test
diff --git 
a/core/src/test/scala/unit/kafka/server/ShareFetchAcknowledgeRequestTest.scala 
b/core/src/test/scala/unit/kafka/server/ShareFetchAcknowledgeRequestTest.scala
index 597d1dc7036..0082b26ecd0 100644
--- 
a/core/src/test/scala/unit/kafka/server/ShareFetchAcknowledgeRequestTest.scala
+++ 
b/core/src/test/scala/unit/kafka/server/ShareFetchAcknowledgeRequestTest.scala
@@ -2203,8 +2203,8 @@ class ShareFetchAcknowledgeRequestTest(cluster: 
ClusterInstance) extends GroupCo
     // Send the second share fetch request to fetch the records produced above
     var shareSessionEpoch = 
ShareRequestMetadata.nextEpoch(ShareRequestMetadata.INITIAL_EPOCH)
     var metadata = new ShareRequestMetadata(MEMBER_ID, shareSessionEpoch)
-    val acknowledgementsMap = util.Map.of[TopicIdPartition, 
util.List[ShareFetchRequestData.AcknowledgementBatch]]
-    var shareFetchRequest = createShareFetchRequest(GROUP_ID, metadata, send, 
util.List.of, acknowledgementsMap)
+    val acknowledgementsMapEmpty = util.Map.of[TopicIdPartition, 
util.List[ShareFetchRequestData.AcknowledgementBatch]]
+    var shareFetchRequest = createShareFetchRequest(GROUP_ID, metadata, send, 
util.List.of, acknowledgementsMapEmpty)
 
     // For the multi partition fetch request, the response may not be 
available in the first attempt
     // as the share partitions might not be initialized yet. So, we retry 
until we get the response.
@@ -2231,14 +2231,14 @@ class ShareFetchAcknowledgeRequestTest(cluster: 
ClusterInstance) extends GroupCo
     produceData(topicIdPartition1, 10)
     produceData(topicIdPartition2, 10)
 
-    // Send another share fetch request with forget list populated with 
topicIdPartition2
+    // Send another share fetch request with forget list populated with 
topicIdPartition1
     shareSessionEpoch = ShareRequestMetadata.nextEpoch(shareSessionEpoch)
     metadata = new ShareRequestMetadata(MEMBER_ID, shareSessionEpoch)
     val forget = util.List.of(topicIdPartition1)
-    shareFetchRequest = createShareFetchRequest(GROUP_ID, metadata, 
util.List.of, forget, acknowledgementsMap)
-    val shareFetchResponse = 
IntegrationTestUtils.sendAndReceive[ShareFetchResponse](shareFetchRequest, 
socket)
+    shareFetchRequest = createShareFetchRequest(GROUP_ID, metadata, 
util.List.of, forget, acknowledgementsMapEmpty)
+    var shareFetchResponse = 
IntegrationTestUtils.sendAndReceive[ShareFetchResponse](shareFetchRequest, 
socket)
 
-    val shareFetchResponseData = shareFetchResponse.data()
+    var shareFetchResponseData = shareFetchResponse.data()
     assertEquals(Errors.NONE.code, shareFetchResponseData.errorCode)
     assertEquals(30000, shareFetchResponseData.acquisitionLockTimeoutMs)
     assertEquals(1, shareFetchResponseData.responses().size())
@@ -2253,6 +2253,70 @@ class ShareFetchAcknowledgeRequestTest(cluster: 
ClusterInstance) extends GroupCo
 
     val partitionData = 
shareFetchResponseData.responses().stream().findFirst().get().partitions().get(0)
     compareFetchResponsePartitions(expectedPartitionData, partitionData)
+
+    // Forgetting topicIdPartition1 is idempotent and there are no more 
records to fetch from topicIdPartition2
+    shareSessionEpoch = ShareRequestMetadata.nextEpoch(shareSessionEpoch)
+    metadata = new ShareRequestMetadata(MEMBER_ID, shareSessionEpoch)
+    shareFetchRequest = createShareFetchRequest(GROUP_ID, metadata, 
util.List.of, forget, acknowledgementsMapEmpty, maxWaitMs = 500)
+    shareFetchResponse = 
IntegrationTestUtils.sendAndReceive[ShareFetchResponse](shareFetchRequest, 
socket)
+
+    shareFetchResponseData = shareFetchResponse.data()
+    assertEquals(Errors.NONE.code, shareFetchResponseData.errorCode)
+    assertEquals(30000, shareFetchResponseData.acquisitionLockTimeoutMs)
+    assertEquals(0, shareFetchResponseData.responses().size())
+
+    // Now acknowledge some of the records received on topicIdPartition1 even 
though the share session does not contain the partition
+    // because forgetting a partition just stops fetching records and does not 
affect records which are already acquired
+    shareSessionEpoch = ShareRequestMetadata.nextEpoch(shareSessionEpoch)
+    metadata = new ShareRequestMetadata(MEMBER_ID, shareSessionEpoch)
+    val acknowledgementsMapForAcknowledge: util.Map[TopicIdPartition, 
util.List[ShareAcknowledgeRequestData.AcknowledgementBatch]] =
+      util.Map.of(topicIdPartition1, util.List.of(new 
ShareAcknowledgeRequestData.AcknowledgementBatch()
+          .setFirstOffset(5)
+          .setLastOffset(9)
+          .setAcknowledgeTypes(util.List.of(1.toByte)))) // Accept the records
+
+    var shareAcknowledgeRequest = createShareAcknowledgeRequest(GROUP_ID, 
metadata, acknowledgementsMapForAcknowledge)
+    var shareAcknowledgeResponse = 
IntegrationTestUtils.sendAndReceive[ShareAcknowledgeResponse](shareAcknowledgeRequest,
 socket)
+
+    var shareAcknowledgeResponseData = shareAcknowledgeResponse.data()
+    assertEquals(Errors.NONE.code, shareAcknowledgeResponseData.errorCode)
+    assertEquals(30000, shareAcknowledgeResponseData.acquisitionLockTimeoutMs)
+    assertEquals(1, shareAcknowledgeResponseData.responses().size())
+    assertEquals(topicId, 
shareAcknowledgeResponseData.responses().stream().findFirst().get().topicId())
+    assertEquals(1, 
shareAcknowledgeResponseData.responses().stream().findFirst().get().partitions().size())
+
+    var expectedAcknowledgePartitionData = new 
ShareAcknowledgeResponseData.PartitionData()
+      .setPartitionIndex(partition1)
+      .setErrorCode(Errors.NONE.code())
+
+    var acknowledgePartitionData = 
shareAcknowledgeResponseData.responses().stream().findFirst().get().partitions().get(0)
+    compareAcknowledgeResponsePartitions(expectedAcknowledgePartitionData, 
acknowledgePartitionData)
+
+    // Finally acknowledge the rest of the records, including one which has 
already been acknowledged
+    shareSessionEpoch = ShareRequestMetadata.nextEpoch(shareSessionEpoch)
+    metadata = new ShareRequestMetadata(MEMBER_ID, shareSessionEpoch)
+    val acknowledgementsMapForAcknowledgeInvalid: util.Map[TopicIdPartition, 
util.List[ShareAcknowledgeRequestData.AcknowledgementBatch]] =
+      util.Map.of(topicIdPartition1, util.List.of(new 
ShareAcknowledgeRequestData.AcknowledgementBatch()
+        .setFirstOffset(0)
+        .setLastOffset(5)
+        .setAcknowledgeTypes(util.List.of(1.toByte)))) // Accept the records
+
+    shareAcknowledgeRequest = createShareAcknowledgeRequest(GROUP_ID, 
metadata, acknowledgementsMapForAcknowledgeInvalid)
+    shareAcknowledgeResponse = 
IntegrationTestUtils.sendAndReceive[ShareAcknowledgeResponse](shareAcknowledgeRequest,
 socket)
+
+    shareAcknowledgeResponseData = shareAcknowledgeResponse.data()
+    assertEquals(Errors.NONE.code, shareAcknowledgeResponseData.errorCode)
+    assertEquals(30000, shareAcknowledgeResponseData.acquisitionLockTimeoutMs)
+    assertEquals(1, shareAcknowledgeResponseData.responses().size())
+    assertEquals(topicId, 
shareAcknowledgeResponseData.responses().stream().findFirst().get().topicId())
+    assertEquals(1, 
shareAcknowledgeResponseData.responses().stream().findFirst().get().partitions().size())
+
+    expectedAcknowledgePartitionData = new 
ShareAcknowledgeResponseData.PartitionData()
+      .setPartitionIndex(partition1)
+      .setErrorCode(Errors.INVALID_RECORD_STATE.code())
+
+    acknowledgePartitionData = 
shareAcknowledgeResponseData.responses().stream().findFirst().get().partitions().get(0)
+    compareAcknowledgeResponsePartitions(expectedAcknowledgePartitionData, 
acknowledgePartitionData)
   }
 
   @ClusterTests(

Reply via email to