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

apoorvmittal10 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 1189910a17a KAFKA-20610: Minor - Renaming variables to not include 
replica manager reference (5/N) (#22602)
1189910a17a is described below

commit 1189910a17a7c8af36bc8e1d9633a835c2001fde
Author: Apoorv Mittal <[email protected]>
AuthorDate: Wed Jun 17 15:28:29 2026 +0100

    KAFKA-20610: Minor - Renaming variables to not include replica manager 
reference (5/N) (#22602)
    
    Minor refactoring to streamline variable names for fetched data from
    log.
    
    Reviewers: Andrew Schofield <[email protected]>
---
 .../java/kafka/server/share/DelayedShareFetch.java | 40 +++++++++++-----------
 1 file changed, 20 insertions(+), 20 deletions(-)

diff --git a/core/src/main/java/kafka/server/share/DelayedShareFetch.java 
b/core/src/main/java/kafka/server/share/DelayedShareFetch.java
index a770da515e8..0ac16bea2e5 100644
--- a/core/src/main/java/kafka/server/share/DelayedShareFetch.java
+++ b/core/src/main/java/kafka/server/share/DelayedShareFetch.java
@@ -315,13 +315,13 @@ public class DelayedShareFetch extends DelayedOperation {
      * Hence, we require to set offsetMetadata to null for this fetch offset, 
which would cause tryComplete to update
      * fetchOffsetMetadata and thereby we will identify this partition for 
remote storage fetch.
      * @param topicPartitionData - Map containing the fetch offset for the 
topic partitions.
-     * @param replicaManagerReadResponse - Map containing the readFromLog 
response from replicaManager for the topic partitions.
+     * @param readResponse - Map containing the readFromLog response for the 
topic partitions.
      */
     private void resetFetchOffsetMetadataForRemoteFetchPartitions(
         LinkedHashMap<TopicIdPartition, Long> topicPartitionData,
-        LinkedHashMap<TopicIdPartition, LogReadResult> 
replicaManagerReadResponse
+        LinkedHashMap<TopicIdPartition, LogReadResult> readResponse
     ) {
-        replicaManagerReadResponse.forEach((topicIdPartition, logReadResult) 
-> {
+        readResponse.forEach((topicIdPartition, logReadResult) -> {
             if (logReadResult.info().delayedRemoteStorageFetch.isPresent()) {
                 SharePartition sharePartition = 
sharePartitions.get(topicIdPartition);
                 sharePartition.updateFetchOffsetMetadata(
@@ -348,19 +348,19 @@ public class DelayedShareFetch extends DelayedOperation {
                 // Update the metric to record the time taken to acquire the 
locks for the share partitions.
                 updateAcquireElapsedTimeMetric();
                 // In case, fetch offset metadata doesn't exist for one or 
more topic partitions, we do a
-                // replicaManager.readFromLog to populate the offset metadata 
and update the fetch offset metadata for
+                // readFromLog to populate the offset metadata and update the 
fetch offset metadata for
                 // those topic partitions.
-                LinkedHashMap<TopicIdPartition, LogReadResult> 
replicaManagerReadResponse = maybeReadFromLog(topicPartitionData);
+                LinkedHashMap<TopicIdPartition, LogReadResult> readResponse = 
maybeReadFromLog(topicPartitionData);
                 // Store the remote fetch info for the topic partitions for 
which we need to perform remote fetch.
-                LinkedHashMap<TopicIdPartition, LogReadResult> 
remoteStorageFetchInfoMap = 
maybePrepareRemoteStorageFetchInfo(topicPartitionData, 
replicaManagerReadResponse);
+                LinkedHashMap<TopicIdPartition, LogReadResult> 
remoteStorageFetchInfoMap = 
maybePrepareRemoteStorageFetchInfo(topicPartitionData, readResponse);
 
                 if (!remoteStorageFetchInfoMap.isEmpty()) {
                     return maybeProcessRemoteFetch(topicPartitionData, 
remoteStorageFetchInfoMap);
                 }
-                maybeUpdateFetchOffsetMetadata(topicPartitionData, 
replicaManagerReadResponse);
-                if (anyPartitionHasLogReadError(replicaManagerReadResponse) || 
isMinBytesSatisfied(topicPartitionData, 
partitionMaxBytesStrategy.maxBytes(shareFetch.fetchParams().maxBytes, 
topicPartitionData.keySet(), topicPartitionData.size()))) {
+                maybeUpdateFetchOffsetMetadata(topicPartitionData, 
readResponse);
+                if (anyPartitionHasLogReadError(readResponse) || 
isMinBytesSatisfied(topicPartitionData, 
partitionMaxBytesStrategy.maxBytes(shareFetch.fetchParams().maxBytes, 
topicPartitionData.keySet(), topicPartitionData.size()))) {
                     partitionsAcquired = topicPartitionData;
-                    localPartitionsAlreadyFetched = replicaManagerReadResponse;
+                    localPartitionsAlreadyFetched = readResponse;
                     return forceComplete();
                 } else {
                     log.debug("minBytes is not satisfied for the share fetch 
request for group {}, member {}, " +
@@ -453,19 +453,19 @@ public class DelayedShareFetch extends DelayedOperation {
     }
 
     private void 
maybeUpdateFetchOffsetMetadata(LinkedHashMap<TopicIdPartition, Long> 
topicPartitionData,
-                                                
LinkedHashMap<TopicIdPartition, LogReadResult> replicaManagerReadResponseData) {
-        for (Map.Entry<TopicIdPartition, LogReadResult> entry : 
replicaManagerReadResponseData.entrySet()) {
+                                                
LinkedHashMap<TopicIdPartition, LogReadResult> readResponseData) {
+        for (Map.Entry<TopicIdPartition, LogReadResult> entry : 
readResponseData.entrySet()) {
             TopicIdPartition topicIdPartition = entry.getKey();
             SharePartition sharePartition = 
sharePartitions.get(topicIdPartition);
-            LogReadResult replicaManagerLogReadResult = entry.getValue();
-            if (replicaManagerLogReadResult.error().code() != 
Errors.NONE.code()) {
-                log.debug("Replica manager read log result {} errored out for 
topic partition {}",
-                    replicaManagerLogReadResult, topicIdPartition);
+            LogReadResult logReadResult = entry.getValue();
+            if (logReadResult.error().code() != Errors.NONE.code()) {
+                log.debug("Log read result {} errored out for topic partition 
{}",
+                    logReadResult, topicIdPartition);
                 continue;
             }
             sharePartition.updateFetchOffsetMetadata(
                 topicPartitionData.get(topicIdPartition),
-                replicaManagerLogReadResult.info().fetchOffsetMetadata);
+                logReadResult.info().fetchOffsetMetadata);
         }
     }
 
@@ -530,8 +530,8 @@ public class DelayedShareFetch extends DelayedOperation {
         return logReader.read(shareFetch.fetchParams(), partitionsToFetch, 
topicPartitionFetchOffsets, partitionMaxBytes);
     }
 
-    private boolean 
anyPartitionHasLogReadError(LinkedHashMap<TopicIdPartition, LogReadResult> 
replicaManagerReadResponse) {
-        return replicaManagerReadResponse.values().stream()
+    private boolean 
anyPartitionHasLogReadError(LinkedHashMap<TopicIdPartition, LogReadResult> 
readResponse) {
+        return readResponse.values().stream()
             .anyMatch(logReadResult -> logReadResult.error().code() != 
Errors.NONE.code());
     }
 
@@ -628,10 +628,10 @@ public class DelayedShareFetch extends DelayedOperation {
 
     private LinkedHashMap<TopicIdPartition, LogReadResult> 
maybePrepareRemoteStorageFetchInfo(
         LinkedHashMap<TopicIdPartition, Long> topicPartitionData,
-        LinkedHashMap<TopicIdPartition, LogReadResult> 
replicaManagerReadResponse
+        LinkedHashMap<TopicIdPartition, LogReadResult> readResponse
     ) {
         LinkedHashMap<TopicIdPartition, LogReadResult> 
remoteStorageFetchInfoMap = new LinkedHashMap<>();
-        for (Map.Entry<TopicIdPartition, LogReadResult> entry : 
replicaManagerReadResponse.entrySet()) {
+        for (Map.Entry<TopicIdPartition, LogReadResult> entry : 
readResponse.entrySet()) {
             TopicIdPartition topicIdPartition = entry.getKey();
             LogReadResult logReadResult = entry.getValue();
             if (logReadResult.info().delayedRemoteStorageFetch.isPresent()) {

Reply via email to