From 80c0521a4577233735b4abd8a520dd9bcaf48eb8 Mon Sep 17 00:00:00 2001 From: Christo Date: Fri, 8 Sep 2023 13:48:16 +0100 Subject: [PATCH 1/3] KAFKA-15352: Update log-start-offset before initiating deletion of remote segments --- .../kafka/log/remote/RemoteLogManager.java | 97 +++++++----- .../log/remote/RemoteLogManagerTest.java | 144 +++++++++++++++++- 2 files changed, 200 insertions(+), 41 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index 4a35abf6a115f..171bf0f659c27 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -832,70 +832,64 @@ public RemoteLogRetentionHandler(Optional retentionSizeData, remainingBreachedSize = retentionSizeData.map(sizeData -> sizeData.remainingBreachedSize).orElse(0L); } - private boolean deleteRetentionSizeBreachedSegments(RemoteLogSegmentMetadata metadata) throws RemoteStorageException, ExecutionException, InterruptedException { + private boolean deleteRetentionSizeBreachedSegments(RemoteLogSegmentMetadata metadata) { if (!retentionSizeData.isPresent()) { return false; } - boolean isSegmentDeleted = deleteRemoteLogSegment(metadata, ignored -> { - // Assumption that segments contain size >= 0 - if (remainingBreachedSize > 0) { - long remainingBytes = remainingBreachedSize - metadata.segmentSizeInBytes(); - if (remainingBytes >= 0) { - remainingBreachedSize = remainingBytes; - return true; - } + boolean shouldDeleteSegment = false; + + // Assumption that segments contain size >= 0 + if (remainingBreachedSize > 0) { + long remainingBytes = remainingBreachedSize - metadata.segmentSizeInBytes(); + if (remainingBytes >= 0) { + remainingBreachedSize = remainingBytes; + shouldDeleteSegment = true; } + } - return false; - }); - if (isSegmentDeleted) { + if (shouldDeleteSegment) { logStartOffset = OptionalLong.of(metadata.endOffset() + 1); - logger.info("Deleted remote log segment {} due to retention size {} breach. Log size after deletion will be {}.", + logger.info("About to delete remote log segment {} due to retention size {} breach. Log size after deletion will be {}.", metadata.remoteLogSegmentId(), retentionSizeData.get().retentionSize, remainingBreachedSize + retentionSizeData.get().retentionSize); } - return isSegmentDeleted; + return shouldDeleteSegment; } - public boolean deleteRetentionTimeBreachedSegments(RemoteLogSegmentMetadata metadata) - throws RemoteStorageException, ExecutionException, InterruptedException { + public boolean deleteRetentionTimeBreachedSegments(RemoteLogSegmentMetadata metadata) { if (!retentionTimeData.isPresent()) { return false; } - boolean isSegmentDeleted = deleteRemoteLogSegment(metadata, - ignored -> metadata.maxTimestampMs() <= retentionTimeData.get().cleanupUntilMs); - if (isSegmentDeleted) { + boolean shouldDeleteSegment = metadata.maxTimestampMs() <= retentionTimeData.get().cleanupUntilMs; + if (shouldDeleteSegment) { remainingBreachedSize = Math.max(0, remainingBreachedSize - metadata.segmentSizeInBytes()); // It is fine to have logStartOffset as `metadata.endOffset() + 1` as the segment offset intervals // are ascending with in an epoch. logStartOffset = OptionalLong.of(metadata.endOffset() + 1); - logger.info("Deleted remote log segment {} due to retention time {}ms breach based on the largest record timestamp in the segment", + logger.info("About to delete remote log segment {} due to retention time {}ms breach based on the largest record timestamp in the segment", metadata.remoteLogSegmentId(), retentionTimeData.get().retentionMs); } - return isSegmentDeleted; + return shouldDeleteSegment; } private boolean deleteLogStartOffsetBreachedSegments(RemoteLogSegmentMetadata metadata, long logStartOffset, - NavigableMap leaderEpochEntries) - throws RemoteStorageException, ExecutionException, InterruptedException { - boolean isSegmentDeleted = deleteRemoteLogSegment(metadata, ignored -> { - if (!leaderEpochEntries.isEmpty()) { - // Note that `logStartOffset` and `leaderEpochEntries.firstEntry().getValue()` should be same - Integer firstEpoch = leaderEpochEntries.firstKey(); - return metadata.segmentLeaderEpochs().keySet().stream().allMatch(epoch -> epoch <= firstEpoch) - && metadata.endOffset() < logStartOffset; - } - return false; - }); - if (isSegmentDeleted) { - logger.info("Deleted remote log segment {} due to log-start-offset {} breach. " + + NavigableMap leaderEpochEntries) { + boolean shouldDeleteSegment = false; + if (!leaderEpochEntries.isEmpty()) { + // Note that `logStartOffset` and `leaderEpochEntries.firstEntry().getValue()` should be same + Integer firstEpoch = leaderEpochEntries.firstKey(); + shouldDeleteSegment = metadata.segmentLeaderEpochs().keySet().stream().allMatch(epoch -> epoch <= firstEpoch) + && metadata.endOffset() < logStartOffset; + } + if (shouldDeleteSegment) { + logger.info("About to delete remote log segment {} due to log-start-offset {} breach. " + "Current earliest-epoch-entry: {}, segment-end-offset: {} and segment-epochs: {}", metadata.remoteLogSegmentId(), logStartOffset, leaderEpochEntries.firstEntry(), metadata.endOffset(), metadata.segmentLeaderEpochs()); } - return isSegmentDeleted; + return shouldDeleteSegment; } // It removes the segments beyond the current leader's earliest epoch. Those segments are considered as @@ -989,6 +983,7 @@ private void cleanupExpiredRemoteLogSegments() throws RemoteStorageException, Ex RemoteLogRetentionHandler remoteLogRetentionHandler = new RemoteLogRetentionHandler(retentionSizeData, retentionTimeData); Iterator epochIterator = epochWithOffsets.navigableKeySet().iterator(); boolean canProcess = true; + List segmentsToDelete = new ArrayList<>(); while (canProcess && epochIterator.hasNext()) { Integer epoch = epochIterator.next(); Iterator segmentsIterator = remoteLogMetadataManager.listRemoteLogSegments(topicIdPartition, epoch); @@ -1004,19 +999,25 @@ private void cleanupExpiredRemoteLogSegments() throws RemoteStorageException, Ex // remote log segments won't be removed. The `isRemoteSegmentWithinLeaderEpoch` validates whether // the epochs present in the segment lies in the checkpoint file. It will always return false // since the checkpoint file was already truncated. - boolean isSegmentDeleted = remoteLogRetentionHandler.deleteLogStartOffsetBreachedSegments( + boolean shouldDeleteSegment = remoteLogRetentionHandler.deleteLogStartOffsetBreachedSegments( metadata, logStartOffset, epochWithOffsets); + if (shouldDeleteSegment) { + segmentsToDelete.add(metadata); + } boolean isValidSegment = false; - if (!isSegmentDeleted) { + if (!shouldDeleteSegment) { // check whether the segment contains the required epoch range with in the current leader epoch lineage. isValidSegment = isRemoteSegmentWithinLeaderEpochs(metadata, logEndOffset, epochWithOffsets); if (isValidSegment) { - isSegmentDeleted = + shouldDeleteSegment = remoteLogRetentionHandler.deleteRetentionTimeBreachedSegments(metadata) || remoteLogRetentionHandler.deleteRetentionSizeBreachedSegments(metadata); + if (shouldDeleteSegment) { + segmentsToDelete.add(metadata); + } } } - canProcess = isSegmentDeleted || !isValidSegment; + canProcess = shouldDeleteSegment || !isValidSegment; } } @@ -1045,6 +1046,24 @@ private void cleanupExpiredRemoteLogSegments() throws RemoteStorageException, Ex // Update log start offset with the computed value after retention cleanup is done remoteLogRetentionHandler.logStartOffset.ifPresent(offset -> handleLogStartOffsetUpdate(topicIdPartition.topicPartition(), offset)); + + // At this point in time we have updated the log start offsets, but not initiated a deletion. + // Either a follower has picked up the changes to the log start offset, or they have not. + // If the follower HAS picked up the changes, and they become the leader this replica won't successfully complete + // the deletion. + // However, the new leader will correctly pick up all breaching segments as log start offset breaching ones + // and delete them accordingly. + // If the follower HAS NOT picked up the changes, and they become the leader then they will go through this process + // again and delete them with the original deletion reason i.e. size, time or log start offset breach. + List undeletedSegments = new ArrayList<>(); + for (RemoteLogSegmentMetadata segmentMetadata : segmentsToDelete) { + if (!remoteLogRetentionHandler.deleteRemoteLogSegment(segmentMetadata, x -> !isCancelled() && isLeader())) { + undeletedSegments.add(segmentMetadata.remoteLogSegmentId().toString()); + } + } + if (!undeletedSegments.isEmpty()) { + logger.info("The following remote segments could not be deleted: {}", String.join(",", undeletedSegments)); + } } private Optional buildRetentionTimeData(long retentionMs) { diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 92a6c63537cd4..4df025eb8e03c 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -59,6 +59,7 @@ import org.apache.kafka.storage.internals.epoch.LeaderEpochFileCache; import org.apache.kafka.storage.internals.log.EpochEntry; import org.apache.kafka.storage.internals.log.LazyIndex; +import org.apache.kafka.storage.internals.log.LogConfig; import org.apache.kafka.storage.internals.log.OffsetIndex; import org.apache.kafka.storage.internals.log.ProducerStateManager; import org.apache.kafka.storage.internals.log.TimeIndex; @@ -179,6 +180,8 @@ public List read() { private final UnifiedLog mockLog = mock(UnifiedLog.class); + private final List> events = new ArrayList<>(); + @BeforeEach void setUp() throws Exception { topicIds.put(leaderTopicIdPartition.topicPartition().topic(), leaderTopicIdPartition.topicId()); @@ -191,7 +194,7 @@ void setUp() throws Exception { kafka.utils.TestUtils.clearYammerMetrics(); remoteLogManager = new RemoteLogManager(remoteLogManagerConfig, brokerId, logDir, clusterId, time, tp -> Optional.of(mockLog), - (topicPartition, offset) -> { }, + (topicPartition, offset) -> events.add(Collections.singletonMap(topicPartition, offset)), brokerTopicStats) { public RemoteStorageManager createRemoteStorageManager() { return remoteStorageManager; @@ -1508,16 +1511,153 @@ public RemoteLogMetadataManager createRemoteLogMetadataManager() { } } + @Test + public void testDeleteRetentionSizeBreachingSegments() throws RemoteStorageException { + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + task.convertToLeader(0); + + when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); + when(mockLog.logEndOffset()).thenReturn(200L); + + List epochEntries = Collections.singletonList(epochEntry0); + + List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); + + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + + checkpoint.write(epochEntries); + LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); + when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); + + Map logProps = new HashMap<>(); + logProps.put("retention.bytes", 0L); + logProps.put("retention.ms", -1L); + LogConfig mockLogConfig = new LogConfig(logProps); + when(mockLog.config()).thenReturn(mockLogConfig); + + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> { + // assert that log-start-offset has been moved accordingly + // we skip the first entry as it is the local replica ensuring it has the correct log start offset + assertEquals(200, events.get(1).get(leaderTopicIdPartition.topicPartition())); + return CompletableFuture.runAsync(() -> { }); + }); + + task.run(); + + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + } + + @Test + public void testDeleteRetentionMsBreachingSegments() throws RemoteStorageException { + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + task.convertToLeader(0); + + when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); + when(mockLog.logEndOffset()).thenReturn(200L); + + List epochEntries = Collections.singletonList(epochEntry0); + + List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); + + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + + checkpoint.write(epochEntries); + LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); + when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); + + Map logProps = new HashMap<>(); + logProps.put("retention.bytes", -1L); + logProps.put("retention.ms", 0L); + LogConfig mockLogConfig = new LogConfig(logProps); + when(mockLog.config()).thenReturn(mockLogConfig); + + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> { + // assert that log-start-offset has been moved accordingly + // we skip the first entry as it is the local replica ensuring it has the correct log start offset + assertEquals(200, events.get(1).get(leaderTopicIdPartition.topicPartition())); + return CompletableFuture.runAsync(() -> { }); + }); + + task.run(); + + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + } + + @Test + public void testDeleteRetentionMsBeingCancelledBeforeSecondDelete() throws RemoteStorageException { + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + task.convertToLeader(0); + + when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); + when(mockLog.logEndOffset()).thenReturn(200L); + + List epochEntries = Collections.singletonList(epochEntry0); + + List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); + + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + + checkpoint.write(epochEntries); + LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); + when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); + + Map logProps = new HashMap<>(); + logProps.put("retention.bytes", -1L); + logProps.put("retention.ms", 0L); + LogConfig mockLogConfig = new LogConfig(logProps); + when(mockLog.config()).thenReturn(mockLogConfig); + + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> { + // assert that log-start-offset has been moved accordingly + // we skip the first entry as it is the local replica ensuring it has the correct log start offset + assertEquals(200, events.get(1).get(leaderTopicIdPartition.topicPartition())); + // cancel the task so that we don't delete the second segment + task.cancel(); + return CompletableFuture.runAsync(() -> { }); + }); + + task.run(); + + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager, never()).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + } + private List listRemoteLogSegmentMetadata(TopicIdPartition topicIdPartition, int segmentCount, int recordsPerSegment, int segmentSize) { + return listRemoteLogSegmentMetadata(topicIdPartition, segmentCount, recordsPerSegment, segmentSize, Collections.emptyList()); + } + + private List listRemoteLogSegmentMetadata(TopicIdPartition topicIdPartition, + int segmentCount, + int recordsPerSegment, + int segmentSize, + List epochEntries) { List segmentMetadataList = new ArrayList<>(); for (int idx = 0; idx < segmentCount; idx++) { long timestamp = time.milliseconds(); long startOffset = (long) idx * recordsPerSegment; long endOffset = startOffset + recordsPerSegment - 1; - Map segmentLeaderEpochs = truncateAndGetLeaderEpochs(totalEpochEntries, startOffset, endOffset); + List localTotalEpochEntries = epochEntries.isEmpty() ? totalEpochEntries : epochEntries; + Map segmentLeaderEpochs = truncateAndGetLeaderEpochs(localTotalEpochEntries, startOffset, endOffset); segmentMetadataList.add(new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), startOffset, endOffset, timestamp, brokerId, timestamp, segmentSize, segmentLeaderEpochs)); From 63da3241e1fdf0f09aca6106426c328cf21e24c3 Mon Sep 17 00:00:00 2001 From: Christo Date: Fri, 8 Sep 2023 17:59:01 +0100 Subject: [PATCH 2/3] Address comments from second round of review --- .../log/remote/RemoteLogManagerTest.java | 225 ++++++++++-------- 1 file changed, 127 insertions(+), 98 deletions(-) diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 4df025eb8e03c..898e610cc1bfc 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -180,8 +180,6 @@ public List read() { private final UnifiedLog mockLog = mock(UnifiedLog.class); - private final List> events = new ArrayList<>(); - @BeforeEach void setUp() throws Exception { topicIds.put(leaderTopicIdPartition.topicPartition().topic(), leaderTopicIdPartition.topicId()); @@ -194,7 +192,7 @@ void setUp() throws Exception { kafka.utils.TestUtils.clearYammerMetrics(); remoteLogManager = new RemoteLogManager(remoteLogManagerConfig, brokerId, logDir, clusterId, time, tp -> Optional.of(mockLog), - (topicPartition, offset) -> events.add(Collections.singletonMap(topicPartition, offset)), + (topicPartition, offset) -> { }, brokerTopicStats) { public RemoteStorageManager createRemoteStorageManager() { return remoteStorageManager; @@ -1512,131 +1510,162 @@ public RemoteLogMetadataManager createRemoteLogMetadataManager() { } @Test - public void testDeleteRetentionSizeBreachingSegments() throws RemoteStorageException { - RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); - task.convertToLeader(0); + public void testDeleteRetentionSizeBreachingSegments() throws RemoteStorageException, IOException { + AtomicLong logStartOffset = new AtomicLong(0); + try (RemoteLogManager remoteLogManager = new RemoteLogManager(remoteLogManagerConfig, brokerId, logDir, clusterId, time, + tp -> Optional.of(mockLog), + (topicPartition, offset) -> logStartOffset.set(offset), + brokerTopicStats) { + public RemoteStorageManager createRemoteStorageManager() { + return remoteStorageManager; + } + public RemoteLogMetadataManager createRemoteLogMetadataManager() { + return remoteLogMetadataManager; + } + }) { + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + task.convertToLeader(0); - when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); - when(mockLog.logEndOffset()).thenReturn(200L); + when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); + when(mockLog.logEndOffset()).thenReturn(200L); - List epochEntries = Collections.singletonList(epochEntry0); + List epochEntries = Collections.singletonList(epochEntry0); - List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); + List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); - when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) - .thenReturn(remoteLogSegmentMetadatas.iterator()); - when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) - .thenReturn(remoteLogSegmentMetadatas.iterator()) - .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()); - checkpoint.write(epochEntries); - LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); - when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); + checkpoint.write(epochEntries); + LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); + when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); - Map logProps = new HashMap<>(); - logProps.put("retention.bytes", 0L); - logProps.put("retention.ms", -1L); - LogConfig mockLogConfig = new LogConfig(logProps); - when(mockLog.config()).thenReturn(mockLogConfig); - - when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) - .thenAnswer(answer -> { - // assert that log-start-offset has been moved accordingly - // we skip the first entry as it is the local replica ensuring it has the correct log start offset - assertEquals(200, events.get(1).get(leaderTopicIdPartition.topicPartition())); - return CompletableFuture.runAsync(() -> { }); - }); + Map logProps = new HashMap<>(); + logProps.put("retention.bytes", 0L); + logProps.put("retention.ms", -1L); + LogConfig mockLogConfig = new LogConfig(logProps); + when(mockLog.config()).thenReturn(mockLogConfig); - task.run(); + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> CompletableFuture.runAsync(() -> { })); - verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); - verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + task.run(); + + assertEquals(200L, logStartOffset.get()); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + } } @Test - public void testDeleteRetentionMsBreachingSegments() throws RemoteStorageException { - RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); - task.convertToLeader(0); + public void testDeleteRetentionMsBreachingSegments() throws RemoteStorageException, IOException { + AtomicLong logStartOffset = new AtomicLong(0); + try (RemoteLogManager remoteLogManager = new RemoteLogManager(remoteLogManagerConfig, brokerId, logDir, clusterId, time, + tp -> Optional.of(mockLog), + (topicPartition, offset) -> logStartOffset.set(offset), + brokerTopicStats) { + public RemoteStorageManager createRemoteStorageManager() { + return remoteStorageManager; + } + public RemoteLogMetadataManager createRemoteLogMetadataManager() { + return remoteLogMetadataManager; + } + }) { + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + task.convertToLeader(0); - when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); - when(mockLog.logEndOffset()).thenReturn(200L); + when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); + when(mockLog.logEndOffset()).thenReturn(200L); - List epochEntries = Collections.singletonList(epochEntry0); + List epochEntries = Collections.singletonList(epochEntry0); - List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); + List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); - when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) - .thenReturn(remoteLogSegmentMetadatas.iterator()); - when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) - .thenReturn(remoteLogSegmentMetadatas.iterator()) - .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()); - checkpoint.write(epochEntries); - LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); - when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); + checkpoint.write(epochEntries); + LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); + when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); - Map logProps = new HashMap<>(); - logProps.put("retention.bytes", -1L); - logProps.put("retention.ms", 0L); - LogConfig mockLogConfig = new LogConfig(logProps); - when(mockLog.config()).thenReturn(mockLogConfig); - - when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) - .thenAnswer(answer -> { - // assert that log-start-offset has been moved accordingly - // we skip the first entry as it is the local replica ensuring it has the correct log start offset - assertEquals(200, events.get(1).get(leaderTopicIdPartition.topicPartition())); - return CompletableFuture.runAsync(() -> { }); - }); + Map logProps = new HashMap<>(); + logProps.put("retention.bytes", -1L); + logProps.put("retention.ms", 0L); + LogConfig mockLogConfig = new LogConfig(logProps); + when(mockLog.config()).thenReturn(mockLogConfig); - task.run(); + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> CompletableFuture.runAsync(() -> { })); + + task.run(); - verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); - verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + assertEquals(200L, logStartOffset.get()); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + } } @Test - public void testDeleteRetentionMsBeingCancelledBeforeSecondDelete() throws RemoteStorageException { - RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); - task.convertToLeader(0); + public void testDeleteRetentionMsBeingCancelledBeforeSecondDelete() throws RemoteStorageException, IOException { + AtomicLong logStartOffset = new AtomicLong(0); + try (RemoteLogManager remoteLogManager = new RemoteLogManager(remoteLogManagerConfig, brokerId, logDir, clusterId, time, + tp -> Optional.of(mockLog), + (topicPartition, offset) -> logStartOffset.set(offset), + brokerTopicStats) { + public RemoteStorageManager createRemoteStorageManager() { + return remoteStorageManager; + } + public RemoteLogMetadataManager createRemoteLogMetadataManager() { + return remoteLogMetadataManager; + } + }) { + RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + task.convertToLeader(0); - when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); - when(mockLog.logEndOffset()).thenReturn(200L); + when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); + when(mockLog.logEndOffset()).thenReturn(200L); - List epochEntries = Collections.singletonList(epochEntry0); + List epochEntries = Collections.singletonList(epochEntry0); - List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); + List remoteLogSegmentMetadatas = listRemoteLogSegmentMetadata(leaderTopicIdPartition, 2, 100, 1024, epochEntries); - when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) - .thenReturn(remoteLogSegmentMetadatas.iterator()); - when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) - .thenReturn(remoteLogSegmentMetadatas.iterator()) - .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition)) + .thenReturn(remoteLogSegmentMetadatas.iterator()); + when(remoteLogMetadataManager.listRemoteLogSegments(leaderTopicIdPartition, 0)) + .thenReturn(remoteLogSegmentMetadatas.iterator()) + .thenReturn(remoteLogSegmentMetadatas.iterator()); - checkpoint.write(epochEntries); - LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); - when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); + checkpoint.write(epochEntries); + LeaderEpochFileCache cache = new LeaderEpochFileCache(tp, checkpoint); + when(mockLog.leaderEpochCache()).thenReturn(Option.apply(cache)); - Map logProps = new HashMap<>(); - logProps.put("retention.bytes", -1L); - logProps.put("retention.ms", 0L); - LogConfig mockLogConfig = new LogConfig(logProps); - when(mockLog.config()).thenReturn(mockLogConfig); - - when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) - .thenAnswer(answer -> { - // assert that log-start-offset has been moved accordingly - // we skip the first entry as it is the local replica ensuring it has the correct log start offset - assertEquals(200, events.get(1).get(leaderTopicIdPartition.topicPartition())); - // cancel the task so that we don't delete the second segment - task.cancel(); - return CompletableFuture.runAsync(() -> { }); - }); + Map logProps = new HashMap<>(); + logProps.put("retention.bytes", -1L); + logProps.put("retention.ms", 0L); + LogConfig mockLogConfig = new LogConfig(logProps); + when(mockLog.config()).thenReturn(mockLogConfig); - task.run(); + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> { + // cancel the task so that we don't delete the second segment + task.cancel(); + return CompletableFuture.runAsync(() -> { + }); + }); - verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); - verify(remoteStorageManager, never()).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + task.run(); + + assertEquals(200L, logStartOffset.get()); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager, never()).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + } } private List listRemoteLogSegmentMetadata(TopicIdPartition topicIdPartition, From e96e0964988324bdc1d32394d0421ef445996f62 Mon Sep 17 00:00:00 2001 From: Christo Date: Mon, 11 Sep 2023 14:56:09 +0100 Subject: [PATCH 3/3] Address comments from third round of review --- .../kafka/log/remote/RemoteLogManager.java | 9 ++--- .../log/remote/RemoteLogManagerTest.java | 34 ++++++++++++++++--- 2 files changed, 33 insertions(+), 10 deletions(-) diff --git a/core/src/main/java/kafka/log/remote/RemoteLogManager.java b/core/src/main/java/kafka/log/remote/RemoteLogManager.java index 171bf0f659c27..d8f2144b3e353 100644 --- a/core/src/main/java/kafka/log/remote/RemoteLogManager.java +++ b/core/src/main/java/kafka/log/remote/RemoteLogManager.java @@ -1001,9 +1001,6 @@ private void cleanupExpiredRemoteLogSegments() throws RemoteStorageException, Ex // since the checkpoint file was already truncated. boolean shouldDeleteSegment = remoteLogRetentionHandler.deleteLogStartOffsetBreachedSegments( metadata, logStartOffset, epochWithOffsets); - if (shouldDeleteSegment) { - segmentsToDelete.add(metadata); - } boolean isValidSegment = false; if (!shouldDeleteSegment) { // check whether the segment contains the required epoch range with in the current leader epoch lineage. @@ -1012,11 +1009,11 @@ private void cleanupExpiredRemoteLogSegments() throws RemoteStorageException, Ex shouldDeleteSegment = remoteLogRetentionHandler.deleteRetentionTimeBreachedSegments(metadata) || remoteLogRetentionHandler.deleteRetentionSizeBreachedSegments(metadata); - if (shouldDeleteSegment) { - segmentsToDelete.add(metadata); - } } } + if (shouldDeleteSegment) { + segmentsToDelete.add(metadata); + } canProcess = shouldDeleteSegment || !isValidSegment; } } diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index 898e610cc1bfc..bb66994b273b4 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -1626,8 +1626,8 @@ public RemoteLogMetadataManager createRemoteLogMetadataManager() { return remoteLogMetadataManager; } }) { - RemoteLogManager.RLMTask task = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); - task.convertToLeader(0); + RemoteLogManager.RLMTask leaderTask = remoteLogManager.new RLMTask(leaderTopicIdPartition, 128); + leaderTask.convertToLeader(0); when(mockLog.topicPartition()).thenReturn(leaderTopicIdPartition.topicPartition()); when(mockLog.logEndOffset()).thenReturn(200L); @@ -1655,16 +1655,42 @@ public RemoteLogMetadataManager createRemoteLogMetadataManager() { when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) .thenAnswer(answer -> { // cancel the task so that we don't delete the second segment - task.cancel(); + leaderTask.cancel(); return CompletableFuture.runAsync(() -> { }); }); - task.run(); + leaderTask.run(); assertEquals(200L, logStartOffset.get()); verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); verify(remoteStorageManager, never()).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); + + // test that the 2nd log segment will be deleted by the new leader + RemoteLogManager.RLMTask newLeaderTask = remoteLogManager.new RLMTask(followerTopicIdPartition, 128); + newLeaderTask.convertToLeader(1); + + Iterator firstIterator = remoteLogSegmentMetadatas.iterator(); + firstIterator.next(); + Iterator secondIterator = remoteLogSegmentMetadatas.iterator(); + secondIterator.next(); + Iterator thirdIterator = remoteLogSegmentMetadatas.iterator(); + thirdIterator.next(); + + when(remoteLogMetadataManager.listRemoteLogSegments(followerTopicIdPartition)) + .thenReturn(firstIterator); + when(remoteLogMetadataManager.listRemoteLogSegments(followerTopicIdPartition, 0)) + .thenReturn(secondIterator) + .thenReturn(thirdIterator); + + when(remoteLogMetadataManager.updateRemoteLogSegmentMetadata(any(RemoteLogSegmentMetadataUpdate.class))) + .thenAnswer(answer -> CompletableFuture.runAsync(() -> { })); + + newLeaderTask.run(); + + assertEquals(200L, logStartOffset.get()); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(0)); + verify(remoteStorageManager).deleteLogSegmentData(remoteLogSegmentMetadatas.get(1)); } }