Skip to content
11 changes: 8 additions & 3 deletions raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -1055,14 +1055,14 @@ private boolean handleFetchResponse(
partitionResponse.snapshotId().endOffset(),
partitionResponse.snapshotId().epoch()
);
return false;
return true;
Comment thread
jsancio marked this conversation as resolved.
Outdated
} else if (partitionResponse.snapshotId().endOffset() < 0) {
logger.error(
"The leader sent a snapshot id with a valid epoch {} but with an invalid end offset {}",
partitionResponse.snapshotId().epoch(),
partitionResponse.snapshotId().endOffset()
);
return false;
return true;
} else {
OffsetAndEpoch snapshotId = new OffsetAndEpoch(
partitionResponse.snapshotId().endOffset(),
Expand Down Expand Up @@ -1274,7 +1274,12 @@ private boolean handleFetchSnapshotResponse(
/* The leader deleted the snapshot before the follower could download it. Start over by
Comment thread
jsancio marked this conversation as resolved.
* reseting the fetching snapshot state and sending another fetch request.
*/
logger.trace("Leader doesn't know about snapshot id {}, returned error {} and snapshot id {}", state.fetchingSnapshot(), partitionSnapshot.errorCode(), partitionSnapshot.snapshotId());
logger.trace(
"Leader doesn't know about snapshot id {}, returned error {} and snapshot id {}",
state.fetchingSnapshot(),
partitionSnapshot.errorCode(),
partitionSnapshot.snapshotId()
);
state.setFetchingSnapshot(Optional.empty());
state.resetFetchTimeout(currentTimeMs);
return true;
Expand Down
Loading