From df79e31e3abaee3379c12ec197c9f6a54eb75957 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Sat, 29 Jul 2023 19:03:54 +0530 Subject: [PATCH 01/11] KAFKA-15272: Fix the logic which finds candidate log segments to upload it to tiered storage --- .../kafka/log/remote/RemoteLogManager.java | 42 ++++++------------ .../log/remote/RemoteLogManagerTest.java | 44 ++++++++++++++++--- 2 files changed, 52 insertions(+), 34 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index ea065b8c8126f..a4b9b146bccfb 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -84,12 +84,10 @@ import java.security.PrivilegedAction; import java.util.ArrayList; import java.util.Collections; -import java.util.Comparator; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; import java.util.List; -import java.util.ListIterator; import java.util.Map; import java.util.Optional; import java.util.OptionalInt; @@ -538,35 +536,33 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException if (lso < 0) { logger.warn("lastStableOffset for partition {} is {}, which should not be negative.", topicIdPartition, lso); } else if (lso > 0 && copiedOffset < lso) { - // Copy segments only till the last-stable-offset as remote storage should contain only committed/acked - // messages - long toOffset = lso; - logger.debug("Checking for segments to copy, copiedOffset: {} and toOffset: {}", copiedOffset, toOffset); long activeSegBaseOffset = log.activeSegment().baseOffset(); // log-start-offset can be ahead of the read-offset, when: // 1) log-start-offset gets incremented via delete-records API (or) // 2) enabling the remote log for the first time long fromOffset = Math.max(copiedOffset + 1, log.logStartOffset()); - ArrayList sortedSegments = new ArrayList<>(JavaConverters.asJavaCollection(log.logSegments(fromOffset, toOffset))); - sortedSegments.sort(Comparator.comparingLong(LogSegment::baseOffset)); - List sortedBaseOffsets = sortedSegments.stream().map(LogSegment::baseOffset).collect(Collectors.toList()); - int activeSegIndex = Collections.binarySearch(sortedBaseOffsets, activeSegBaseOffset); - // sortedSegments becomes empty list when fromOffset and toOffset are same, and activeSegIndex becomes -1 - if (activeSegIndex < 0) { + // Segments which match the following criteria are eligible for copying to remote storage: + // 1) Segment is not the active segment and + // 2) Segment end-offset is less than the last-stable-offset as remote storage should contain only + // committed/acked messages + List candidateSegments = JavaConverters.asJavaCollection(log.logSegments(fromOffset, Long.MAX_VALUE)) + .stream() + .filter(segment -> segment.baseOffset() != activeSegBaseOffset && segment.readNextOffset() <= lso) + .collect(Collectors.toList()); + logger.debug("Checking for segments to copy, copiedOffset: {} and lso: {}, candidateSegments: {}", + copiedOffset, lso, candidateSegments); + if (candidateSegments.isEmpty()) { logger.debug("No segments found to be copied for partition {} with copiedOffset: {} and active segment's base-offset: {}", topicIdPartition, copiedOffset, activeSegBaseOffset); } else { - ListIterator logSegmentsIter = sortedSegments.subList(0, activeSegIndex).listIterator(); - while (logSegmentsIter.hasNext()) { - LogSegment segment = logSegmentsIter.next(); + for (LogSegment segment : candidateSegments) { if (isCancelled() || !isLeader()) { logger.info("Skipping copying log segments as the current task state is changed, cancelled: {} leader:{}", isCancelled(), isLeader()); return; } - - copyLogSegment(log, segment, getNextSegmentBaseOffset(activeSegBaseOffset, logSegmentsIter)); + copyLogSegment(log, segment, segment.readNextOffset()); } } } else { @@ -583,18 +579,6 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException } } - private long getNextSegmentBaseOffset(long activeSegBaseOffset, ListIterator logSegmentsIter) { - long nextSegmentBaseOffset; - if (logSegmentsIter.hasNext()) { - nextSegmentBaseOffset = logSegmentsIter.next().baseOffset(); - logSegmentsIter.previous(); - } else { - nextSegmentBaseOffset = activeSegBaseOffset; - } - - return nextSegmentBaseOffset; - } - private void copyLogSegment(UnifiedLog log, LogSegment segment, long nextSegmentBaseOffset) throws InterruptedException, ExecutionException, RemoteStorageException, IOException { File logFile = segment.log().file(); String logFileName = logFile.getName(); diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index c51a34fe3648d..b9c8e07f80f70 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -269,15 +269,38 @@ void testStartup() { void testCopyLogSegmentsToRemoteShouldCopyExpectedLogSegment() throws Exception { long oldSegmentStartOffset = 0L; long nextSegmentStartOffset = 150L; - long oldSegmentEndOffset = nextSegmentStartOffset - 1; + long lso = 250L; + long leo = 300L; + assertCopyExpectedLogSegmentsToRemote(oldSegmentStartOffset, nextSegmentStartOffset, lso, leo); + } + + /** + * The following values will be equal when the active segment gets rotated to passive when there are no new messages: + * last-stable-offset = high-water-mark = log-end-offset = base-offset-of-active-segment + * + * This test asserts that the active log segment that was rotated after log.roll.ms are copied to remote storage. + */ + @Test + void testCopyLogSegmentToRemoteForStaleTopic() throws Exception { + long oldSegmentStartOffset = 0L; + long nextSegmentStartOffset = 150L; + long lso = 150L; + long leo = 150L; + assertCopyExpectedLogSegmentsToRemote(oldSegmentStartOffset, nextSegmentStartOffset, lso, leo); + } + private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, + long nextSegmentStartOffset, + long lastStableOffset, + long logEndOffset) throws Exception { + long oldSegmentEndOffset = nextSegmentStartOffset - 1; when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); // leader epoch preparation checkpoint.write(totalEpochEntries); LeaderEpochFileCache cache = new LeaderEpochFileCache(leaderTopicIdPartition.topicPartition(), checkpoint); when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); - when(remoteLogMetadataManager.highestOffsetForEpoch(any(TopicIdPartition.class), anyInt())).thenReturn(Optional.of(0L)); + when(remoteLogMetadataManager.highestOffsetForEpoch(any(TopicIdPartition.class), anyInt())).thenReturn(Optional.of(-1L)); File tempFile = TestUtils.tempFile(); File mockProducerSnapshotIndex = TestUtils.tempFile(); @@ -287,7 +310,9 @@ void testCopyLogSegmentsToRemoteShouldCopyExpectedLogSegment() throws Exception LogSegment activeSegment = mock(LogSegment.class); when(oldSegment.baseOffset()).thenReturn(oldSegmentStartOffset); + when(oldSegment.readNextOffset()).thenReturn(nextSegmentStartOffset); when(activeSegment.baseOffset()).thenReturn(nextSegmentStartOffset); + when(activeSegment.readNextOffset()).thenReturn(logEndOffset); FileRecords fileRecords = mock(FileRecords.class); when(oldSegment.log()).thenReturn(fileRecords); @@ -297,12 +322,21 @@ void testCopyLogSegmentsToRemoteShouldCopyExpectedLogSegment() throws Exception when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.logSegments(anyLong(), anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment, activeSegment))); + when(mockLog.logSegments(anyLong(), anyLong())) + .thenAnswer(invocation -> { + Object[] args = invocation.getArguments(); + long fromOffset = (Long) args[0]; // inclusive + long toOffset = (Long) args[1]; // exclusive + return JavaConverters.collectionAsScalaIterable( + Arrays.asList(oldSegment, activeSegment).stream() + .filter(segment -> segment.baseOffset() >= fromOffset && segment.readNextOffset() < toOffset) + .collect(Collectors.toList())); + }); ProducerStateManager mockStateManager = mock(ProducerStateManager.class); when(mockLog.producerStateManager()).thenReturn(mockStateManager); when(mockStateManager.fetchSnapshot(anyLong())).thenReturn(Optional.of(mockProducerSnapshotIndex)); - when(mockLog.lastStableOffset()).thenReturn(250L); + when(mockLog.lastStableOffset()).thenReturn(lastStableOffset); LazyIndex idx = LazyIndex.forOffset(UnifiedLog.offsetIndexFile(tempDir, oldSegmentStartOffset, ""), oldSegmentStartOffset, 1000); LazyIndex timeIdx = LazyIndex.forTime(UnifiedLog.timeIndexFile(tempDir, oldSegmentStartOffset, ""), oldSegmentStartOffset, 1500); @@ -349,7 +383,7 @@ void testCopyLogSegmentsToRemoteShouldCopyExpectedLogSegment() throws Exception assertEquals(remoteLogSegmentMetadataArg.getValue(), remoteLogSegmentMetadataArg2.getValue()); // The old segment should only contain leader epoch [0->0, 1->100] since its offset range is [0, 149] verifyLogSegmentData(logSegmentDataArg.getValue(), idx, timeIdx, txnIndex, tempFile, mockProducerSnapshotIndex, - Arrays.asList(epochEntry0, epochEntry1)); + Arrays.asList(epochEntry0, epochEntry1)); // verify remoteLogMetadataManager did add the expected RemoteLogSegmentMetadataUpdate ArgumentCaptor remoteLogSegmentMetadataUpdateArg = ArgumentCaptor.forClass(RemoteLogSegmentMetadataUpdate.class); From c2cd2de6b70477c2c71920cd15c738df95a3a609 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 10:50:43 +0530 Subject: [PATCH 02/11] Avoid calling expensive LogSegment#readNextOffset() operation. --- .../kafka/log/remote/RemoteLogManager.java | 56 ++++++++++++++-- .../log/remote/RemoteLogManagerTest.java | 66 +++++++++++++++---- 2 files changed, 104 insertions(+), 18 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index a4b9b146bccfb..e8f3539ffaf5f 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -89,6 +89,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Optional; import java.util.OptionalInt; import java.util.OptionalLong; @@ -523,6 +524,26 @@ private void maybeUpdateReadOffset(UnifiedLog log) throws RemoteStorageException } } + List enrichedLogSegments(UnifiedLog log, Long fromOffset, Long lastStableOffset) { + List enrichedLogSegments = new ArrayList<>(); + List segments = JavaConverters.seqAsJavaList(log.nonActiveLogSegmentsFrom(fromOffset).toSeq()); + if (!segments.isEmpty()) { + int idx = 1; + for (; idx < segments.size(); idx++) { + LogSegment previous = segments.get(idx - 1); + LogSegment current = segments.get(idx); + enrichedLogSegments.add(new EnrichedLogSegment(previous, current.baseOffset())); + } + // LogSegment#readNextOffset() is an expensive call, so we only call it when necessary. + int lastIdx = idx - 1; + if (segments.get(lastIdx).baseOffset() < lastStableOffset) { + LogSegment last = segments.get(lastIdx); + enrichedLogSegments.add(new EnrichedLogSegment(last, last.readNextOffset())); + } + } + return enrichedLogSegments; + } + public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException { if (isCancelled()) return; @@ -536,7 +557,6 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException if (lso < 0) { logger.warn("lastStableOffset for partition {} is {}, which should not be negative.", topicIdPartition, lso); } else if (lso > 0 && copiedOffset < lso) { - long activeSegBaseOffset = log.activeSegment().baseOffset(); // log-start-offset can be ahead of the read-offset, when: // 1) log-start-offset gets incremented via delete-records API (or) // 2) enabling the remote log for the first time @@ -546,23 +566,23 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException // 1) Segment is not the active segment and // 2) Segment end-offset is less than the last-stable-offset as remote storage should contain only // committed/acked messages - List candidateSegments = JavaConverters.asJavaCollection(log.logSegments(fromOffset, Long.MAX_VALUE)) + List candidateSegments = enrichedLogSegments(log, fromOffset, lso) .stream() - .filter(segment -> segment.baseOffset() != activeSegBaseOffset && segment.readNextOffset() <= lso) + .filter(enriched -> enriched.segment.readNextOffset() <= lso) .collect(Collectors.toList()); logger.debug("Checking for segments to copy, copiedOffset: {} and lso: {}, candidateSegments: {}", copiedOffset, lso, candidateSegments); if (candidateSegments.isEmpty()) { logger.debug("No segments found to be copied for partition {} with copiedOffset: {} and active segment's base-offset: {}", - topicIdPartition, copiedOffset, activeSegBaseOffset); + topicIdPartition, copiedOffset, log.activeSegment().baseOffset()); } else { - for (LogSegment segment : candidateSegments) { + for (EnrichedLogSegment enrichedLogSegment : candidateSegments) { if (isCancelled() || !isLeader()) { logger.info("Skipping copying log segments as the current task state is changed, cancelled: {} leader:{}", isCancelled(), isLeader()); return; } - copyLogSegment(log, segment, segment.readNextOffset()); + copyLogSegment(log, enrichedLogSegment.segment, enrichedLogSegment.readNextOffset); } } } else { @@ -986,4 +1006,28 @@ public void close() { } } + static class EnrichedLogSegment { + private final LogSegment segment; + private final long readNextOffset; + + public EnrichedLogSegment(LogSegment segment, + long readNextOffset) { + this.segment = segment; + this.readNextOffset = readNextOffset; + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + EnrichedLogSegment that = (EnrichedLogSegment) o; + return readNextOffset == that.readNextOffset && Objects.equals(segment, that.segment); + } + + @Override + public int hashCode() { + return Objects.hash(segment, readNextOffset); + } + } + } diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index b9c8e07f80f70..6d66650ee56a0 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -322,16 +322,8 @@ private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.logSegments(anyLong(), anyLong())) - .thenAnswer(invocation -> { - Object[] args = invocation.getArguments(); - long fromOffset = (Long) args[0]; // inclusive - long toOffset = (Long) args[1]; // exclusive - return JavaConverters.collectionAsScalaIterable( - Arrays.asList(oldSegment, activeSegment).stream() - .filter(segment -> segment.baseOffset() >= fromOffset && segment.readNextOffset() < toOffset) - .collect(Collectors.toList())); - }); + when(mockLog.nonActiveLogSegmentsFrom(anyLong())) + .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment))); ProducerStateManager mockStateManager = mock(ProducerStateManager.class); when(mockLog.producerStateManager()).thenReturn(mockStateManager); @@ -514,7 +506,7 @@ void testMetricsUpdateOnCopyLogSegmentsFailure() throws Exception { when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.logSegments(anyLong(), anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment, activeSegment))); + when(mockLog.nonActiveLogSegmentsFrom(anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment))); ProducerStateManager mockStateManager = mock(ProducerStateManager.class); when(mockLog.producerStateManager()).thenReturn(mockStateManager); @@ -578,7 +570,7 @@ void testCopyLogSegmentsToRemoteShouldNotCopySegmentForFollower() throws Excepti when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.logSegments(anyLong(), anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment, activeSegment))); + when(mockLog.nonActiveLogSegmentsFrom(anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment))); when(mockLog.lastStableOffset()).thenReturn(250L); RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); @@ -857,6 +849,56 @@ public RemoteLogMetadataManager createRemoteLogMetadataManager() { } } + @Test + public void testEnrichedLogSegments() { + UnifiedLog log = mock(UnifiedLog.class); + LogSegment segment1 = mock(LogSegment.class); + LogSegment segment2 = mock(LogSegment.class); + LogSegment segment3 = mock(LogSegment.class); + + when(segment1.baseOffset()).thenReturn(5L); + when(segment2.baseOffset()).thenReturn(10L); + when(segment3.baseOffset()).thenReturn(15L); + when(segment3.readNextOffset()).thenReturn(22L); + + when(log.nonActiveLogSegmentsFrom(5L)) + .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, segment3))); + + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); + List expected = + Arrays.asList( + new RemoteLogManager.EnrichedLogSegment(segment1, 10L), + new RemoteLogManager.EnrichedLogSegment(segment2, 15L), + new RemoteLogManager.EnrichedLogSegment(segment3, 22L) + ); + + assertEquals(expected, task.enrichedLogSegments(log, 5L, 20L)); + } + + @Test + public void testEnrichedLogSegmentsSkipsNextReadOffsetCall() { + UnifiedLog log = mock(UnifiedLog.class); + LogSegment segment1 = mock(LogSegment.class); + LogSegment segment2 = mock(LogSegment.class); + LogSegment segment3 = mock(LogSegment.class); + + when(segment1.baseOffset()).thenReturn(5L); + when(segment2.baseOffset()).thenReturn(10L); + when(segment3.baseOffset()).thenReturn(15L); + + when(log.nonActiveLogSegmentsFrom(5L)) + .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, segment3))); + + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); + List expected = + Arrays.asList( + new RemoteLogManager.EnrichedLogSegment(segment1, 10L), + new RemoteLogManager.EnrichedLogSegment(segment2, 15L) + ); + + assertEquals(expected, task.enrichedLogSegments(log, 5L, 15L)); + } + private Partition mockPartition(TopicIdPartition topicIdPartition) { TopicPartition tp = topicIdPartition.topicPartition(); Partition partition = mock(Partition.class); From 3dab4ce01eb6649359af9c129754a79a623886f0 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 13:55:08 +0530 Subject: [PATCH 03/11] Addressed the review comments. --- core/src/main/java/kafka/log/remote/RemoteLogManager.java | 2 +- .../src/test/java/kafka/log/remote/RemoteLogManagerTest.java | 5 ++--- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index e8f3539ffaf5f..01c5414d09884 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -568,7 +568,7 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException // committed/acked messages List candidateSegments = enrichedLogSegments(log, fromOffset, lso) .stream() - .filter(enriched -> enriched.segment.readNextOffset() <= lso) + .filter(enrichedSegment -> enrichedSegment.readNextOffset <= lso) .collect(Collectors.toList()); logger.debug("Checking for segments to copy, copiedOffset: {} and lso: {}, candidateSegments: {}", copiedOffset, lso, candidateSegments); diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 6d66650ee56a0..c0017a00ffeb1 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -275,9 +275,8 @@ void testCopyLogSegmentsToRemoteShouldCopyExpectedLogSegment() throws Exception } /** - * The following values will be equal when the active segment gets rotated to passive when there are no new messages: - * last-stable-offset = high-water-mark = log-end-offset = base-offset-of-active-segment - * + * The following values will be equal when the active segment gets rotated to passive and there are no new messages: + * last-stable-offset = high-water-mark = log-end-offset = base-offset-of-active-segment. * This test asserts that the active log segment that was rotated after log.roll.ms are copied to remote storage. */ @Test From 882f7727fb55ca957f9c166df02c5b1364633dee Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 14:37:27 +0530 Subject: [PATCH 04/11] Addressed the review-2 comments. --- .../kafka/log/remote/RemoteLogManager.java | 45 ++++++++---------- .../log/remote/RemoteLogManagerTest.java | 47 ++++--------------- 2 files changed, 30 insertions(+), 62 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index 01c5414d09884..6ce255fbf9c25 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -524,22 +524,16 @@ private void maybeUpdateReadOffset(UnifiedLog log) throws RemoteStorageException } } - List enrichedLogSegments(UnifiedLog log, Long fromOffset, Long lastStableOffset) { + List enrichedLogSegments(UnifiedLog log, Long fromOffset) { List enrichedLogSegments = new ArrayList<>(); - List segments = JavaConverters.seqAsJavaList(log.nonActiveLogSegmentsFrom(fromOffset).toSeq()); + List segments = JavaConverters.seqAsJavaList(log.logSegments(fromOffset, Long.MAX_VALUE).toSeq()); if (!segments.isEmpty()) { - int idx = 1; - for (; idx < segments.size(); idx++) { - LogSegment previous = segments.get(idx - 1); - LogSegment current = segments.get(idx); - enrichedLogSegments.add(new EnrichedLogSegment(previous, current.baseOffset())); - } - // LogSegment#readNextOffset() is an expensive call, so we only call it when necessary. - int lastIdx = idx - 1; - if (segments.get(lastIdx).baseOffset() < lastStableOffset) { - LogSegment last = segments.get(lastIdx); - enrichedLogSegments.add(new EnrichedLogSegment(last, last.readNextOffset())); + for (int idx = 1; idx < segments.size(); idx++) { + LogSegment previousSeg = segments.get(idx - 1); + LogSegment currentSeg = segments.get(idx); + enrichedLogSegments.add(new EnrichedLogSegment(previousSeg, currentSeg.baseOffset())); } + // Discard the last active segment } return enrichedLogSegments; } @@ -561,20 +555,21 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException // 1) log-start-offset gets incremented via delete-records API (or) // 2) enabling the remote log for the first time long fromOffset = Math.max(copiedOffset + 1, log.logStartOffset()); + long activeSegmentBaseOffset = log.activeSegment().baseOffset(); // Segments which match the following criteria are eligible for copying to remote storage: // 1) Segment is not the active segment and // 2) Segment end-offset is less than the last-stable-offset as remote storage should contain only // committed/acked messages - List candidateSegments = enrichedLogSegments(log, fromOffset, lso) + List candidateSegments = enrichedLogSegments(log, fromOffset) .stream() - .filter(enrichedSegment -> enrichedSegment.readNextOffset <= lso) + .filter(enrichedSegment -> enrichedSegment.logSegment.baseOffset() != activeSegmentBaseOffset && enrichedSegment.nextSegmentOffset <= lso) .collect(Collectors.toList()); logger.debug("Checking for segments to copy, copiedOffset: {} and lso: {}, candidateSegments: {}", copiedOffset, lso, candidateSegments); if (candidateSegments.isEmpty()) { logger.debug("No segments found to be copied for partition {} with copiedOffset: {} and active segment's base-offset: {}", - topicIdPartition, copiedOffset, log.activeSegment().baseOffset()); + topicIdPartition, copiedOffset, activeSegmentBaseOffset); } else { for (EnrichedLogSegment enrichedLogSegment : candidateSegments) { if (isCancelled() || !isLeader()) { @@ -582,7 +577,7 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException isCancelled(), isLeader()); return; } - copyLogSegment(log, enrichedLogSegment.segment, enrichedLogSegment.readNextOffset); + copyLogSegment(log, enrichedLogSegment.logSegment, enrichedLogSegment.nextSegmentOffset); } } } else { @@ -1007,13 +1002,13 @@ public void close() { } static class EnrichedLogSegment { - private final LogSegment segment; - private final long readNextOffset; + private final LogSegment logSegment; + private final long nextSegmentOffset; - public EnrichedLogSegment(LogSegment segment, - long readNextOffset) { - this.segment = segment; - this.readNextOffset = readNextOffset; + public EnrichedLogSegment(LogSegment logSegment, + long nextSegmentOffset) { + this.logSegment = logSegment; + this.nextSegmentOffset = nextSegmentOffset; } @Override @@ -1021,12 +1016,12 @@ public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; EnrichedLogSegment that = (EnrichedLogSegment) o; - return readNextOffset == that.readNextOffset && Objects.equals(segment, that.segment); + return nextSegmentOffset == that.nextSegmentOffset && Objects.equals(logSegment, that.logSegment); } @Override public int hashCode() { - return Objects.hash(segment, readNextOffset); + return Objects.hash(logSegment, nextSegmentOffset); } } diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index c0017a00ffeb1..b0d543586c69b 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -321,8 +321,7 @@ private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.nonActiveLogSegmentsFrom(anyLong())) - .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment))); + when(mockLog.logSegments(anyLong(), anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment, activeSegment))); ProducerStateManager mockStateManager = mock(ProducerStateManager.class); when(mockLog.producerStateManager()).thenReturn(mockStateManager); @@ -374,7 +373,7 @@ private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, assertEquals(remoteLogSegmentMetadataArg.getValue(), remoteLogSegmentMetadataArg2.getValue()); // The old segment should only contain leader epoch [0->0, 1->100] since its offset range is [0, 149] verifyLogSegmentData(logSegmentDataArg.getValue(), idx, timeIdx, txnIndex, tempFile, mockProducerSnapshotIndex, - Arrays.asList(epochEntry0, epochEntry1)); + Arrays.asList(epochEntry0, epochEntry1)); // verify remoteLogMetadataManager did add the expected RemoteLogSegmentMetadataUpdate ArgumentCaptor remoteLogSegmentMetadataUpdateArg = ArgumentCaptor.forClass(RemoteLogSegmentMetadataUpdate.class); @@ -505,7 +504,7 @@ void testMetricsUpdateOnCopyLogSegmentsFailure() throws Exception { when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.nonActiveLogSegmentsFrom(anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment))); + when(mockLog.logSegments(anyLong(), anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment, activeSegment))); ProducerStateManager mockStateManager = mock(ProducerStateManager.class); when(mockLog.producerStateManager()).thenReturn(mockStateManager); @@ -569,7 +568,7 @@ void testCopyLogSegmentsToRemoteShouldNotCopySegmentForFollower() throws Excepti when(mockLog.activeSegment()).thenReturn(activeSegment); when(mockLog.logStartOffset()).thenReturn(oldSegmentStartOffset); - when(mockLog.nonActiveLogSegmentsFrom(anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment))); + when(mockLog.logSegments(anyLong(), anyLong())).thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(oldSegment, activeSegment))); when(mockLog.lastStableOffset()).thenReturn(250L); RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); @@ -853,40 +852,14 @@ public void testEnrichedLogSegments() { UnifiedLog log = mock(UnifiedLog.class); LogSegment segment1 = mock(LogSegment.class); LogSegment segment2 = mock(LogSegment.class); - LogSegment segment3 = mock(LogSegment.class); - - when(segment1.baseOffset()).thenReturn(5L); - when(segment2.baseOffset()).thenReturn(10L); - when(segment3.baseOffset()).thenReturn(15L); - when(segment3.readNextOffset()).thenReturn(22L); - - when(log.nonActiveLogSegmentsFrom(5L)) - .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, segment3))); - - RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); - List expected = - Arrays.asList( - new RemoteLogManager.EnrichedLogSegment(segment1, 10L), - new RemoteLogManager.EnrichedLogSegment(segment2, 15L), - new RemoteLogManager.EnrichedLogSegment(segment3, 22L) - ); - - assertEquals(expected, task.enrichedLogSegments(log, 5L, 20L)); - } - - @Test - public void testEnrichedLogSegmentsSkipsNextReadOffsetCall() { - UnifiedLog log = mock(UnifiedLog.class); - LogSegment segment1 = mock(LogSegment.class); - LogSegment segment2 = mock(LogSegment.class); - LogSegment segment3 = mock(LogSegment.class); + LogSegment activeSegment = mock(LogSegment.class); when(segment1.baseOffset()).thenReturn(5L); when(segment2.baseOffset()).thenReturn(10L); - when(segment3.baseOffset()).thenReturn(15L); + when(activeSegment.baseOffset()).thenReturn(15L); - when(log.nonActiveLogSegmentsFrom(5L)) - .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, segment3))); + when(log.logSegments(5L, Long.MAX_VALUE)) + .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, activeSegment))); RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); List expected = @@ -894,8 +867,8 @@ public void testEnrichedLogSegmentsSkipsNextReadOffsetCall() { new RemoteLogManager.EnrichedLogSegment(segment1, 10L), new RemoteLogManager.EnrichedLogSegment(segment2, 15L) ); - - assertEquals(expected, task.enrichedLogSegments(log, 5L, 15L)); + List actual = task.enrichedLogSegments(log, 5L); + assertEquals(expected, actual); } private Partition mockPartition(TopicIdPartition topicIdPartition) { From 154c52509737bd5e13e8a30866a373b77b2eecc6 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 15:37:17 +0530 Subject: [PATCH 05/11] Added toString --- core/src/main/java/kafka/log/remote/RemoteLogManager.java | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index 6ce255fbf9c25..564e1ea0464ea 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -1023,6 +1023,14 @@ public boolean equals(Object o) { public int hashCode() { return Objects.hash(logSegment, nextSegmentOffset); } + + @Override + public String toString() { + return "EnrichedLogSegment{" + + "logSegment=" + logSegment + + ", nextSegmentOffset=" + nextSegmentOffset + + '}'; + } } } From 272443cb0f129766dd8f9ab60335088ee06afbcb Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 20:57:32 +0530 Subject: [PATCH 06/11] Renamed EnrichedLogSegment to CandidateLogSegment and addressed the review comments. --- .../kafka/log/remote/RemoteLogManager.java | 43 +++++++++++-------- .../log/remote/RemoteLogManagerTest.java | 36 +++++++++++++--- 2 files changed, 55 insertions(+), 24 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index 564e1ea0464ea..dbb6005f91606 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -524,18 +524,30 @@ private void maybeUpdateReadOffset(UnifiedLog log) throws RemoteStorageException } } - List enrichedLogSegments(UnifiedLog log, Long fromOffset) { - List enrichedLogSegments = new ArrayList<>(); + /** + * Segments which match the following criteria are eligible for copying to remote storage: + * 1) Segment is not the active segment and + * 2) Segment end-offset is less than the last-stable-offset as remote storage should contain only + * committed/acked messages + * @param log The log from which the segments are to be copied + * @param fromOffset The offset from which the segments are to be copied + * @param lastStableOffset The last stable offset of the log + * @return candidate log segments to be copied to remote storage + */ + List candidateLogSegments(UnifiedLog log, Long fromOffset, Long lastStableOffset) { + List candidateLogSegments = new ArrayList<>(); List segments = JavaConverters.seqAsJavaList(log.logSegments(fromOffset, Long.MAX_VALUE).toSeq()); if (!segments.isEmpty()) { for (int idx = 1; idx < segments.size(); idx++) { LogSegment previousSeg = segments.get(idx - 1); LogSegment currentSeg = segments.get(idx); - enrichedLogSegments.add(new EnrichedLogSegment(previousSeg, currentSeg.baseOffset())); + if (currentSeg.baseOffset() <= lastStableOffset) { + candidateLogSegments.add(new CandidateLogSegment(previousSeg, currentSeg.baseOffset())); + } } // Discard the last active segment } - return enrichedLogSegments; + return candidateLogSegments; } public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException { @@ -557,27 +569,20 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException long fromOffset = Math.max(copiedOffset + 1, log.logStartOffset()); long activeSegmentBaseOffset = log.activeSegment().baseOffset(); - // Segments which match the following criteria are eligible for copying to remote storage: - // 1) Segment is not the active segment and - // 2) Segment end-offset is less than the last-stable-offset as remote storage should contain only - // committed/acked messages - List candidateSegments = enrichedLogSegments(log, fromOffset) - .stream() - .filter(enrichedSegment -> enrichedSegment.logSegment.baseOffset() != activeSegmentBaseOffset && enrichedSegment.nextSegmentOffset <= lso) - .collect(Collectors.toList()); + List candidateSegments = candidateLogSegments(log, fromOffset, lso); logger.debug("Checking for segments to copy, copiedOffset: {} and lso: {}, candidateSegments: {}", copiedOffset, lso, candidateSegments); if (candidateSegments.isEmpty()) { logger.debug("No segments found to be copied for partition {} with copiedOffset: {} and active segment's base-offset: {}", topicIdPartition, copiedOffset, activeSegmentBaseOffset); } else { - for (EnrichedLogSegment enrichedLogSegment : candidateSegments) { + for (CandidateLogSegment candidateLogSegment : candidateSegments) { if (isCancelled() || !isLeader()) { logger.info("Skipping copying log segments as the current task state is changed, cancelled: {} leader:{}", isCancelled(), isLeader()); return; } - copyLogSegment(log, enrichedLogSegment.logSegment, enrichedLogSegment.nextSegmentOffset); + copyLogSegment(log, candidateLogSegment.logSegment, candidateLogSegment.nextSegmentOffset); } } } else { @@ -1001,12 +1006,12 @@ public void close() { } } - static class EnrichedLogSegment { + static class CandidateLogSegment { private final LogSegment logSegment; private final long nextSegmentOffset; - public EnrichedLogSegment(LogSegment logSegment, - long nextSegmentOffset) { + public CandidateLogSegment(LogSegment logSegment, + long nextSegmentOffset) { this.logSegment = logSegment; this.nextSegmentOffset = nextSegmentOffset; } @@ -1015,7 +1020,7 @@ public EnrichedLogSegment(LogSegment logSegment, public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; - EnrichedLogSegment that = (EnrichedLogSegment) o; + CandidateLogSegment that = (CandidateLogSegment) o; return nextSegmentOffset == that.nextSegmentOffset && Objects.equals(logSegment, that.logSegment); } @@ -1026,7 +1031,7 @@ public int hashCode() { @Override public String toString() { - return "EnrichedLogSegment{" + + return "CandidateLogSegment{" + "logSegment=" + logSegment + ", nextSegmentOffset=" + nextSegmentOffset + '}'; diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index b0d543586c69b..2b72d28ba4fd1 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -848,7 +848,7 @@ public RemoteLogMetadataManager createRemoteLogMetadataManager() { } @Test - public void testEnrichedLogSegments() { + public void testCandidateLogSegmentsSkipsActiveSegment() { UnifiedLog log = mock(UnifiedLog.class); LogSegment segment1 = mock(LogSegment.class); LogSegment segment2 = mock(LogSegment.class); @@ -862,12 +862,38 @@ public void testEnrichedLogSegments() { .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, activeSegment))); RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); - List expected = + List expected = Arrays.asList( - new RemoteLogManager.EnrichedLogSegment(segment1, 10L), - new RemoteLogManager.EnrichedLogSegment(segment2, 15L) + new RemoteLogManager.CandidateLogSegment(segment1, 10L), + new RemoteLogManager.CandidateLogSegment(segment2, 15L) ); - List actual = task.enrichedLogSegments(log, 5L); + List actual = task.candidateLogSegments(log, 5L, 15L); + assertEquals(expected, actual); + } + + @Test + public void testCandidateLogSegmentsSkipsSegmentsBelowLastStableOffset() { + UnifiedLog log = mock(UnifiedLog.class); + LogSegment segment1 = mock(LogSegment.class); + LogSegment segment2 = mock(LogSegment.class); + LogSegment segment3 = mock(LogSegment.class); + LogSegment activeSegment = mock(LogSegment.class); + + when(segment1.baseOffset()).thenReturn(5L); + when(segment2.baseOffset()).thenReturn(10L); + when(segment3.baseOffset()).thenReturn(15L); + when(activeSegment.baseOffset()).thenReturn(20L); + + when(log.logSegments(5L, Long.MAX_VALUE)) + .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, segment3, activeSegment))); + + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); + List expected = + Arrays.asList( + new RemoteLogManager.CandidateLogSegment(segment1, 10L), + new RemoteLogManager.CandidateLogSegment(segment2, 15L) + ); + List actual = task.candidateLogSegments(log, 5L, 15L); assertEquals(expected, actual); } From c800c6cf4e2fc395d066baa7320718fd1e9e59b8 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 21:04:12 +0530 Subject: [PATCH 07/11] Update the logger. --- .../main/java/kafka/log/remote/RemoteLogManager.java | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index dbb6005f91606..99f206055d261 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -563,18 +563,16 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException if (lso < 0) { logger.warn("lastStableOffset for partition {} is {}, which should not be negative.", topicIdPartition, lso); } else if (lso > 0 && copiedOffset < lso) { - // log-start-offset can be ahead of the read-offset, when: + // log-start-offset can be ahead of the copied-offset, when: // 1) log-start-offset gets incremented via delete-records API (or) // 2) enabling the remote log for the first time long fromOffset = Math.max(copiedOffset + 1, log.logStartOffset()); - long activeSegmentBaseOffset = log.activeSegment().baseOffset(); - List candidateSegments = candidateLogSegments(log, fromOffset, lso); - logger.debug("Checking for segments to copy, copiedOffset: {} and lso: {}, candidateSegments: {}", - copiedOffset, lso, candidateSegments); + logger.debug("Candidate log segments, logStartOffset: {}, copiedOffset: {}, fromOffset: {}, lso: {} " + + "and candidateSegments: {}", log.logStartOffset(), copiedOffset, fromOffset, lso, candidateSegments); if (candidateSegments.isEmpty()) { logger.debug("No segments found to be copied for partition {} with copiedOffset: {} and active segment's base-offset: {}", - topicIdPartition, copiedOffset, activeSegmentBaseOffset); + topicIdPartition, copiedOffset, log.activeSegment().baseOffset()); } else { for (CandidateLogSegment candidateLogSegment : candidateSegments) { if (isCancelled() || !isLeader()) { From 63dd63f3a123d5e49a5ecd8ac1bbf3215162b36b Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Mon, 31 Jul 2023 21:11:35 +0530 Subject: [PATCH 08/11] Update the test parameter. --- core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 2b72d28ba4fd1..3a18d58f40576 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -867,7 +867,7 @@ public void testCandidateLogSegmentsSkipsActiveSegment() { new RemoteLogManager.CandidateLogSegment(segment1, 10L), new RemoteLogManager.CandidateLogSegment(segment2, 15L) ); - List actual = task.candidateLogSegments(log, 5L, 15L); + List actual = task.candidateLogSegments(log, 5L, 20L); assertEquals(expected, actual); } From fe6ccd334bb0775a5e31a6295698bf36d88423f4 Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Tue, 1 Aug 2023 12:03:15 +0530 Subject: [PATCH 09/11] Rename CandidateLogSegment to EnrichedLogSegment --- .../kafka/log/remote/RemoteLogManager.java | 25 ++++++++++--------- .../log/remote/RemoteLogManagerTest.java | 16 ++++++------ 2 files changed, 21 insertions(+), 20 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index 99f206055d261..0b9da1012c71f 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -534,15 +534,15 @@ private void maybeUpdateReadOffset(UnifiedLog log) throws RemoteStorageException * @param lastStableOffset The last stable offset of the log * @return candidate log segments to be copied to remote storage */ - List candidateLogSegments(UnifiedLog log, Long fromOffset, Long lastStableOffset) { - List candidateLogSegments = new ArrayList<>(); + List candidateLogSegments(UnifiedLog log, Long fromOffset, Long lastStableOffset) { + List candidateLogSegments = new ArrayList<>(); List segments = JavaConverters.seqAsJavaList(log.logSegments(fromOffset, Long.MAX_VALUE).toSeq()); if (!segments.isEmpty()) { for (int idx = 1; idx < segments.size(); idx++) { LogSegment previousSeg = segments.get(idx - 1); LogSegment currentSeg = segments.get(idx); if (currentSeg.baseOffset() <= lastStableOffset) { - candidateLogSegments.add(new CandidateLogSegment(previousSeg, currentSeg.baseOffset())); + candidateLogSegments.add(new EnrichedLogSegment(previousSeg, currentSeg.baseOffset())); } } // Discard the last active segment @@ -567,14 +567,14 @@ public void copyLogSegmentsToRemote(UnifiedLog log) throws InterruptedException // 1) log-start-offset gets incremented via delete-records API (or) // 2) enabling the remote log for the first time long fromOffset = Math.max(copiedOffset + 1, log.logStartOffset()); - List candidateSegments = candidateLogSegments(log, fromOffset, lso); + List candidateLogSegments = candidateLogSegments(log, fromOffset, lso); logger.debug("Candidate log segments, logStartOffset: {}, copiedOffset: {}, fromOffset: {}, lso: {} " + - "and candidateSegments: {}", log.logStartOffset(), copiedOffset, fromOffset, lso, candidateSegments); - if (candidateSegments.isEmpty()) { + "and candidateLogSegments: {}", log.logStartOffset(), copiedOffset, fromOffset, lso, candidateLogSegments); + if (candidateLogSegments.isEmpty()) { logger.debug("No segments found to be copied for partition {} with copiedOffset: {} and active segment's base-offset: {}", topicIdPartition, copiedOffset, log.activeSegment().baseOffset()); } else { - for (CandidateLogSegment candidateLogSegment : candidateSegments) { + for (EnrichedLogSegment candidateLogSegment : candidateLogSegments) { if (isCancelled() || !isLeader()) { logger.info("Skipping copying log segments as the current task state is changed, cancelled: {} leader:{}", isCancelled(), isLeader()); @@ -1004,12 +1004,13 @@ public void close() { } } - static class CandidateLogSegment { + // Visible for testing + static class EnrichedLogSegment { private final LogSegment logSegment; private final long nextSegmentOffset; - public CandidateLogSegment(LogSegment logSegment, - long nextSegmentOffset) { + public EnrichedLogSegment(LogSegment logSegment, + long nextSegmentOffset) { this.logSegment = logSegment; this.nextSegmentOffset = nextSegmentOffset; } @@ -1018,7 +1019,7 @@ public CandidateLogSegment(LogSegment logSegment, public boolean equals(Object o) { if (this == o) return true; if (o == null || getClass() != o.getClass()) return false; - CandidateLogSegment that = (CandidateLogSegment) o; + EnrichedLogSegment that = (EnrichedLogSegment) o; return nextSegmentOffset == that.nextSegmentOffset && Objects.equals(logSegment, that.logSegment); } @@ -1029,7 +1030,7 @@ public int hashCode() { @Override public String toString() { - return "CandidateLogSegment{" + + return "EnrichedLogSegment{" + "logSegment=" + logSegment + ", nextSegmentOffset=" + nextSegmentOffset + '}'; diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 3a18d58f40576..b9ea95130c7d0 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -862,12 +862,12 @@ public void testCandidateLogSegmentsSkipsActiveSegment() { .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, activeSegment))); RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); - List expected = + List expected = Arrays.asList( - new RemoteLogManager.CandidateLogSegment(segment1, 10L), - new RemoteLogManager.CandidateLogSegment(segment2, 15L) + new RemoteLogManager.EnrichedLogSegment(segment1, 10L), + new RemoteLogManager.EnrichedLogSegment(segment2, 15L) ); - List actual = task.candidateLogSegments(log, 5L, 20L); + List actual = task.candidateLogSegments(log, 5L, 20L); assertEquals(expected, actual); } @@ -888,12 +888,12 @@ public void testCandidateLogSegmentsSkipsSegmentsBelowLastStableOffset() { .thenReturn(JavaConverters.collectionAsScalaIterable(Arrays.asList(segment1, segment2, segment3, activeSegment))); RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition); - List expected = + List expected = Arrays.asList( - new RemoteLogManager.CandidateLogSegment(segment1, 10L), - new RemoteLogManager.CandidateLogSegment(segment2, 15L) + new RemoteLogManager.EnrichedLogSegment(segment1, 10L), + new RemoteLogManager.EnrichedLogSegment(segment2, 15L) ); - List actual = task.candidateLogSegments(log, 5L, 15L); + List actual = task.candidateLogSegments(log, 5L, 15L); assertEquals(expected, actual); } From 8fc4e2e2b2023c80c6862632b077b2416e85b2bc Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Tue, 1 Aug 2023 12:07:32 +0530 Subject: [PATCH 10/11] Update the test --- core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index b9ea95130c7d0..14d22d30441b2 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -309,9 +309,7 @@ private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, LogSegment activeSegment = mock(LogSegment.class); when(oldSegment.baseOffset()).thenReturn(oldSegmentStartOffset); - when(oldSegment.readNextOffset()).thenReturn(nextSegmentStartOffset); when(activeSegment.baseOffset()).thenReturn(nextSegmentStartOffset); - when(activeSegment.readNextOffset()).thenReturn(logEndOffset); FileRecords fileRecords = mock(FileRecords.class); when(oldSegment.log()).thenReturn(fileRecords); @@ -872,7 +870,7 @@ public void testCandidateLogSegmentsSkipsActiveSegment() { } @Test - public void testCandidateLogSegmentsSkipsSegmentsBelowLastStableOffset() { + public void testCandidateLogSegmentsSkipsSegmentsAfterLastStableOffset() { UnifiedLog log = mock(UnifiedLog.class); LogSegment segment1 = mock(LogSegment.class); LogSegment segment2 = mock(LogSegment.class); From c8f06a7c3afebbcdb392dec9f18d86ee8b2056bb Mon Sep 17 00:00:00 2001 From: Kamal Chandraprakash Date: Tue, 1 Aug 2023 12:18:14 +0530 Subject: [PATCH 11/11] Update the test --- core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 14d22d30441b2..36a41f9671475 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -310,6 +310,8 @@ private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, when(oldSegment.baseOffset()).thenReturn(oldSegmentStartOffset); when(activeSegment.baseOffset()).thenReturn(nextSegmentStartOffset); + verify(oldSegment, times(0)).readNextOffset(); + verify(activeSegment, times(0)).readNextOffset(); FileRecords fileRecords = mock(FileRecords.class); when(oldSegment.log()).thenReturn(fileRecords); @@ -325,6 +327,7 @@ private void assertCopyExpectedLogSegmentsToRemote(long oldSegmentStartOffset, when(mockLog.producerStateManager()).thenReturn(mockStateManager); when(mockStateManager.fetchSnapshot(anyLong())).thenReturn(Optional.of(mockProducerSnapshotIndex)); when(mockLog.lastStableOffset()).thenReturn(lastStableOffset); + when(mockLog.logEndOffset()).thenReturn(logEndOffset); LazyIndex idx = LazyIndex.forOffset(UnifiedLog.offsetIndexFile(tempDir, oldSegmentStartOffset, ""), oldSegmentStartOffset, 1000); LazyIndex timeIdx = LazyIndex.forTime(UnifiedLog.timeIndexFile(tempDir, oldSegmentStartOffset, ""), oldSegmentStartOffset, 1500);