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 33ea66bf7ef KAFKA-20736: Improve share session handling on leader 
change (#22766)
33ea66bf7ef is described below

commit 33ea66bf7ef3821aaacb983308080ec0883a4f83
Author: Andrew Schofield <[email protected]>
AuthorDate: Mon Jul 20 11:02:50 2026 +0100

    KAFKA-20736: Improve share session handling on leader change (#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 <[email protected]>, Shivsundar R
     <[email protected]>, Apoorv Mittal <[email protected]>
---
 .../kafka/clients/consumer/ShareConsumerTest.java  |  11 +-
 .../internals/ShareConsumeRequestManager.java      | 370 +++++++++++++-------
 .../consumer/internals/ShareSessionHandler.java    |  30 +-
 .../internals/ShareConsumeRequestManagerTest.java  | 384 +++++++++++++++++++--
 .../internals/ShareSessionHandlerTest.java         |  81 +++++
 5 files changed, 710 insertions(+), 166 deletions(-)

diff --git 
a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/consumer/ShareConsumerTest.java
 
b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/consumer/ShareConsumerTest.java
index d4865ab6cf8..5b664c22ec3 100644
--- 
a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/consumer/ShareConsumerTest.java
+++ 
b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/consumer/ShareConsumerTest.java
@@ -33,10 +33,9 @@ import org.apache.kafka.common.Uuid;
 import org.apache.kafka.common.errors.GroupMaxSizeReachedException;
 import org.apache.kafka.common.errors.InterruptException;
 import org.apache.kafka.common.errors.InvalidTopicException;
-import org.apache.kafka.common.errors.NotLeaderOrFollowerException;
+import org.apache.kafka.common.errors.NetworkException;
 import org.apache.kafka.common.errors.RecordDeserializationException;
 import org.apache.kafka.common.errors.SerializationException;
-import org.apache.kafka.common.errors.ShareSessionNotFoundException;
 import org.apache.kafka.common.errors.UnknownTopicIdException;
 import org.apache.kafka.common.errors.WakeupException;
 import org.apache.kafka.common.header.Header;
@@ -1054,7 +1053,7 @@ public class ShareConsumerTest extends 
ShareConsumerTestBase {
 
             AtomicBoolean callbackCalled = new AtomicBoolean(false);
             
shareConsumer.setAcknowledgementCommitCallback((offsetsByTopicPartition, 
exception) -> {
-                assertInstanceOf(NotLeaderOrFollowerException.class, 
exception);
+                assertInstanceOf(NetworkException.class, exception);
                 callbackCalled.set(true);
             });
 
@@ -1083,7 +1082,7 @@ public class ShareConsumerTest extends 
ShareConsumerTestBase {
             assertEquals(1, commitResult.size());
             TopicIdPartition tidp = commitResult.keySet().iterator().next();
             assertTrue(commitResult.get(tidp).isPresent());
-            assertInstanceOf(NotLeaderOrFollowerException.class, 
commitResult.get(tidp).get());
+            assertInstanceOf(NetworkException.class, 
commitResult.get(tidp).get());
 
             assertTrue(callbackCalled.get());
         }
@@ -1098,7 +1097,7 @@ public class ShareConsumerTest extends 
ShareConsumerTestBase {
 
             AtomicBoolean callbackCalled = new AtomicBoolean(false);
             
shareConsumer.setAcknowledgementCommitCallback((offsetsByTopicPartition, 
exception) -> {
-                assertInstanceOf(NotLeaderOrFollowerException.class, 
exception);
+                assertInstanceOf(NetworkException.class, exception);
                 callbackCalled.set(true);
             });
 
@@ -1147,7 +1146,7 @@ public class ShareConsumerTest extends 
ShareConsumerTestBase {
 
             AtomicBoolean callbackCalled = new AtomicBoolean(false);
             
shareConsumer.setAcknowledgementCommitCallback((offsetsByTopicPartition, 
exception) -> {
-                assertInstanceOf(ShareSessionNotFoundException.class, 
exception);
+                assertInstanceOf(NetworkException.class, exception);
                 callbackCalled.set(true);
             });
 
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 779a68b93b0..d9399148784 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
@@ -71,6 +71,15 @@ import java.util.stream.Collectors;
  * {@code ShareConsumeRequestManager} is responsible for generating {@link 
ShareFetchRequest} and
  * {@link ShareAcknowledgeRequest} to fetch and acknowledge records being 
delivered for a consumer
  * in a share group.
+ *
+ * <p>The request manager keeps track of which share sessions contain which 
topic-partitions. This information will
+ * generally match the leaders in the metadata, but when the fetchable 
partitions or leadership change, there can be
+ * short-lived differences. This is intentional to ensure that 
topic-partitions which are no longer being
+ * fetched and topic-partitions which have changed leaders are cleanly removed 
from their former share sessions.
+ * It is all slightly complicated by the fact that the fetchable partitions 
are represented by TopicPartition
+ * (without topic ID) and, when a topic is recreated, its topic ID changes. 
All state in share sessions is tracked by
+ * topic ID so the current mapping from topic name to topic ID as represented 
in the share sessions is also maintained
+ * in the request manager, and again there can be short-lived differences 
between this and the cluster metadata.
  */
 @SuppressWarnings({"NPathComplexity", "CyclomaticComplexity"})
 public class ShareConsumeRequestManager implements RequestManager, 
MemberStateListener, Closeable {
@@ -90,6 +99,8 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
     private Uuid memberId;
     private boolean fetchMoreRecords = false;
     private final AtomicInteger fetchRecordsNodeId = new AtomicInteger(-1);
+    private final Map<TopicPartition, TopicIdPartition> shareSessionTopicIdMap;
+    private final Map<TopicIdPartition, LeaderIdAndEpoch> 
shareSessionLeaderMap;
     private final Map<Integer, Map<TopicIdPartition, Acknowledgements>> 
fetchAcknowledgementsToSend;
     private final Map<Integer, Map<TopicIdPartition, Acknowledgements>> 
fetchAcknowledgementsInFlight;
     private final Map<Integer, Tuple<AcknowledgeRequestState>> 
acknowledgeRequestStates;
@@ -127,6 +138,8 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         this.sessionHandlers = new HashMap<>();
         this.nodesWithPendingRequests = new HashSet<>();
         this.acknowledgeRequestStates = new HashMap<>();
+        this.shareSessionTopicIdMap = new HashMap<>();
+        this.shareSessionLeaderMap = new HashMap<>();
         this.fetchAcknowledgementsToSend = new HashMap<>();
         this.fetchAcknowledgementsInFlight = new HashMap<>();
         this.closeFuture = new CompletableFuture<>();
@@ -151,121 +164,102 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
             return PollResult.EMPTY;
         }
 
+        // Iterate over the partitions to fetch, building a map from partition 
to leader node ID
         Map<Node, ShareSessionHandler> handlerMap = new HashMap<>();
+        Cluster cluster = metadata.fetch();
         Map<String, Uuid> topicIds = metadata.topicIds();
         for (TopicPartition partition : partitionsToFetch()) {
-            Optional<Node> leaderOpt = 
metadata.currentLeader(partition).leader;
+            TopicIdPartition tip = shareSessionTopicIdMap.get(partition);
+            if (tip == null) {
+                Uuid topicId = topicIds.get(partition.topic());
+                if (topicId == null) {
+                    log.debug("Requesting metadata update for partition {} 
since topic ID is missing", partition);
+                    metadata.requestUpdate(false);
+                    continue;
+                }
 
-            if (leaderOpt.isEmpty()) {
-                log.debug("Requesting metadata update for partition {} since 
current leader node is missing", partition);
-                metadata.requestUpdate(false);
-                continue;
+                tip = new TopicIdPartition(topicId, partition);
+                shareSessionTopicIdMap.put(partition, tip);
             }
 
-            Uuid topicId = topicIds.get(partition.topic());
-            if (topicId == null) {
-                log.debug("Requesting metadata update for partition {} since 
topic ID is missing", partition);
-                metadata.requestUpdate(false);
-                continue;
+            LeaderIdAndEpoch leader = shareSessionLeaderMap.get(tip);
+            if (leader == null || cluster.nodeById(leader.leaderId) == null) {
+                Metadata.LeaderAndEpoch leaderOpt = 
metadata.currentLeader(partition);
+                if (leaderOpt.leader.isEmpty() || 
cluster.nodeById(leaderOpt.leader.get().id()) == null) {
+                    log.debug("Requesting metadata update for partition {} 
since current leader node is missing", partition);
+                    metadata.requestUpdate(false);
+                    shareSessionLeaderMap.remove(tip);
+                    continue;
+                }
+
+                shareSessionLeaderMap.put(tip, new 
LeaderIdAndEpoch(leaderOpt.leader.get().id(), leaderOpt.epoch.orElse(-1)));
             }
+        }
 
-            Node node = leaderOpt.get();
-            if (nodesWithPendingRequests.contains(node.id())) {
-                log.trace("Skipping fetch for partition {} because previous 
fetch request to {} has not been processed", partition, node.id());
+        Set<Integer> missingNodes = new HashSet<>();
+        for (TopicPartition partition : partitionsToFetch()) {
+            TopicIdPartition tip = shareSessionTopicIdMap.get(partition);
+            if (tip == null) {
+                continue;
+            }
+            LeaderIdAndEpoch leader = shareSessionLeaderMap.get(tip);
+            if (leader == null) {
+                continue;
+            }
+            int nodeId = leader.leaderId;
+            Node node = cluster.nodeById(nodeId);
+            if (node == null) {
+                shareSessionLeaderMap.remove(tip);
+                missingNodes.add(nodeId);
+                continue;
+            }
+            if (nodesWithPendingRequests.contains(nodeId)) {
+                log.trace("Skipping fetch for partition {} because previous 
request to {} has not been processed", partition, nodeId);
             } else {
-                // If there is a leader and no in-flight requests, issue a new 
fetch.
                 ShareSessionHandler handler = handlerMap.computeIfAbsent(node,
-                        k -> sessionHandlers.computeIfAbsent(node.id(), n -> 
new ShareSessionHandler(logContext, n, memberId)));
+                    k -> sessionHandlers.computeIfAbsent(nodeId, n -> new 
ShareSessionHandler(logContext, n, memberId)));
 
-                TopicIdPartition tip = new TopicIdPartition(topicId, 
partition);
                 Acknowledgements acknowledgementsToSend = null;
-                boolean canSendAcknowledgements = true;
 
-                Map<TopicIdPartition, Acknowledgements> nodeAcksFromFetchMap = 
fetchAcknowledgementsToSend.get(node.id());
+                Map<TopicIdPartition, Acknowledgements> nodeAcksFromFetchMap = 
fetchAcknowledgementsToSend.get(nodeId);
                 if (nodeAcksFromFetchMap != null) {
                     acknowledgementsToSend = nodeAcksFromFetchMap.remove(tip);
 
                     if (acknowledgementsToSend != null) {
                         // Check if the share session epoch is valid for 
sending acknowledgements.
                         if (!maybeAddAcknowledgements(handler, node, tip, 
acknowledgementsToSend)) {
-                            canSendAcknowledgements = false;
+                            acknowledgementsToSend = null;
                         }
                     }
                 }
 
-                if (canSendAcknowledgements) {
-                    handler.addPartitionToFetch(tip, acknowledgementsToSend);
-                } else {
-                    handler.addPartitionToFetch(tip, null);
-                }
+                handler.addPartitionToFetch(tip, acknowledgementsToSend);
                 topicNamesMap.putIfAbsent(new IdAndPartition(tip.topicId(), 
tip.partition()), tip.topic());
 
-                // If we have not chosen a node for fetching records yet,
-                // choose now, and rotate the assigned partitions so the next 
poll starts on a different partition.
-                // This is only applicable for record_limit mode.
-                if (isShareAcquireModeRecordLimit() && 
fetchRecordsNodeId.compareAndSet(-1, node.id())) {
+                // If we have not chosen a node for fetching records yet, 
choose now, and rotate the assigned partitions
+                // so the next poll starts on a different partition. This is 
only applicable for record_limit mode.
+                if (isShareAcquireModeRecordLimit() && 
fetchRecordsNodeId.compareAndSet(-1, nodeId)) {
                     subscriptions.movePartitionToEnd(partition);
                 }
 
-                log.debug("Added fetch request for partition {} to node {}", 
tip, node.id());
+                log.debug("Added fetch request for partition {} to node {}", 
tip, nodeId);
             }
         }
 
         // Iterate over the session handlers to see if there are 
acknowledgements to be sent for partitions
-        // which are no longer part of the current subscription.
-        // We fail acknowledgements for records fetched from a previous leader.
-        Cluster cluster = metadata.fetch();
+        // which are no longer part of the current subscription or which have 
disappeared from the metadata.
+        // Also include session handlers which have no fetchable partitions so 
the share sessions can be
+        // brought up to date.
         sessionHandlers.forEach((nodeId, sessionHandler) -> {
             Node node = cluster.nodeById(nodeId);
-            if (node != null) {
-                if (nodesWithPendingRequests.contains(node.id())) {
-                    log.trace("Skipping fetch because previous fetch request 
to {} has not been processed", nodeId);
-                } else {
-                    Map<TopicIdPartition, Acknowledgements> 
nodeAcksFromFetchMap = fetchAcknowledgementsToSend.get(nodeId);
-                    if (nodeAcksFromFetchMap != null) {
-                        nodeAcksFromFetchMap.forEach((tip, acks) -> {
-                            if (!isLeaderKnownToHaveChanged(nodeId, tip)) {
-                                // Check if the share session epoch is valid 
for sending acknowledgements.
-                                if (!maybeAddAcknowledgements(sessionHandler, 
node, tip, acks)) {
-                                    return;
-                                }
-
-                                
sessionHandler.addPartitionToAcknowledgeOnly(tip, acks);
-                                handlerMap.put(node, sessionHandler);
-
-                                topicNamesMap.putIfAbsent(new 
IdAndPartition(tip.topicId(), tip.partition()), tip.topic());
-                                log.debug("Added fetch request for previously 
subscribed partition {} to node {}", tip, nodeId);
-                            } else {
-                                log.debug("Leader for the partition is down or 
has changed, failing acknowledgements for partition {}", tip);
-                                
acks.complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
-                                maybeSendShareAcknowledgementEvent(Map.of(tip, 
acks), true, Optional.empty());
-                            }
-                        });
-
-                        nodeAcksFromFetchMap.clear();
-                    }
-                }
+            if (node == null) {
+                missingNodes.add(nodeId);
+                return;
             }
-        });
-
-        // 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);
-                    }
-                }
+            if (nodesWithPendingRequests.contains(nodeId)) {
+                log.trace("Skipping fetch because previous request to {} has 
not been processed", nodeId);
+            } else {
+                prepareRemainingSessionHandlerForRequest(node, sessionHandler, 
handlerMap);
             }
         });
 
@@ -274,10 +268,9 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
             Node target = entry.getKey();
             ShareSessionHandler handler = entry.getValue();
 
-            // For record_limit mode, we only send a full ShareFetch to a 
single node at a time.
-            // We prepare to build ShareFetch requests for all nodes with 
session handlers to permit
-            // piggybacking of acknowledgements, and also to adjust the 
topic-partitions
-            // in the share session, but if the request would contain neither 
of those, it can be skipped.
+            // For record_limit mode, we only send a full ShareFetch to a 
single node at a time. We prepare to build ShareFetch
+            // requests for all nodes with session handlers to permit 
piggybacking of acknowledgements, and also to adjust the
+            // topic-partitions in the share session, but if the request would 
contain neither of those, it can be skipped.
             boolean canSkipIfRequestEmpty = isShareAcquireModeRecordLimit() && 
target.id() != fetchRecordsNodeId.get();
 
             ShareFetchRequest.Builder requestBuilder = 
handler.newShareFetchBuilder(groupId, shareFetchConfig, canSkipIfRequestEmpty);
@@ -300,6 +293,11 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
             return new UnsentRequest(requestBuilder, 
Optional.of(target)).whenComplete(responseHandler);
         }).filter(Objects::nonNull).collect(Collectors.toList());
 
+        // Remove session handlers for nodes which are now missing from the 
cluster metadata.
+        if (!missingNodes.isEmpty()) {
+            removeSessionHandlersForMissingNodes(missingNodes);
+        }
+
         return new PollResult(requests);
     }
 
@@ -309,7 +307,7 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
 
     /**
      * Add acknowledgements for a topic-partition to the node's in-flight 
acknowledgements.
-     * If we cannot add acknowledgements, they are completed with {@link 
Errors#NOT_LEADER_OR_FOLLOWER} exception.
+     * If we cannot add acknowledgements, they are completed with {@link 
Errors#NETWORK_EXCEPTION} exception.
      * This probably indicates the connection to the leader broker was lost, 
but then re-established without a
      * leadership change, in which case the acknowledgements fail.
      *
@@ -322,7 +320,7 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         if (handler.isNewSession()) {
             // Failing the acknowledgements as we cannot have piggybacked 
acknowledgements in the initial ShareFetchRequest.
             log.debug("Cannot send acknowledgements on initial epoch for 
ShareSession for partition {}", tip);
-            
acknowledgements.complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
+            acknowledgements.complete(Errors.NETWORK_EXCEPTION.exception());
             maybeSendShareAcknowledgementEvent(Map.of(tip, acknowledgements), 
true, Optional.empty());
             return false;
         } else {
@@ -332,6 +330,65 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         }
     }
 
+    /**
+     * Iterate over the session handlers to see if there are acknowledgements 
to be sent for partitions
+     * which are no longer part of the current subscription or which have 
disappeared from the metadata.
+     * Also include session handlers which have no fetchable partitions so the 
share sessions can be
+     * brought up to date.
+     */
+    private void prepareRemainingSessionHandlerForRequest(Node node, 
ShareSessionHandler sessionHandler, Map<Node, ShareSessionHandler> handlerMap) {
+        Map<TopicIdPartition, Acknowledgements> nodeAcksFromFetchMap = 
fetchAcknowledgementsToSend.get(node.id());
+        if (nodeAcksFromFetchMap != null) {
+            nodeAcksFromFetchMap.forEach((tip, acks) -> {
+                if (!isLeaderKnownToHaveChanged(node.id(), tip)) {
+                    // Check if the share session epoch is valid for sending 
acknowledgements.
+                    if (!maybeAddAcknowledgements(sessionHandler, node, tip, 
acks)) {
+                        return;
+                    }
+
+                    sessionHandler.addPartitionToAcknowledgeOnly(tip, acks);
+
+                    topicNamesMap.putIfAbsent(new 
IdAndPartition(tip.topicId(), tip.partition()), tip.topic());
+                    log.debug("Added fetch request for previously subscribed 
partition {} to node {}", tip, node.id());
+                } else {
+                    log.debug("Leader for the partition is down or has 
changed, failing acknowledgements for partition {}", tip);
+                    acks.complete(Errors.NETWORK_EXCEPTION.exception());
+                    maybeSendShareAcknowledgementEvent(Map.of(tip, acks), 
true, Optional.empty());
+                }
+            });
+
+            nodeAcksFromFetchMap.clear();
+            handlerMap.put(node, sessionHandler);
+        } else {
+            handlerMap.putIfAbsent(node, sessionHandler);
+        }
+    }
+
+    /**
+     * Remove session handlers for nodes which are now missing from the 
cluster metadata. Only nodes with no
+     * in-flight requests are removed, since those requests will complete. 
Piggyback acknowledgements for the
+     * missing nodes are completed with {@link Errors#NETWORK_EXCEPTION}.
+     */
+    private void removeSessionHandlersForMissingNodes(Set<Integer> 
missingNodes) {
+        Cluster cluster = metadata.fetch();
+        missingNodes.forEach(nodeId -> {
+            Node node = cluster.nodeById(nodeId);
+            if (node == null && !nodesWithPendingRequests.contains(nodeId)) {
+                Map<TopicIdPartition, Acknowledgements> nodeAcksFromFetchMap = 
fetchAcknowledgementsToSend.remove(nodeId);
+                if (nodeAcksFromFetchMap != null) {
+                    nodeAcksFromFetchMap.forEach((tip, acks) -> {
+                        log.debug("Node {} is no longer in the cluster 
metadata, failing acknowledgements for partition {}", nodeId, tip);
+                        acks.complete(Errors.NETWORK_EXCEPTION.exception());
+                        maybeSendShareAcknowledgementEvent(Map.of(tip, acks), 
true, Optional.empty());
+                    });
+                }
+
+                log.debug("Removing session handler for node {} which is no 
longer in the cluster metadata", nodeId);
+                sessionHandlers.remove(nodeId);
+            }
+        });
+    }
+
     public void fetch(Map<TopicIdPartition, NodeAcknowledgements> 
acknowledgementsMap) {
         if (!fetchMoreRecords) {
             log.debug("Fetch more data");
@@ -570,11 +627,14 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
 
         Map<Integer, Map<TopicIdPartition, Acknowledgements>> 
acknowledgementsMapAllNodes = new HashMap<>();
         Map<TopicIdPartition, Acknowledgements> acknowledgementsMapCannotSend 
= new HashMap<>();
+        Map<TopicIdPartition, Errors> acknowledgementsMapCannotSendErrors = 
new HashMap<>();
         acknowledgementsMap.forEach((tip, nodeAcks) -> {
             if ((cluster.nodeById(nodeAcks.nodeId()) == null) || 
isLeaderKnownToHaveChanged(nodeAcks.nodeId(), tip)) {
                 Acknowledgements prevAcks = 
acknowledgementsMapCannotSend.putIfAbsent(tip, nodeAcks.acknowledgements());
                 if (prevAcks != null) {
                     prevAcks.merge(nodeAcks.acknowledgements());
+                } else {
+                    acknowledgementsMapCannotSendErrors.put(tip, 
acknowledgementsCannotBeSentError(nodeAcks.nodeId(), tip));
                 }
             } else {
                 Map<TopicIdPartition, Acknowledgements> acksMap = 
acknowledgementsMapAllNodes.computeIfAbsent(nodeAcks.nodeId(), k -> new 
HashMap<>());
@@ -630,7 +690,7 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         });
 
         acknowledgementsMapCannotSend.forEach((tip, acks) -> {
-            acks.complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
+            
acks.complete(acknowledgementsMapCannotSendErrors.get(tip).exception());
             resultHandler.complete(tip, acks, 
AcknowledgeRequestType.COMMIT_SYNC, true, Optional.empty());
         });
 
@@ -655,7 +715,7 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         acknowledgementsMap.forEach((tip, nodeAcks) -> {
             if ((cluster.nodeById(nodeAcks.nodeId()) == null) || 
isLeaderKnownToHaveChanged(nodeAcks.nodeId(), tip)) {
                 log.debug("Leader for the partition is down or has changed, 
failing acknowledgements for partition {}", tip);
-                
nodeAcks.acknowledgements().complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
+                
nodeAcks.acknowledgements().complete(acknowledgementsCannotBeSentError(nodeAcks.nodeId(),
 tip).exception());
                 maybeSendShareAcknowledgementEvent(Map.of(tip, 
nodeAcks.acknowledgements()), true, Optional.empty());
             } else {
                 Map<TopicIdPartition, Acknowledgements> acksMap = 
acknowledgementsMapAllNodes.computeIfAbsent(nodeAcks.nodeId(), k -> new 
HashMap<>());
@@ -735,7 +795,7 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         Map<Integer, Map<TopicIdPartition, Acknowledgements>> 
acknowledgementsMapAllNodes = new HashMap<>();
         acknowledgementsMap.forEach((tip, nodeAcks) -> {
             if ((cluster.nodeById(nodeAcks.nodeId()) == null) || 
isLeaderKnownToHaveChanged(nodeAcks.nodeId(), tip)) {
-                
nodeAcks.acknowledgements().complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
+                
nodeAcks.acknowledgements().complete(acknowledgementsCannotBeSentError(nodeAcks.nodeId(),
 tip).exception());
                 maybeSendShareAcknowledgementEvent(Map.of(tip, 
nodeAcks.acknowledgements()), true, Optional.empty());
             } else {
                 Map<TopicIdPartition, Acknowledgements> acksMap = 
acknowledgementsMapAllNodes.computeIfAbsent(nodeAcks.nodeId(), k -> new 
HashMap<>());
@@ -750,7 +810,7 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         fetchAcknowledgementsToSend.forEach((nodeId, nodeAcks) ->
             nodeAcks.forEach((tip, acks) -> {
                 if ((cluster.nodeById(nodeId) == null) || 
isLeaderKnownToHaveChanged(nodeId, tip)) {
-                    acks.complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
+                    acks.complete(acknowledgementsCannotBeSentError(nodeId, 
tip).exception());
                     maybeSendShareAcknowledgementEvent(Map.of(tip, acks), 
true, Optional.empty());
                 } else {
                     Map<TopicIdPartition, Acknowledgements> acksMap = 
acknowledgementsMapAllNodes.computeIfAbsent(nodeId, k -> new HashMap<>());
@@ -852,11 +912,18 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
                 if (response.error() == Errors.UNKNOWN_TOPIC_ID) {
                     metadata.requestUpdate(false);
                 }
-                // Complete any in-flight acknowledgements with the error code 
from the response.
+                // Complete any in-flight acknowledgements with the error code 
from the response, unless the share session was lost,
+                // in which case they are failed with NETWORK_EXCEPTION 
reflecting a loss of connectivity.
+                final Errors ackError;
+                if (response.error() == Errors.SHARE_SESSION_NOT_FOUND || 
response.error() == Errors.INVALID_SHARE_SESSION_EPOCH) {
+                    ackError = Errors.NETWORK_EXCEPTION;
+                } else {
+                    ackError = response.error();
+                }
                 Map<TopicIdPartition, Acknowledgements> 
nodeAcknowledgementsInFlight = 
fetchAcknowledgementsInFlight.remove(fetchTarget.id());
-                if (nodeAcknowledgementsInFlight != null) {
+                if (nodeAcknowledgementsInFlight != null && 
!nodeAcknowledgementsInFlight.isEmpty()) {
                     nodeAcknowledgementsInFlight.forEach((tip, acks) -> {
-                        
acks.complete(Errors.forCode(response.error().code()).exception());
+                        acks.complete(ackError.exception());
                         
metricsManager.recordFailedAcknowledgements(acks.size());
                     });
                     
maybeSendShareAcknowledgementEvent(nodeAcknowledgementsInFlight, 
requestData.isRenewAck(), Optional.empty());
@@ -909,7 +976,14 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
                     if (partitionData.currentLeader().leaderId() != -1 && 
partitionData.currentLeader().leaderEpoch() != -1) {
                         
partitionsWithUpdatedLeaderInfo.put(tip.topicPartition(), new 
Metadata.LeaderIdAndEpoch(
                             
Optional.of(partitionData.currentLeader().leaderId()), 
Optional.of(partitionData.currentLeader().leaderEpoch())));
+
+                        maybeUpdateLeaderCache(tip, 
partitionData.currentLeader().leaderId(), 
partitionData.currentLeader().leaderEpoch());
+                    } else {
+                        shareSessionLeaderMap.remove(tip);
                     }
+                } else if (partitionError == Errors.UNKNOWN_TOPIC_OR_PARTITION 
|| partitionError == Errors.UNKNOWN_TOPIC_ID) {
+                    shareSessionLeaderMap.remove(tip);
+                    shareSessionTopicIdMap.remove(tip.topicPartition());
                 }
 
                 completedFetches.add(
@@ -976,7 +1050,6 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
                 Map<TopicIdPartition, Acknowledgements> 
nodeAcknowledgementsInFlight = 
fetchAcknowledgementsInFlight.get(fetchTarget.id());
                 if (nodeAcknowledgementsInFlight != null) {
                     Acknowledgements acks = 
nodeAcknowledgementsInFlight.remove(tip);
-
                     if (acks != null) {
                         
metricsManager.recordFailedAcknowledgements(acks.size());
                         if (error instanceof KafkaException) {
@@ -1119,13 +1192,14 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
                                       Optional<Integer> 
acquisitionLockTimeoutMs) {
         if (partitionError.exception() != null) {
             boolean retry = false;
-            if (partitionError == Errors.NOT_LEADER_OR_FOLLOWER || 
partitionError == Errors.FENCED_LEADER_EPOCH ||
-                partitionError == Errors.UNKNOWN_TOPIC_OR_PARTITION || 
partitionError == Errors.UNKNOWN_TOPIC_ID) {
+            if (partitionError == Errors.NOT_LEADER_OR_FOLLOWER || 
partitionError == Errors.FENCED_LEADER_EPOCH) {
                 // If the leader has changed, there's no point in retrying the 
operation because the acquisition locks
-                // will have been released.
+                // will have been released. Instead, these records will be 
re-delivered once they get timed out on the broker.
                 // If the topic or partition has been deleted, we do not retry 
the failed acknowledgements.
-                // Instead, these records will be re-delivered once they get 
timed out on the broker.
-                updateLeaderInfoMap(partitionData, 
partitionsWithUpdatedLeaderInfo, partitionError, tip.topicPartition());
+                updateLeaderInfoMap(partitionData, 
partitionsWithUpdatedLeaderInfo, partitionError, tip);
+            } else if (partitionError == Errors.UNKNOWN_TOPIC_OR_PARTITION || 
partitionError == Errors.UNKNOWN_TOPIC_ID) {
+                shareSessionLeaderMap.remove(tip);
+                shareSessionTopicIdMap.remove(tip.topicPartition());
             } else if (partitionError.exception() instanceof 
RetriableException) {
                 retry = true;
             }
@@ -1161,18 +1235,43 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
     private void 
updateLeaderInfoMap(ShareAcknowledgeResponseData.PartitionData partitionData,
                                   Map<TopicPartition, 
Metadata.LeaderIdAndEpoch> partitionsWithUpdatedLeaderInfo,
                                   Errors partitionError,
-                                  TopicPartition tp) {
+                                  TopicIdPartition tip) {
 
-        log.debug("For {}, received error {}, with leaderIdAndEpoch {} in 
ShareAcknowledge", tp, partitionError, partitionData.currentLeader());
+        log.debug("For {}, received error {}, with leaderIdAndEpoch {} in 
ShareAcknowledge", tip, partitionError, partitionData.currentLeader());
         if (partitionData.currentLeader().leaderId() != -1 && 
partitionData.currentLeader().leaderEpoch() != -1) {
-            partitionsWithUpdatedLeaderInfo.put(tp,
+            partitionsWithUpdatedLeaderInfo.put(tip.topicPartition(),
                 new Metadata.LeaderIdAndEpoch(
                     Optional.of(partitionData.currentLeader().leaderId()),
                     Optional.of(partitionData.currentLeader().leaderEpoch())
-            ));
+                ));
+
+            maybeUpdateLeaderCache(tip, 
partitionData.currentLeader().leaderId(), 
partitionData.currentLeader().leaderEpoch());
+        } else {
+            shareSessionLeaderMap.remove(tip);
+        }
+    }
+
+    /**
+     * Update the cache leader for a partition in a share session, only if the 
new leader epoch is newer than the
+     * currently cached epoch. This mirrors the rules applied by {@link 
Metadata#updateLastSeenEpochIfNewer(TopicPartition, int)}.
+     * A stale entry is never overwritten by an older epoch.
+     */
+    private void maybeUpdateLeaderCache(TopicIdPartition tip, int leaderId, 
int leaderEpoch) {
+        LeaderIdAndEpoch oldLeader = shareSessionLeaderMap.get(tip);
+        if ((oldLeader == null) || (leaderEpoch > oldLeader.epoch)) {
+            shareSessionLeaderMap.put(tip, new LeaderIdAndEpoch(leaderId, 
leaderEpoch));
         }
     }
 
+    /**
+     * Chooses the error used to fail acknowledgements which could not be sent 
to the node hosting the share session.
+     * If leadership has moved to a different live broker, the error is {@link 
Errors#NOT_LEADER_OR_FOLLOWER}.
+     * Otherwise, the error is {@link Errors#NETWORK_EXCEPTION}.
+     */
+    private Errors acknowledgementsCannotBeSentError(int nodeId, 
TopicIdPartition tip) {
+        return isLeaderKnownToHaveChanged(nodeId, tip) ? 
Errors.NOT_LEADER_OR_FOLLOWER : Errors.NETWORK_EXCEPTION;
+    }
+
     private TopicIdPartition lookupTopicId(Uuid topicId, int partitionIndex) {
         String topicName = metadata.topicNames().get(topicId);
         if (topicName == null) {
@@ -1288,6 +1387,13 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         }
 
         UnsentRequest buildRequest() {
+            // If the node is no longer in the cluster metadata, we can never 
send to it
+            Node nodeToSend = metadata.fetch().nodeById(nodeId);
+            if (nodeToSend == null) {
+                failPendingAcknowledgementsCannotBeSent();
+                return null;
+            }
+
             // If this is the closing request, close the share session by 
setting the final epoch
             if (isCloseRequest()) {
                 sessionHandler.notifyClose();
@@ -1303,35 +1409,32 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
             ShareAcknowledgeRequest.Builder requestBuilder = 
sessionHandler.newShareAcknowledgeBuilder(groupId);
 
             isProcessed = false;
-            Node nodeToSend = metadata.fetch().nodeById(nodeId);
 
             if (requestBuilder == null) {
-                handleNewShareSessionNotLeaderOrFollower();
+                failPendingAcknowledgementsCannotBeSent();
                 return null;
-            } else if (nodeToSend != null) {
-                nodesWithPendingRequests.add(nodeId);
+            }
 
-                log.trace("Building acknowledgements to send : {}", 
finalAcknowledgementsToSend);
+            nodesWithPendingRequests.add(nodeId);
 
-                inFlightAcknowledgements.putAll(finalAcknowledgementsToSend);
-                if (incompleteAcknowledgements.isEmpty()) {
-                    acknowledgementsToSend.clear();
-                } else {
-                    incompleteAcknowledgements.clear();
-                }
+            log.trace("Building acknowledgements to send : {}", 
finalAcknowledgementsToSend);
 
-                UnsentRequest unsentRequest = new 
UnsentRequest(requestBuilder, Optional.of(nodeToSend));
-                BiConsumer<ClientResponse, Throwable> responseHandler = 
(clientResponse, error) -> {
-                    if (error != null) {
-                        handleShareAcknowledgeFailure(nodeToSend, 
requestBuilder.data(), this, error, unsentRequest.handler().completionTimeMs());
-                    } else {
-                        handleShareAcknowledgeSuccess(nodeToSend, 
requestBuilder.data(), this, clientResponse, 
unsentRequest.handler().completionTimeMs());
-                    }
-                };
-                return unsentRequest.whenComplete(responseHandler);
+            inFlightAcknowledgements.putAll(finalAcknowledgementsToSend);
+            if (incompleteAcknowledgements.isEmpty()) {
+                acknowledgementsToSend.clear();
+            } else {
+                incompleteAcknowledgements.clear();
             }
 
-            return null;
+            UnsentRequest unsentRequest = new UnsentRequest(requestBuilder, 
Optional.of(nodeToSend));
+            BiConsumer<ClientResponse, Throwable> responseHandler = 
(clientResponse, error) -> {
+                if (error != null) {
+                    handleShareAcknowledgeFailure(nodeToSend, 
requestBuilder.data(), this, error, unsentRequest.handler().completionTimeMs());
+                } else {
+                    handleShareAcknowledgeSuccess(nodeToSend, 
requestBuilder.data(), this, clientResponse, 
unsentRequest.handler().completionTimeMs());
+                }
+            };
+            return unsentRequest.whenComplete(responseHandler);
         }
 
         int getInFlightAcknowledgementsCount(TopicIdPartition tip) {
@@ -1407,16 +1510,17 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         }
 
         /**
-         * Set the error code for all remaining acknowledgements in the event 
that a new share session
-         * needs to be started which prevents the remaining acknowledgements 
from being sent.
+         * Fail all remaining acknowledgements when they cannot be sent. This 
can happen when network connectivity
+         * is lost with the leader and the share session is lost, they are 
failed with {@link Errors#NETWORK_EXCEPTION}.
+         * If the leader has moved to a different live broker, they are failed 
with {@link Errors#NOT_LEADER_OR_FOLLOWER}.
          */
-        void handleNewShareSessionNotLeaderOrFollower() {
+        void failPendingAcknowledgementsCannotBeSent() {
             Map<TopicIdPartition, Acknowledgements> acknowledgementsMapToClear 
=
                 incompleteAcknowledgements.isEmpty() ? acknowledgementsToSend 
: incompleteAcknowledgements;
 
             acknowledgementsMapToClear.forEach((tip, acks) -> {
                 if (acks != null) {
-                    acks.complete(Errors.NOT_LEADER_OR_FOLLOWER.exception());
+                    acks.complete(acknowledgementsCannotBeSentError(nodeId, 
tip).exception());
                 }
                 // We do not know whether this is a renew ack, but handling 
the error as if it were, will ensure
                 // that we do not leave dangling acknowledgements
@@ -1614,6 +1718,16 @@ public class ShareConsumeRequestManager implements 
RequestManager, MemberStateLi
         }
     }
 
+    static class LeaderIdAndEpoch {
+        private final int leaderId;
+        private final int epoch;
+
+        LeaderIdAndEpoch(int leaderId, int epoch) {
+            this.leaderId = leaderId;
+            this.epoch = epoch;
+        }
+    }
+
     public enum AcknowledgeRequestType {
         COMMIT_ASYNC((byte) 0),
         COMMIT_SYNC((byte) 1),
diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandler.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandler.java
index 5b92e871e0e..c89716ad622 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandler.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandler.java
@@ -166,12 +166,22 @@ public class ShareSessionHandler {
         Map<TopicIdPartition, 
List<ShareFetchRequestData.AcknowledgementBatch>> acknowledgementBatches = new 
HashMap<>();
         if (!nextAcknowledgements.isEmpty()) {
             for (Map.Entry<TopicIdPartition, Acknowledgements> partitionsAcks 
: nextAcknowledgements.entrySet()) {
+                TopicIdPartition tip = partitionsAcks.getKey();
                 List<AcknowledgementBatch> partitionAckBatches = 
partitionsAcks.getValue().getAcknowledgementBatches();
                 for (AcknowledgementBatch ackBatch : partitionAckBatches) {
                     if 
(ackBatch.acknowledgeTypes().contains(AcknowledgeType.RENEW.id)) {
                         hasRenewAcknowledgements = true;
                     }
-                    
acknowledgementBatches.computeIfAbsent(partitionsAcks.getKey(), k -> new 
ArrayList<>()).add(ackBatch.toShareFetchRequest());
+                    acknowledgementBatches.computeIfAbsent(tip, k -> new 
ArrayList<>()).add(ackBatch.toShareFetchRequest());
+                }
+
+                // If the partition is only being included in the request to 
send acknowledgements, we need to
+                // remove it so that it doesn't get added into the share 
session for fetching
+                TopicIdPartition sessionTip = 
sessionPartitions.get(tip.topicPartition());
+                if ((sessionTip == null) || !sessionTip.equals(tip)) {
+                    if (!removed.contains(tip)) {
+                        removed.add(tip);
+                    }
                 }
             }
         }
@@ -179,16 +189,26 @@ public class ShareSessionHandler {
         nextPartitions = new LinkedHashMap<>();
         nextAcknowledgements = new LinkedHashMap<>();
 
-        if (canSkipIfRequestEmpty && added.isEmpty() && removed.isEmpty() && 
acknowledgementBatches.isEmpty()) {
-            return null;
+        // If there are no changes to the share session and no 
acknowledgements, we can sometimes skip sending an empty request
+        if (added.isEmpty() && removed.isEmpty() && 
acknowledgementBatches.isEmpty()) {
+            // If the share session is empty, there are no partitions to fetch 
from and we can always skip
+            if (sessionPartitions.isEmpty()) {
+                log.debug("Skipping sending empty ShareFetch because share 
partitions empty");
+                return null;
+            }
+
+            // If the share session is not empty, but we do not want to fetch 
records for this node, we can skip
+            if (canSkipIfRequestEmpty) {
+                log.debug("Skipping sending empty ShareFetch because no share 
session changes or acknowledgements");
+                return null;
+            }
         }
 
         if (log.isDebugEnabled()) {
-            log.debug("Build ShareFetch {} for node {}. Added {}, removed {}, 
replaced {} out of {}",
+            log.debug("Build ShareFetch {} for node {}. Added {}, removed {} 
out of {}",
                 nextMetadata, node,
                 topicIdPartitionsToLogString(added),
                 topicIdPartitionsToLogString(removed),
-                topicIdPartitionsToLogString(replaced),
                 topicIdPartitionsToLogString(sessionPartitions.values()));
         }
 
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 67864ee2d16..ce756395a37 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
@@ -40,6 +40,7 @@ import org.apache.kafka.common.errors.ApiException;
 import org.apache.kafka.common.errors.AuthenticationException;
 import org.apache.kafka.common.errors.DisconnectException;
 import org.apache.kafka.common.errors.InvalidRecordStateException;
+import org.apache.kafka.common.errors.NetworkException;
 import org.apache.kafka.common.errors.NotLeaderOrFollowerException;
 import org.apache.kafka.common.errors.TopicAuthorizationException;
 import org.apache.kafka.common.errors.UnknownServerException;
@@ -330,7 +331,7 @@ public class ShareConsumeRequestManagerTest {
     }
 
     @Test
-    public void testServerDisconnectedOnShareAcknowledge() throws 
InterruptedException {
+    public void testServerDisconnectedOnShareAcknowledge() {
         buildRequestManager();
         // Enabling the config so that background event is sent when the 
acknowledgement response is received.
         
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
@@ -371,7 +372,7 @@ public class ShareConsumeRequestManagerTest {
         assertNull(shareConsumeRequestManager.requestStates(0));
         // The callback for these unsent acknowledgements will be invoked with 
an error code.
         assertEquals(Map.of(tip0, acknowledgements2), 
completedAcknowledgements.get(0));
-        assertInstanceOf(NotLeaderOrFollowerException.class, 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+        assertInstanceOf(NetworkException.class, 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
 
         // Attempt a normal fetch to check if nodesWithPendingRequests is 
empty.
         assertEquals(1, sendFetches());
@@ -645,7 +646,7 @@ public class ShareConsumeRequestManagerTest {
     }
 
     @Test
-    public void testRetryAcknowledgements() throws InterruptedException {
+    public void testRetryAcknowledgements() {
         buildRequestManager();
 
         assignFromSubscribed(Set.of(tp0));
@@ -1264,7 +1265,7 @@ public class ShareConsumeRequestManagerTest {
     }
 
     @Test
-    public void testCallbackHandlerConfig() throws InterruptedException {
+    public void testCallbackHandlerConfig() {
         buildRequestManager();
         
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
 
@@ -1439,7 +1440,7 @@ public class ShareConsumeRequestManagerTest {
     }
 
     @Test
-    public void testShareAcknowledgeInvalidResponse() throws 
InterruptedException {
+    public void testShareAcknowledgeInvalidResponse() {
         buildRequestManager();
         
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
 
@@ -1537,11 +1538,11 @@ public class ShareConsumeRequestManagerTest {
         assertEquals(0, 
builder.data().topics().find(tip0.topicId()).partitions().find(0).acknowledgementBatches().size());
 
         assertEquals(3, completedAcknowledgements.get(0).get(tip0).size());
-        assertEquals(Errors.NOT_LEADER_OR_FOLLOWER.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+        assertEquals(Errors.NETWORK_EXCEPTION.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
     }
 
     @Test
-    public void 
testPiggybackAcknowledgementsOnInitialShareSessionErrorSubscriptionChange() {
+    public void 
testPiggybackAcknowledgementsOnInitialShareSessionErrorTopicRemovedFromMetadata()
 {
         buildRequestManager();
         
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
 
@@ -1567,13 +1568,15 @@ public class ShareConsumeRequestManagerTest {
 
         assertEquals(0, completedAcknowledgements.size());
 
-        // Next fetch would not include any acknowledgements.
+        // Next fetch would not include any acknowledgements, but it will 
include tip-0 because it's a new share session.
         NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
-        assertEquals(0, pollResult.unsentRequests.size());
+        assertEquals(1, pollResult.unsentRequests.size());
+        ShareFetchRequest.Builder builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+        assertEquals(1, builder.data().topics().size());
 
         // We should fail any waiting acknowledgements for tip-0 as it would 
have a share session epoch equal to 0.
         assertEquals(3, completedAcknowledgements.get(0).get(tip0).size());
-        assertEquals(Errors.NOT_LEADER_OR_FOLLOWER.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+        assertEquals(Errors.NETWORK_EXCEPTION.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
     }
 
     @Test
@@ -1600,7 +1603,7 @@ public class ShareConsumeRequestManagerTest {
 
         // We would complete these acknowledgements with the error code from 
the response.
         assertEquals(3, completedAcknowledgements.get(0).get(tip0).size());
-        assertEquals(Errors.SHARE_SESSION_NOT_FOUND.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+        assertEquals(Errors.NETWORK_EXCEPTION.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
 
         // Next fetch would proceed as expected and would not include any 
acknowledgements.
         NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
@@ -1957,8 +1960,8 @@ public class ShareConsumeRequestManagerTest {
 
         assertNotEquals(startingClusterMetadata, metadata.fetch());
 
-        // And now the partitions are on the same leader so only one fetch is 
sent
-        assertEquals(1, sendFetches());
+        // And now the partitions are on the same leader but a fetch is still 
sent to the former leader to remove the partition from the share session
+        assertEquals(2, sendFetches());
         assertFalse(shareConsumeRequestManager.hasCompletedFetches());
 
         partitionData = buildPartitionDataMap(tip0, records, 
ShareCompletedFetchTest.acquiredRecords(2L, 1), Errors.NONE, Errors.NONE);
@@ -1969,6 +1972,11 @@ public class ShareConsumeRequestManagerTest {
                 
.setAcquiredRecords(ShareCompletedFetchTest.acquiredRecords(1L, 1))
                 .setAcknowledgeErrorCode(Errors.NONE.code()));
         client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(), 0), nodeId0);
+        partitionData = new LinkedHashMap<>();
+        partitionData.put(tip1,
+            new ShareFetchResponseData.PartitionData()
+                .setPartitionIndex(tip1.topicPartition().partition()));
+        client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(), 0), nodeId1);
         networkClientDelegate.poll(time.timer(0));
         assertTrue(shareConsumeRequestManager.hasCompletedFetches());
 
@@ -2042,8 +2050,8 @@ public class ShareConsumeRequestManagerTest {
         // Validate metadata update is still requested even though the current 
leader was returned
         assertTrue(metadata.updateRequested());
 
-        // And now the partitions are on the same leader so only one fetch is 
sent
-        assertEquals(1, sendFetches());
+        // And now the partitions are on the same leader but a fetch is still 
sent to the former leader to remove the partition from the share session
+        assertEquals(2, sendFetches());
         assertFalse(shareConsumeRequestManager.hasCompletedFetches());
 
         partitionData = buildPartitionDataMap(tip0, records, 
ShareCompletedFetchTest.acquiredRecords(2L, 1), Errors.NONE, Errors.NONE);
@@ -2055,6 +2063,11 @@ public class ShareConsumeRequestManagerTest {
                 .setAcknowledgeErrorCode(Errors.NONE.code()));
         client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(), 0), nodeId0);
         networkClientDelegate.poll(time.timer(0));
+        partitionData = new LinkedHashMap<>();
+        partitionData.put(tip1,
+            new ShareFetchResponseData.PartitionData()
+                .setPartitionIndex(tip1.topicPartition().partition()));
+        client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(), 0), nodeId1);
         assertTrue(shareConsumeRequestManager.hasCompletedFetches());
 
         partitionRecords = fetchRecords();
@@ -2126,11 +2139,11 @@ public class ShareConsumeRequestManagerTest {
         assertNotEquals(startingClusterMetadata, metadata.fetch());
 
         // Even though the partitions are on the same leader, records were 
fetched on the previous leader.
-        // We do not send those acknowledgements to the previous leader, we 
fail them with NOT_LEADER_OR_FOLLOWER exception.
-        assertEquals(1, sendFetches());
+        // Since we need to remove the partition from the previous leader, we 
still send those acknowledgement to the previous leader
+        // with the partition in the list of partitions to forget.
+        assertEquals(2, sendFetches());
         assertFalse(shareConsumeRequestManager.hasCompletedFetches());
-        assertEquals(acknowledgements, 
completedAcknowledgements.get(0).get(tip0));
-        assertEquals(Errors.NOT_LEADER_OR_FOLLOWER.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+        assertTrue(completedAcknowledgements.isEmpty());
 
         partitionData.clear();
         partitionData.put(tip0,
@@ -2149,6 +2162,8 @@ public class ShareConsumeRequestManagerTest {
         client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(), 0), nodeId1);
         networkClientDelegate.poll(time.timer(0));
         assertTrue(shareConsumeRequestManager.hasCompletedFetches());
+        assertEquals(acknowledgements, 
completedAcknowledgements.get(0).get(tip0));
+        assertEquals(error.exception(), 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
 
         partitionRecords = fetchRecords();
         assertTrue(partitionRecords.containsKey(tp0));
@@ -2587,13 +2602,9 @@ public class ShareConsumeRequestManagerTest {
         List<ConsumerRecord<byte[], byte[]>> fetchedRecords = 
partitionRecords.get(tp0);
         assertEquals(1, fetchedRecords.size());
 
-        Acknowledgements acknowledgements = Acknowledgements.empty();
-        acknowledgements.add(1, AcknowledgeType.ACCEPT);
-        shareConsumeRequestManager.fetch(Map.of(tip0, new 
NodeAcknowledgements(0, acknowledgements)));
-
         assertEquals(startingClusterMetadata, metadata.fetch());
 
-        acknowledgements = Acknowledgements.empty();
+        Acknowledgements acknowledgements = Acknowledgements.empty();
         acknowledgements.add(1, AcknowledgeType.ACCEPT);
         shareConsumeRequestManager.fetch(Map.of(tip0, new 
NodeAcknowledgements(0, acknowledgements)));
 
@@ -2618,7 +2629,7 @@ public class ShareConsumeRequestManagerTest {
         networkClientDelegate.poll(time.timer(0));
         assertTrue(shareConsumeRequestManager.hasCompletedFetches());
 
-        // The node was disconnected, so the acknowledgement failed
+        // The node was disconnected, so the acknowledgements for tp0 failed
         assertInstanceOf(DisconnectException.class, 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
         completedAcknowledgements.clear();
 
@@ -2636,7 +2647,7 @@ public class ShareConsumeRequestManagerTest {
 
         shareConsumeRequestManager.fetch(Map.of(tip1, new 
NodeAcknowledgements(1, acknowledgements)));
 
-        assertEquals(1, sendFetches());
+        assertEquals(2, sendFetches());
         assertFalse(shareConsumeRequestManager.hasCompletedFetches());
 
         partitionData.clear();
@@ -2651,6 +2662,8 @@ public class ShareConsumeRequestManagerTest {
             new ShareFetchResponseData.PartitionData()
                 .setPartitionIndex(tip1.topicPartition().partition()));
         client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(), 0), nodeId1);
+        partitionData.clear();
+        
client.prepareResponseFrom(ShareFetchResponse.of(Errors.SHARE_SESSION_NOT_FOUND,
 0, partitionData, List.of(), 0), nodeId0);
         networkClientDelegate.poll(time.timer(0));
         assertTrue(shareConsumeRequestManager.hasCompletedFetches());
 
@@ -2664,6 +2677,323 @@ public class ShareConsumeRequestManagerTest {
         assertEquals(1, fetchedRecords.size());
     }
 
+    /**
+     * When a broker is fenced, it is dropped from the {@code brokers} list of 
subsequent MetadataResponses,
+     * so {@code Cluster.nodeById()} returns null for it, while leadership for 
its partitions fails over to
+     * another live broker. This test checks that once the cached leader node 
disappears from the metadata
+     * and leadership moves to a still-present node, the next fetch is routed 
to the new leader.
+     */
+    @Test
+    public void testFetchRecoversWhenCachedLeaderNodeDisappearsFromMetadata() {
+        buildRequestManager();
+
+        subscriptions.subscribeToShareGroup(Set.of(topicName));
+        subscriptions.assignFromSubscribed(List.of(tp0));
+
+        // Two-node cluster; tp0 is led by node 0, so the share session and 
the cached leader are on node 0.
+        client.updateMetadata(
+                RequestTestUtils.metadataUpdateWithIds(2, Map.of(topicName, 2),
+                        tp -> validLeaderEpoch, topicIds, false));
+        Node nodeId0 = metadata.fetch().nodeById(0);
+        Node nodeId1 = metadata.fetch().nodeById(1);
+        assertEquals(nodeId0, metadata.fetch().leaderFor(tp0));
+
+        // Establish the share session and cache the leader (node 0) for tp0.
+        assertEquals(1, sendFetches());
+        client.prepareResponseFrom(
+                ShareFetchResponse.of(Errors.NONE, 0,
+                        buildPartitionDataMap(tip0, records, acquiredRecords, 
Errors.NONE, Errors.NONE),
+                        List.of(), 0),
+                nodeId0);
+        networkClientDelegate.poll(time.timer(0));
+        assertTrue(shareConsumeRequestManager.hasCompletedFetches());
+        fetchRecords();
+
+        // Node 0 is fenced: it disappears from the brokers list of the next 
MetadataResponse and leadership
+        // for tp0 fails over to node 1 with a higher leader epoch.
+        MetadataResponse.PartitionMetadata tp0Metadata = new 
MetadataResponse.PartitionMetadata(
+                Errors.NONE, tp0, Optional.of(nodeId1.id()), 
Optional.of(validLeaderEpoch + 1),
+                List.of(nodeId1.id()), List.of(nodeId1.id()), List.of());
+        MetadataResponse.TopicMetadata topicMetadata = new 
MetadataResponse.TopicMetadata(
+                Errors.NONE, topicName, topicId, false, List.of(tp0Metadata),
+                MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED);
+        MetadataResponse metadataWithoutNode0 = 
RequestTestUtils.metadataResponse(
+                List.of(nodeId1), "kafka-cluster", 1, List.of(topicMetadata));
+        metadata.updateWithCurrentRequestVersion(metadataWithoutNode0, false, 
time.milliseconds());
+
+        // Node 0 is gone from the cluster; node 1 is now the leader for tp0.
+        assertNull(metadata.fetch().nodeById(0));
+        assertEquals(nodeId1, metadata.fetch().leaderFor(tp0));
+
+        // The next fetch should be routed to the new leader (node 1), not 
silently skipped.
+        NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
+        assertEquals(1, pollResult.unsentRequests.size(),
+                "Expected a fetch to the new leader after the cached leader 
node disappeared from metadata");
+        assertEquals(nodeId1, pollResult.unsentRequests.get(0).node().get());
+        ShareFetchRequest.Builder builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+        assertEquals(1, builder.data().topics().size());
+        assertEquals(tip0.topicId(), 
builder.data().topics().stream().findFirst().get().topicId());
+    }
+
+    /**
+     * When the node backing a share session disappears from the cluster 
metadata, the stale session handler
+     * must be removed and any acknowledgements that were queued to be 
piggybacked on a fetch to that node
+     * must be failed with NETWORK_EXCEPTION.
+     */
+    @Test
+    public void 
testSessionHandlerRemovedAndPiggybackAcksFailedWhenNodeDisappears() {
+        buildRequestManager();
+        
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
+
+        subscriptions.subscribeToShareGroup(Set.of(topicName));
+        subscriptions.assignFromSubscribed(List.of(tp0));
+
+        // Two-node cluster; tp0 is led by node 0, so the share session and 
the cached leader are on node 0.
+        client.updateMetadata(
+                RequestTestUtils.metadataUpdateWithIds(2, Map.of(topicName, 2),
+                        tp -> validLeaderEpoch, topicIds, false));
+        Node nodeId0 = metadata.fetch().nodeById(0);
+        Node nodeId1 = metadata.fetch().nodeById(1);
+        assertEquals(nodeId0, metadata.fetch().leaderFor(tp0));
+
+        // Establish the share session on node 0 and fetch records.
+        assertEquals(1, sendFetches());
+        client.prepareResponseFrom(
+                ShareFetchResponse.of(Errors.NONE, 0,
+                        buildPartitionDataMap(tip0, records, acquiredRecords, 
Errors.NONE, Errors.NONE),
+                        List.of(), 0),
+                nodeId0);
+        networkClientDelegate.poll(time.timer(0));
+        assertTrue(shareConsumeRequestManager.hasCompletedFetches());
+        fetchRecords();
+        assertNotNull(shareConsumeRequestManager.sessionHandler(nodeId0.id()));
+
+        // Queue acknowledgements to be piggybacked on the next fetch to node 
0.
+        Acknowledgements acknowledgements = getAcknowledgements(1,
+                AcknowledgeType.ACCEPT, AcknowledgeType.ACCEPT, 
AcknowledgeType.REJECT);
+        shareConsumeRequestManager.fetch(Map.of(tip0, new 
NodeAcknowledgements(0, acknowledgements)));
+
+        // Node 0 is fenced: it disappears from the metadata and leadership 
for tp0 fails over to node 1.
+        MetadataResponse.PartitionMetadata tp0Metadata = new 
MetadataResponse.PartitionMetadata(
+                Errors.NONE, tp0, Optional.of(nodeId1.id()), 
Optional.of(validLeaderEpoch + 1),
+                List.of(nodeId1.id()), List.of(nodeId1.id()), List.of());
+        MetadataResponse.TopicMetadata topicMetadata = new 
MetadataResponse.TopicMetadata(
+                Errors.NONE, topicName, topicId, false, List.of(tp0Metadata),
+                MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED);
+        MetadataResponse metadataWithoutNode0 = 
RequestTestUtils.metadataResponse(
+                List.of(nodeId1), "kafka-cluster", 1, List.of(topicMetadata));
+        metadata.updateWithCurrentRequestVersion(metadataWithoutNode0, false, 
time.milliseconds());
+        assertNull(metadata.fetch().nodeById(0));
+
+        // The next poll re-routes the fetch to node 1, removes the stale 
session handler for node 0, and fails
+        // the piggyback acknowledgements that could no longer be sent to node 
0.
+        assertEquals(1, sendFetches());
+        assertNull(shareConsumeRequestManager.sessionHandler(nodeId0.id()),
+                "Session handler for the disappeared node should have been 
removed");
+        assertNotNull(shareConsumeRequestManager.sessionHandler(nodeId1.id()));
+
+        assertEquals(1, completedAcknowledgements.size());
+        assertEquals(acknowledgements, 
completedAcknowledgements.get(0).get(tip0));
+        assertInstanceOf(NetworkException.class, 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+    }
+
+    /**
+     * An acknowledge request state is created for the node that led the 
partition at the time of the commit.
+     * If that node then disappears from the cluster metadata before the 
request is sent, the acknowledgements must be failed with
+     * NOT_LEADER_OR_FOLLOWER.
+     */
+    @Test
+    public void testAcknowledgeRequestFailedWhenNodeDisappearsBeforeSend() {
+        buildRequestManager();
+        
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
+
+        subscriptions.subscribeToShareGroup(Set.of(topicName));
+        subscriptions.assignFromSubscribed(List.of(tp0));
+
+        client.updateMetadata(
+                RequestTestUtils.metadataUpdateWithIds(2, Map.of(topicName, 2),
+                        tp -> validLeaderEpoch, topicIds, false));
+        Node nodeId0 = metadata.fetch().nodeById(0);
+        Node nodeId1 = metadata.fetch().nodeById(1);
+
+        // Establish the share session on node 0 and fetch records.
+        assertEquals(1, sendFetches());
+        client.prepareResponseFrom(
+                ShareFetchResponse.of(Errors.NONE, 0,
+                        buildPartitionDataMap(tip0, records, acquiredRecords, 
Errors.NONE, Errors.NONE),
+                        List.of(), 0),
+                nodeId0);
+        networkClientDelegate.poll(time.timer(0));
+        fetchRecords();
+
+        // commitSync enqueues an acknowledge request state for node 0, which 
is still the leader at this point.
+        Acknowledgements acknowledgements = getAcknowledgements(1,
+                AcknowledgeType.ACCEPT, AcknowledgeType.ACCEPT, 
AcknowledgeType.REJECT);
+        CompletableFuture<Map<TopicIdPartition, Acknowledgements>> future =
+                shareConsumeRequestManager.commitSync(Map.of(tip0, new 
NodeAcknowledgements(0, acknowledgements)),
+                        calculateDeadlineMs(time.timer(defaultApiTimeoutMs)));
+        assertFalse(future.isDone());
+
+        // Node 0 is fenced: it disappears from the metadata and leadership 
for tp0 fails over to node 1.
+        MetadataResponse.PartitionMetadata tp0Metadata = new 
MetadataResponse.PartitionMetadata(
+                Errors.NONE, tp0, Optional.of(nodeId1.id()), 
Optional.of(validLeaderEpoch + 1),
+                List.of(nodeId1.id()), List.of(nodeId1.id()), List.of());
+        MetadataResponse.TopicMetadata topicMetadata = new 
MetadataResponse.TopicMetadata(
+                Errors.NONE, topicName, topicId, false, List.of(tp0Metadata),
+                MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED);
+        MetadataResponse metadataWithoutNode0 = 
RequestTestUtils.metadataResponse(
+                List.of(nodeId1), "kafka-cluster", 1, List.of(topicMetadata));
+        metadata.updateWithCurrentRequestVersion(metadataWithoutNode0, false, 
time.milliseconds());
+        assertNull(metadata.fetch().nodeById(0));
+
+        // No request is sent to the vanished node, and the commit completes 
with NOT_LEADER_OR_FOLLOWER
+        // rather than hanging.
+        assertEquals(0, shareConsumeRequestManager.sendAcknowledgements());
+        assertTrue(future.isDone());
+        assertEquals(1, completedAcknowledgements.size());
+        assertEquals(acknowledgements, 
completedAcknowledgements.get(0).get(tip0));
+        assertInstanceOf(NotLeaderOrFollowerException.class, 
completedAcknowledgements.get(0).get(tip0).getAcknowledgeException());
+    }
+
+    /**
+     * When a topic is deleted and recreated, its topic ID changes. The 
request manager caches the mapping
+     * from topic-partition to topic ID (and leader). A partition-level 
UNKNOWN_TOPIC_ID error in a ShareFetch
+     * response must clear that cached mapping so that, once the metadata 
reflects the recreated topic, the stale
+     * topic ID is forgotten from the share session and the new one is fetched.
+     */
+    @Test
+    public void 
testFetchResponseWithUnknownTopicIdRefreshesTopicIdOnRecreation() {
+        buildRequestManager();
+
+        assignFromSubscribed(Set.of(tp0));
+
+        // Establish the share session and cache the topic ID and leader for 
tp0.
+        sendFetchAndVerifyResponse(records, acquiredRecords, Errors.NONE);
+        fetchRecords();
+
+        // A subsequent fetch returns a partition-level UNKNOWN_TOPIC_ID 
error, which causes the request
+        // manager to forget the cached topic ID and leader for the partition.
+        assertEquals(1, sendFetches());
+        client.prepareResponse(fullFetchResponse(tip0, records, 
emptyAcquiredRecords, Errors.UNKNOWN_TOPIC_ID));
+        networkClientDelegate.poll(time.timer(0));
+        assertTrue(shareConsumeRequestManager.hasCompletedFetches());
+        fetchRecords();
+
+        // The topic is recreated with a new topic ID.
+        Uuid recreatedTopicId = Uuid.randomUuid();
+        client.updateMetadata(
+                RequestTestUtils.metadataUpdateWithIds(1, Map.of(topicName, 2),
+                        tp -> validLeaderEpoch, Map.of(topicName, 
recreatedTopicId), false));
+
+        // The next ShareFetch fetches the recreated topic ID and forgets the 
stale one from the share session.
+        NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
+        assertEquals(1, pollResult.unsentRequests.size());
+        ShareFetchRequest.Builder builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+
+        // The recreated topic ID is fetched.
+        assertEquals(1, builder.data().topics().size());
+        assertEquals(recreatedTopicId, 
builder.data().topics().stream().findFirst().get().topicId());
+
+        // The stale topic ID is forgotten.
+        assertEquals(1, builder.data().forgottenTopicsData().size());
+        assertEquals(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));
+    }
+
+    /**
+     * As {@link 
#testFetchResponseWithUnknownTopicIdRefreshesTopicIdOnRecreation()} but the 
UNKNOWN_TOPIC_ID
+     * error arrives in a ShareAcknowledge response rather than a ShareFetch 
response.
+     */
+    @Test
+    public void 
testAcknowledgeResponseWithUnknownTopicIdRefreshesTopicIdOnRecreation() {
+        buildRequestManager();
+        
shareConsumeRequestManager.setAcknowledgementCommitCallbackRegistered(true);
+
+        assignFromSubscribed(Set.of(tp0));
+
+        // Establish the share session and cache the topic ID and leader for 
tp0.
+        sendFetchAndVerifyResponse(records, acquiredRecords, Errors.NONE);
+        fetchRecords();
+
+        // Acknowledge some records and receive a partition-level 
UNKNOWN_TOPIC_ID error, which causes the request
+        // manager to forget the cached topic ID and leader for the partition.
+        Acknowledgements acknowledgements = getAcknowledgements(1, 
AcknowledgeType.ACCEPT, AcknowledgeType.ACCEPT, AcknowledgeType.REJECT);
+        shareConsumeRequestManager.commitAsync(Map.of(tip0, new 
NodeAcknowledgements(0, acknowledgements)),
+                calculateDeadlineMs(time.timer(defaultApiTimeoutMs)));
+
+        assertEquals(1, shareConsumeRequestManager.sendAcknowledgements());
+        client.prepareResponse(fullAcknowledgeResponse(tip0, 
Errors.UNKNOWN_TOPIC_ID));
+        networkClientDelegate.poll(time.timer(0));
+
+        // The acknowledgements are completed with the error rather than 
retried.
+        assertEquals(3, completedAcknowledgements.get(0).get(tip0).size());
+
+        // The topic is recreated with a new topic ID.
+        Uuid recreatedTopicId = Uuid.randomUuid();
+        client.updateMetadata(
+                RequestTestUtils.metadataUpdateWithIds(1, Map.of(topicName, 2),
+                        tp -> validLeaderEpoch, Map.of(topicName, 
recreatedTopicId), false));
+
+        // The next ShareFetch fetches the recreated topic ID and forgets the 
stale one from the share session.
+        NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
+        assertEquals(1, pollResult.unsentRequests.size());
+        ShareFetchRequest.Builder builder = (ShareFetchRequest.Builder) 
pollResult.unsentRequests.get(0).requestBuilder();
+
+        assertEquals(1, builder.data().topics().size());
+        assertEquals(recreatedTopicId, 
builder.data().topics().stream().findFirst().get().topicId());
+
+        assertEquals(1, builder.data().forgottenTopicsData().size());
+        assertEquals(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 cached leader for a partition in the share session should only be 
replaced when the ShareFetch response
+     * carries a newer leader epoch, following the same rules as {@link 
Metadata#updateLastSeenEpochIfNewer}. A stale
+     * leader with an older epoch must not overwrite the cached leader, so 
subsequent fetches continue to go to the
+     * node with the newest known leader epoch.
+     */
+    @Test
+    public void testStaleLeaderEpochDoesNotDowngradeShareSessionLeader() {
+        buildRequestManager();
+
+        subscriptions.subscribeToShareGroup(Set.of(topicName));
+        subscriptions.assignFromSubscribed(Set.of(tp0));
+
+        // tp0's leader is node0 with a high leader epoch.
+        client.updateMetadata(
+                RequestTestUtils.metadataUpdateWithIds(2, Map.of(topicName, 1),
+                        tp -> validLeaderEpoch + 5, topicIds, false));
+        Node nodeId0 = metadata.fetch().nodeById(0);
+        Node nodeId1 = metadata.fetch().nodeById(1);
+        assertEquals(nodeId0, metadata.fetch().leaderFor(tp0));
+
+        // The first fetch goes to node0 and caches the leader (node0, epoch + 
5).
+        assertEquals(1, sendFetches());
+        assertEquals(nodeId0, metadata.fetch().leaderFor(tp0));
+
+        // node0 responds with a leadership error naming node1 as the new 
leader, but with an older leader epoch.
+        LinkedHashMap<TopicIdPartition, ShareFetchResponseData.PartitionData> 
partitionData = new LinkedHashMap<>();
+        partitionData.put(tip0,
+                new ShareFetchResponseData.PartitionData()
+                        .setPartitionIndex(tip0.topicPartition().partition())
+                        .setErrorCode(Errors.NOT_LEADER_OR_FOLLOWER.code())
+                        .setCurrentLeader(new 
ShareFetchResponseData.LeaderIdAndEpoch()
+                                .setLeaderId(nodeId1.id())
+                                .setLeaderEpoch(validLeaderEpoch + 2)));
+        client.prepareResponseFrom(ShareFetchResponse.of(Errors.NONE, 0, 
partitionData, List.of(nodeId1), 0), nodeId0);
+        networkClientDelegate.poll(time.timer(0));
+        assertTrue(shareConsumeRequestManager.hasCompletedFetches());
+        fetchRecords();
+
+        // The stale leader must not have replaced the cached leader, so the 
next fetch still goes only to node0.
+        NetworkClientDelegate.PollResult pollResult = 
shareConsumeRequestManager.sendFetchesReturnPollResult();
+        assertEquals(1, pollResult.unsentRequests.size());
+        assertEquals(nodeId0, pollResult.unsentRequests.get(0).node().get());
+    }
+
     @Test
     public void testFetchOneNodeAtATimeForRecordLimitMode() {
         // We will simulate two nodes, each with one partition. The first node 
will have more records
@@ -2749,7 +3079,7 @@ public class ShareConsumeRequestManagerTest {
     }
 
     @Test
-    public void testCloseInternalClosesShareFetchMetricsManager() throws 
Exception {
+    public void testCloseInternalClosesShareFetchMetricsManager() {
         buildRequestManager();
 
         // Define all sensor names that should be created and removed
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandlerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandlerTest.java
index f81a53f9204..405adff250e 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandlerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareSessionHandlerTest.java
@@ -538,6 +538,87 @@ public class ShareSessionHandlerTest {
         assertNull(builder);
     }
 
+    @Test
+    public void testSkipEmptyShareFetchWhenSessionEmpty() {
+        // Using the default (non record_limit) config, canSkipIfRequestEmpty 
is false.
+        ShareFetchConfig shareFetchConfig = DEFAULT_SHARE_FETCH_CONFIG;
+
+        String groupId = "G1";
+        Uuid memberId = Uuid.randomUuid();
+        ShareSessionHandler handler = new ShareSessionHandler(LOG_CONTEXT, 1, 
memberId);
+
+        Map<Uuid, String> topicNames = new HashMap<>();
+        Uuid fooId = addTopicId(topicNames, "foo");
+        TopicIdPartition foo0 = new TopicIdPartition(fooId, 0, "foo");
+
+        // Adding a partition to the session builds a request.
+        handler.addPartitionToFetch(foo0, null);
+        assertNotNull(handler.newShareFetchBuilder(groupId, shareFetchConfig, 
false));
+
+        ShareFetchResponse resp = ShareFetchResponse.of(Errors.NONE,
+            0,
+            buildResponseData(new RespEntry("foo", 0, fooId)),
+            List.of(),
+            0);
+        handler.handleResponse(resp, ApiKeys.SHARE_FETCH.latestVersion());
+
+        // A non-empty session with no changes still builds a request in the 
default mode, since we want to fetch records.
+        handler.addPartitionToFetch(foo0, null);
+        assertNotNull(handler.newShareFetchBuilder(groupId, shareFetchConfig, 
false));
+        handler.handleResponse(resp, ApiKeys.SHARE_FETCH.latestVersion());
+
+        // Removing the only partition from the session builds a request to 
forget it.
+        assertNotNull(handler.newShareFetchBuilder(groupId, shareFetchConfig, 
false));
+        handler.handleResponse(ShareFetchResponse.of(Errors.NONE, 0, new 
LinkedHashMap<>(), List.of(), 0), ApiKeys.SHARE_FETCH.latestVersion());
+
+        // Once the session is empty, no request is built even though 
canSkipIfRequestEmpty is false.
+        assertTrue(handler.sessionPartitionMap().isEmpty());
+        assertNull(handler.newShareFetchBuilder(groupId, shareFetchConfig, 
false));
+    }
+
+    @ParameterizedTest
+    @MethodSource("shareFetchConfigProvider")
+    public void 
testAcknowledgeOnlyPartitionNotInSessionIsForgotten(ShareFetchConfig 
shareFetchConfig) {
+        String groupId = "G1";
+        Uuid memberId = Uuid.randomUuid();
+        ShareSessionHandler handler = new ShareSessionHandler(LOG_CONTEXT, 1, 
memberId);
+
+        Map<Uuid, String> topicNames = new HashMap<>();
+        Uuid fooId = addTopicId(topicNames, "foo");
+        Uuid barId = addTopicId(topicNames, "bar");
+        TopicIdPartition foo0 = new TopicIdPartition(fooId, 0, "foo");
+        TopicIdPartition bar0 = new TopicIdPartition(barId, 0, "bar");
+
+        // Establish the session with foo0.
+        handler.addPartitionToFetch(foo0, null);
+        handler.newShareFetchBuilder(groupId, shareFetchConfig, false);
+        ShareFetchResponse resp = ShareFetchResponse.of(Errors.NONE,
+            0,
+            buildResponseData(new RespEntry("foo", 0, fooId), new 
RespEntry("bar", 0, barId)),
+            List.of(),
+            0);
+        handler.handleResponse(resp, ApiKeys.SHARE_FETCH.latestVersion());
+
+        // Continue fetching foo0, and send acknowledgements for bar0, which 
is not part of the session.
+        Acknowledgements acknowledgements = Acknowledgements.empty();
+        acknowledgements.add(0L, AcknowledgeType.ACCEPT);
+        handler.addPartitionToFetch(foo0, null);
+        handler.addPartitionToAcknowledgeOnly(bar0, acknowledgements);
+        ShareFetchRequestData requestData = 
handler.newShareFetchBuilder(groupId, shareFetchConfig, false).build().data();
+
+        // bar0 is not added to the session; it is placed in the forgotten 
list so that it is removed from the
+        // share session on the broker once its acknowledgements have been 
sent.
+        assertMapsEqual(reqMap(foo0), handler.sessionPartitionMap());
+        assertEquals(List.of(bar0), reqForgetList(requestData, topicNames));
+
+        // The acknowledgements for bar0 are still included in the request.
+        ShareFetchRequestData.FetchPartition barPartition = 
requestData.topics().stream()
+                .filter(topic -> topic.topicId().equals(barId))
+                .flatMap(topic -> topic.partitions().stream())
+                .findFirst().get();
+        assertEquals(1, barPartition.acknowledgementBatches().size());
+    }
+
     private Uuid addTopicId(Map<Uuid, String> topicNames, String name) {
         Uuid id = Uuid.randomUuid();
         topicNames.put(id, name);

Reply via email to