From aae1e2ba18f0ee79b39e26b1bec65dce010a859c Mon Sep 17 00:00:00 2001 From: Christo Lolov Date: Wed, 19 Jul 2023 14:16:05 +0100 Subject: [PATCH 1/4] KAFKA-14038: Optimise calculation of size for log in remote tier --- .../storage/RemoteLogMetadataManager.java | 9 +++ .../storage/NoOpRemoteLogMetadataManager.java | 5 ++ ...ssLoaderAwareRemoteLogMetadataManager.java | 5 ++ .../TopicBasedRemoteLogMetadataManager.java | 11 ++++ ...opicBasedRemoteLogMetadataManagerTest.java | 61 +++++++++++++++++++ ...eLogMetadataManagerWrapperWithHarness.java | 5 ++ .../InmemoryRemoteLogMetadataManager.java | 12 ++++ 7 files changed, 108 insertions(+) diff --git a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java index 9a29746b292f6..b3a9d6ad57ea3 100644 --- a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java +++ b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java @@ -201,4 +201,13 @@ void onPartitionLeadershipChanges(Set leaderPartitions, * @param partitions topic partitions that have been stopped. */ void onStopPartitions(Set partitions); + + /** + * Returns total size of the log for the given leader epoch in remote storage. + * + * @param topicPartition topic partition for which size needs to be calculated. + * @param leaderEpoch Size will only include segments belonging to this epoch. + * @return Total size of the log stored in remote storage in bytes. + */ + Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException; } \ No newline at end of file diff --git a/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java b/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java index 900d5bd5c695b..802053ceb0301 100644 --- a/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java +++ b/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java @@ -74,6 +74,11 @@ public void onPartitionLeadershipChanges(Set leaderPartitions, public void onStopPartitions(Set partitions) { } + @Override + public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) { + return null; + } + @Override public void close() throws IOException { } diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java index 6b46585158957..d9a2acf492c17 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java @@ -101,6 +101,11 @@ public void onStopPartitions(Set partitions) { }); } + @Override + public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + return withClassLoader(() -> delegate.remoteLogSize(topicPartition, leaderEpoch)); + } + @Override public void configure(Map configs) { withClassLoader(() -> { diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java index ffd6e14503935..d8839e6567966 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java @@ -326,6 +326,17 @@ public void onStopPartitions(Set partitions) { } } + @Override + public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + long remoteLogSize = 0L; + Iterator remoteLogSegmentMetadataIterator = remotePartitionMetadataStore.listRemoteLogSegments(topicPartition, leaderEpoch); + while (remoteLogSegmentMetadataIterator.hasNext()) { + RemoteLogSegmentMetadata remoteLogSegmentMetadata = remoteLogSegmentMetadataIterator.next(); + remoteLogSize += remoteLogSegmentMetadata.segmentSizeInBytes(); + } + return remoteLogSize; + } + @Override public void configure(Map configs) { Objects.requireNonNull(configs, "configs can not be null."); diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java index a41a9a3869979..82311828c7a46 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java @@ -25,6 +25,7 @@ import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentId; import org.apache.kafka.server.log.remote.storage.RemoteLogSegmentMetadata; import org.apache.kafka.server.log.remote.storage.RemoteResourceNotFoundException; +import org.apache.kafka.server.log.remote.storage.RemoteStorageException; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; @@ -168,4 +169,64 @@ private void waitUntilConsumerCatchesUp(TopicIdPartition newLeaderTopicIdPartiti } } + @Test + public void testRemoteLogSizeCalculationForUnknownTopicIdPartitionThrows() { + TopicIdPartition topicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition("singleton", 0)); + Assertions.assertThrows(RemoteResourceNotFoundException.class, () -> topicBasedRlmm().remoteLogSize(topicIdPartition, 0)); + } + + @Test + public void testRemoteLogSizeCalculationWithSegmentsOfTheSameEpoch() throws RemoteStorageException, TimeoutException { + TopicIdPartition topicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition("singleton", 0)); + TopicBasedRemoteLogMetadataManager topicBasedRemoteLogMetadataManager = topicBasedRlmm(); + + RemoteLogSegmentMetadata firstSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 0, 100, -1L, 0, time.milliseconds(), SEG_SIZE, Collections.singletonMap(0, 0L)); + RemoteLogSegmentMetadata secondSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 100, 200, -1L, 0, time.milliseconds(), SEG_SIZE * 2, Collections.singletonMap(0, 0L)); + RemoteLogSegmentMetadata thirdSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 200, 300, -1L, 0, time.milliseconds(), SEG_SIZE * 3, Collections.singletonMap(0, 0L)); + + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(firstSegmentMetadata); + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(secondSegmentMetadata); + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(thirdSegmentMetadata); + + topicBasedRemoteLogMetadataManager.onPartitionLeadershipChanges(Collections.singleton(topicIdPartition), Collections.emptySet()); + + // RemoteLogSegmentMetadata events are already published, and topicBasedRlmm's consumer manager will start + // fetching those events and build the cache. + waitUntilConsumerCatchesup(topicIdPartition, topicIdPartition, 30_000L); + + Long remoteLogSize = topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 0); + + Assertions.assertEquals(SEG_SIZE * 6, remoteLogSize); + } + + @Test + public void testRemoteLogSizeCalculationWithSegmentsOfDifferentEpochs() throws RemoteStorageException, TimeoutException { + TopicIdPartition topicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition("singleton", 0)); + TopicBasedRemoteLogMetadataManager topicBasedRemoteLogMetadataManager = topicBasedRlmm(); + + RemoteLogSegmentMetadata firstSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 0, 100, -1L, 0, time.milliseconds(), SEG_SIZE, Collections.singletonMap(0, 0L)); + RemoteLogSegmentMetadata secondSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 100, 200, -1L, 0, time.milliseconds(), SEG_SIZE * 2, Collections.singletonMap(1, 100L)); + RemoteLogSegmentMetadata thirdSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 200, 300, -1L, 0, time.milliseconds(), SEG_SIZE * 3, Collections.singletonMap(2, 200L)); + + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(firstSegmentMetadata); + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(secondSegmentMetadata); + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(thirdSegmentMetadata); + + topicBasedRemoteLogMetadataManager.onPartitionLeadershipChanges(Collections.singleton(topicIdPartition), Collections.emptySet()); + + // RemoteLogSegmentMetadata events are already published, and topicBasedRlmm's consumer manager will start + // fetching those events and build the cache. + waitUntilConsumerCatchesup(topicIdPartition, topicIdPartition, 30_000L); + + Assertions.assertEquals(SEG_SIZE, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 0)); + Assertions.assertEquals(SEG_SIZE * 2, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 1)); + Assertions.assertEquals(SEG_SIZE * 3, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 2)); + } + } diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java index ef9287d31e4df..05824528c5fef 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java @@ -86,6 +86,11 @@ public void onStopPartitions(Set partitions) { remoteLogMetadataManagerHarness.remoteLogMetadataManager().onStopPartitions(partitions); } + @Override + public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + return remoteLogMetadataManagerHarness.remoteLogMetadataManager().remoteLogSize(topicPartition, leaderEpoch); + } + @Override public void close() throws IOException { remoteLogMetadataManagerHarness.remoteLogMetadataManager().close(); diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java index a970de509ad5c..8fc4a30f6af64 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java @@ -156,6 +156,18 @@ public void onStopPartitions(Set partitions) { // this instance. It does not depend upon stopped partitions. } + @Override + public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + long remoteLogSize = 0L; + RemoteLogMetadataCache remoteLogMetadataCache = getRemoteLogMetadataCache(topicPartition); + Iterator remoteLogSegmentMetadataIterator = remoteLogMetadataCache.listAllRemoteLogSegments(); + while (remoteLogSegmentMetadataIterator.hasNext()) { + RemoteLogSegmentMetadata remoteLogSegmentMetadata = remoteLogSegmentMetadataIterator.next(); + remoteLogSize += remoteLogSegmentMetadata.segmentSizeInBytes(); + } + return remoteLogSize; + } + @Override public void close() throws IOException { // Clearing the references to the map and assigning empty immutable maps. From dc69049dd266c912a89e14e21ea5c0fdf0e2686d Mon Sep 17 00:00:00 2001 From: Christo Lolov Date: Wed, 26 Jul 2023 10:47:13 +0100 Subject: [PATCH 2/4] Address review comments --- .../log/remote/storage/RemoteLogMetadataManager.java | 4 ++-- .../remote/storage/NoOpRemoteLogMetadataManager.java | 4 ++-- .../ClassLoaderAwareRemoteLogMetadataManager.java | 2 +- .../storage/TopicBasedRemoteLogMetadataManager.java | 10 ++++++++-- .../TopicBasedRemoteLogMetadataManagerTest.java | 4 ++-- ...asedRemoteLogMetadataManagerWrapperWithHarness.java | 2 +- .../storage/InmemoryRemoteLogMetadataManager.java | 2 +- 7 files changed, 17 insertions(+), 11 deletions(-) diff --git a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java index b3a9d6ad57ea3..9ae36eb00d8aa 100644 --- a/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java +++ b/storage/api/src/main/java/org/apache/kafka/server/log/remote/storage/RemoteLogMetadataManager.java @@ -205,9 +205,9 @@ void onPartitionLeadershipChanges(Set leaderPartitions, /** * Returns total size of the log for the given leader epoch in remote storage. * - * @param topicPartition topic partition for which size needs to be calculated. + * @param topicIdPartition topic partition for which size needs to be calculated. * @param leaderEpoch Size will only include segments belonging to this epoch. * @return Total size of the log stored in remote storage in bytes. */ - Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException; + long remoteLogSize(TopicIdPartition topicIdPartition, int leaderEpoch) throws RemoteStorageException; } \ No newline at end of file diff --git a/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java b/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java index 802053ceb0301..a60c3d408974a 100644 --- a/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java +++ b/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java @@ -75,8 +75,8 @@ public void onStopPartitions(Set partitions) { } @Override - public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) { - return null; + public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) { + return 0; } @Override diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java index d9a2acf492c17..918cbbd422b10 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java @@ -102,7 +102,7 @@ public void onStopPartitions(Set partitions) { } @Override - public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { return withClassLoader(() -> delegate.remoteLogSize(topicPartition, leaderEpoch)); } diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java index d8839e6567966..144393c0882d3 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java @@ -327,15 +327,21 @@ public void onStopPartitions(Set partitions) { } @Override - public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { long remoteLogSize = 0L; + // This is a simple-to-understand but not the most optimal solution. + // The TopicBasedRemoteLogMetadataManager's remote metadata store is file-based. During design discussions + // at https://lists.apache.org/thread/kxd6fffq02thbpd0p5y4mfbs062g7jr6 + // we reached a consensus that sequential iteration over files on the local file system is performant enough. + // Should this stop being the case, the remote log size could be calculated by incrementing/decrementing + // counters during API calls for a more performant implementation. Iterator remoteLogSegmentMetadataIterator = remotePartitionMetadataStore.listRemoteLogSegments(topicPartition, leaderEpoch); while (remoteLogSegmentMetadataIterator.hasNext()) { RemoteLogSegmentMetadata remoteLogSegmentMetadata = remoteLogSegmentMetadataIterator.next(); remoteLogSize += remoteLogSegmentMetadata.segmentSizeInBytes(); } return remoteLogSize; - } + } @Override public void configure(Map configs) { diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java index 82311828c7a46..95f628f6abbda 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java @@ -195,7 +195,7 @@ public void testRemoteLogSizeCalculationWithSegmentsOfTheSameEpoch() throws Remo // RemoteLogSegmentMetadata events are already published, and topicBasedRlmm's consumer manager will start // fetching those events and build the cache. - waitUntilConsumerCatchesup(topicIdPartition, topicIdPartition, 30_000L); + waitUntilConsumerCatchesUp(topicIdPartition, topicIdPartition, 30_000L); Long remoteLogSize = topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 0); @@ -222,7 +222,7 @@ public void testRemoteLogSizeCalculationWithSegmentsOfDifferentEpochs() throws R // RemoteLogSegmentMetadata events are already published, and topicBasedRlmm's consumer manager will start // fetching those events and build the cache. - waitUntilConsumerCatchesup(topicIdPartition, topicIdPartition, 30_000L); + waitUntilConsumerCatchesUp(topicIdPartition, topicIdPartition, 30_000L); Assertions.assertEquals(SEG_SIZE, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 0)); Assertions.assertEquals(SEG_SIZE * 2, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 1)); diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java index 05824528c5fef..3dcd18da25aef 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java @@ -87,7 +87,7 @@ public void onStopPartitions(Set partitions) { } @Override - public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { return remoteLogMetadataManagerHarness.remoteLogMetadataManager().remoteLogSize(topicPartition, leaderEpoch); } diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java index 8fc4a30f6af64..f4a204dba15cc 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java @@ -157,7 +157,7 @@ public void onStopPartitions(Set partitions) { } @Override - public Long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { long remoteLogSize = 0L; RemoteLogMetadataCache remoteLogMetadataCache = getRemoteLogMetadataCache(topicPartition); Iterator remoteLogSegmentMetadataIterator = remoteLogMetadataCache.listAllRemoteLogSegments(); From 5d54c098713ef53449562c5d1b7c031c39eb0940 Mon Sep 17 00:00:00 2001 From: Christo Date: Thu, 27 Jul 2023 10:49:32 +0100 Subject: [PATCH 3/4] Address out-of-range leader epoch comment --- ...opicBasedRemoteLogMetadataManagerTest.java | 22 +++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java index 95f628f6abbda..eaf62edfea7ee 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerTest.java @@ -229,4 +229,26 @@ public void testRemoteLogSizeCalculationWithSegmentsOfDifferentEpochs() throws R Assertions.assertEquals(SEG_SIZE * 3, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 2)); } + @Test + public void testRemoteLogSizeCalculationWithSegmentsHavingNonExistentEpochs() throws RemoteStorageException, TimeoutException { + TopicIdPartition topicIdPartition = new TopicIdPartition(Uuid.randomUuid(), new TopicPartition("singleton", 0)); + TopicBasedRemoteLogMetadataManager topicBasedRemoteLogMetadataManager = topicBasedRlmm(); + + RemoteLogSegmentMetadata firstSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 0, 100, -1L, 0, time.milliseconds(), SEG_SIZE, Collections.singletonMap(0, 0L)); + RemoteLogSegmentMetadata secondSegmentMetadata = new RemoteLogSegmentMetadata(new RemoteLogSegmentId(topicIdPartition, Uuid.randomUuid()), + 100, 200, -1L, 0, time.milliseconds(), SEG_SIZE * 2, Collections.singletonMap(1, 100L)); + + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(firstSegmentMetadata); + topicBasedRemoteLogMetadataManager.addRemoteLogSegmentMetadata(secondSegmentMetadata); + + topicBasedRemoteLogMetadataManager.onPartitionLeadershipChanges(Collections.singleton(topicIdPartition), Collections.emptySet()); + + // RemoteLogSegmentMetadata events are already published, and topicBasedRlmm's consumer manager will start + // fetching those events and build the cache. + waitUntilConsumerCatchesUp(topicIdPartition, topicIdPartition, 30_000L); + + Assertions.assertEquals(0, topicBasedRemoteLogMetadataManager.remoteLogSize(topicIdPartition, 9001)); + } + } From 365f0657a4f8821bb40997b3a8aa3b8f4320e93d Mon Sep 17 00:00:00 2001 From: Christo Date: Thu, 27 Jul 2023 17:33:23 +0100 Subject: [PATCH 4/4] Address renaming miss --- .../log/remote/storage/NoOpRemoteLogMetadataManager.java | 2 +- .../storage/ClassLoaderAwareRemoteLogMetadataManager.java | 4 ++-- .../metadata/storage/TopicBasedRemoteLogMetadataManager.java | 4 ++-- .../TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java | 4 ++-- .../log/remote/storage/InmemoryRemoteLogMetadataManager.java | 4 ++-- 5 files changed, 9 insertions(+), 9 deletions(-) diff --git a/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java b/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java index a60c3d408974a..718815171828a 100644 --- a/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java +++ b/storage/api/src/test/java/org/apache/kafka/server/log/remote/storage/NoOpRemoteLogMetadataManager.java @@ -75,7 +75,7 @@ public void onStopPartitions(Set partitions) { } @Override - public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) { + public long remoteLogSize(TopicIdPartition topicIdPartition, int leaderEpoch) { return 0; } diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java index 918cbbd422b10..663d06275c5aa 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/ClassLoaderAwareRemoteLogMetadataManager.java @@ -102,8 +102,8 @@ public void onStopPartitions(Set partitions) { } @Override - public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { - return withClassLoader(() -> delegate.remoteLogSize(topicPartition, leaderEpoch)); + public long remoteLogSize(TopicIdPartition topicIdPartition, int leaderEpoch) throws RemoteStorageException { + return withClassLoader(() -> delegate.remoteLogSize(topicIdPartition, leaderEpoch)); } @Override diff --git a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java index 144393c0882d3..9f9ef6a63f364 100644 --- a/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java +++ b/storage/src/main/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManager.java @@ -327,7 +327,7 @@ public void onStopPartitions(Set partitions) { } @Override - public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + public long remoteLogSize(TopicIdPartition topicIdPartition, int leaderEpoch) throws RemoteStorageException { long remoteLogSize = 0L; // This is a simple-to-understand but not the most optimal solution. // The TopicBasedRemoteLogMetadataManager's remote metadata store is file-based. During design discussions @@ -335,7 +335,7 @@ public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) thro // we reached a consensus that sequential iteration over files on the local file system is performant enough. // Should this stop being the case, the remote log size could be calculated by incrementing/decrementing // counters during API calls for a more performant implementation. - Iterator remoteLogSegmentMetadataIterator = remotePartitionMetadataStore.listRemoteLogSegments(topicPartition, leaderEpoch); + Iterator remoteLogSegmentMetadataIterator = remotePartitionMetadataStore.listRemoteLogSegments(topicIdPartition, leaderEpoch); while (remoteLogSegmentMetadataIterator.hasNext()) { RemoteLogSegmentMetadata remoteLogSegmentMetadata = remoteLogSegmentMetadataIterator.next(); remoteLogSize += remoteLogSegmentMetadata.segmentSizeInBytes(); diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java index 3dcd18da25aef..e73ac31b16c8f 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/metadata/storage/TopicBasedRemoteLogMetadataManagerWrapperWithHarness.java @@ -87,8 +87,8 @@ public void onStopPartitions(Set partitions) { } @Override - public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { - return remoteLogMetadataManagerHarness.remoteLogMetadataManager().remoteLogSize(topicPartition, leaderEpoch); + public long remoteLogSize(TopicIdPartition topicIdPartition, int leaderEpoch) throws RemoteStorageException { + return remoteLogMetadataManagerHarness.remoteLogMetadataManager().remoteLogSize(topicIdPartition, leaderEpoch); } @Override diff --git a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java index f4a204dba15cc..7cc4552427e1a 100644 --- a/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java +++ b/storage/src/test/java/org/apache/kafka/server/log/remote/storage/InmemoryRemoteLogMetadataManager.java @@ -157,9 +157,9 @@ public void onStopPartitions(Set partitions) { } @Override - public long remoteLogSize(TopicIdPartition topicPartition, int leaderEpoch) throws RemoteStorageException { + public long remoteLogSize(TopicIdPartition topicIdPartition, int leaderEpoch) throws RemoteStorageException { long remoteLogSize = 0L; - RemoteLogMetadataCache remoteLogMetadataCache = getRemoteLogMetadataCache(topicPartition); + RemoteLogMetadataCache remoteLogMetadataCache = getRemoteLogMetadataCache(topicIdPartition); Iterator remoteLogSegmentMetadataIterator = remoteLogMetadataCache.listAllRemoteLogSegments(); while (remoteLogSegmentMetadataIterator.hasNext()) { RemoteLogSegmentMetadata remoteLogSegmentMetadata = remoteLogSegmentMetadataIterator.next();