Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 12 additions & 12 deletions core/src/main/scala/kafka/raft/KafkaMetadataLog.scala
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ package kafka.raft
import java.io.{File, IOException}
import java.nio.file.{Files, NoSuchFileException}
import java.util.concurrent.ConcurrentSkipListSet
import java.util.{NoSuchElementException, Optional, Properties}
import java.util.{Optional, Properties}

import kafka.api.ApiVersion
import kafka.log.{AppendOrigin, Log, LogConfig, LogOffsetSnapshot, SnapshotGenerated}
Expand Down Expand Up @@ -248,20 +248,20 @@ final class KafkaMetadataLog private (
}

override def latestSnapshotId(): Optional[OffsetAndEpoch] = {
try {
Optional.of(snapshotIds.last)
} catch {
case _: NoSuchElementException =>
Optional.empty()
val descending = snapshotIds.descendingIterator
if (descending.hasNext) {
Optional.of(descending.next)
} else {
Optional.empty()
}
}

override def earliestSnapshotId(): Optional[OffsetAndEpoch] = {
try {
Optional.of(snapshotIds.first)
} catch {
case _: NoSuchElementException =>
Optional.empty()
val ascendingIterator = snapshotIds.iterator
if (ascendingIterator.hasNext) {
Optional.of(ascendingIterator.next)
} else {
Optional.empty()
}
}

Expand Down Expand Up @@ -290,7 +290,7 @@ final class KafkaMetadataLog private (
* Removes all snapshots on the log directory whose epoch and end offset is less than the giving epoch and end offset.
*/
private def removeSnapshotFilesBefore(logStartSnapshotId: OffsetAndEpoch): Unit = {
val expiredSnapshotIdsIter = snapshotIds.headSet(logStartSnapshotId, false).iterator()
val expiredSnapshotIdsIter = snapshotIds.headSet(logStartSnapshotId, false).iterator
while (expiredSnapshotIdsIter.hasNext) {
val snapshotId = expiredSnapshotIdsIter.next()
// If snapshotIds contains a snapshot id, the KafkaRaftClient and Listener can expect that the snapshot exists
Expand Down
12 changes: 7 additions & 5 deletions raft/src/main/java/org/apache/kafka/raft/ReplicatedLog.java
Original file line number Diff line number Diff line change
Expand Up @@ -84,11 +84,13 @@ public interface ReplicatedLog extends Closeable {
default ValidOffsetAndEpoch validateOffsetAndEpoch(long offset, int epoch) {
if (startOffset() == 0 && offset == 0) {
return ValidOffsetAndEpoch.valid(new OffsetAndEpoch(0, 0));
} else if (
earliestSnapshotId().isPresent() &&
((offset < startOffset()) ||
(offset == startOffset() && epoch != earliestSnapshotId().get().epoch) ||
(epoch < earliestSnapshotId().get().epoch))
}

Optional<OffsetAndEpoch> earliestSnapshotId = earliestSnapshotId();
if (earliestSnapshotId.isPresent() &&
((offset < startOffset()) ||
(offset == startOffset() && epoch != earliestSnapshotId.get().epoch) ||
(epoch < earliestSnapshotId.get().epoch))
) {
/* Send a snapshot if the leader has a snapshot at the log start offset and
* 1. the fetch offset is less than the log start offset or
Expand Down