From ea3932bbe3c856d3e0fe089ed6fe7dba0f2959c3 Mon Sep 17 00:00:00 2001 From: Will Perlichek Date: Fri, 8 Nov 2024 07:56:07 -0800 Subject: [PATCH 1/4] Refactor RLM test to reduce flakiness --- .../log/remote/RemoteLogManagerTest.java | 29 +++++++++++++++---- 1 file changed, 24 insertions(+), 5 deletions(-) diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index e75a6ca85d40c..f406daaade834 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -1717,9 +1717,28 @@ void testFetchOffsetByTimestampWithTieredStorageDoesNotFetchIndexWhenExistsLocal FileRecords.TimestampAndOffset expectedRemoteResult = new FileRecords.TimestampAndOffset(timestamp + 999, 999, Optional.of(Integer.MAX_VALUE)); Partition mockFollowerPartition = mockPartition(tpId); - LogSegment logSegment = mockLogSegment(50L, timestamp, null); - LogSegment logSegment1 = mockLogSegment(100L, timestamp + 1, expectedLocalResult); - when(mockLog.logSegments()).thenReturn(Arrays.asList(logSegment, logSegment1)); + LogSegment logSegmentBaseOffset50 = mockLogSegment(50L, timestamp, null); + LogSegment logSegmentBaseOffset100 = mockLogSegment(100L, timestamp + 1, expectedLocalResult); + LogSegment logSegmentBaseOffset101 = mockLogSegment(101L, timestamp + 1, expectedLocalResult); + + final long twoSegmentsBaseOffsets50and100 = 0L; + final long oneSegmentBaseOffset100 = 1L; + final long oneSegmentBaseOffset101 = 2L; + + AtomicLong localLogOffsetState = new AtomicLong(twoSegmentsBaseOffsets50and100); + + when(mockLog.logSegments()).thenAnswer(invocation -> { + if (localLogOffsetState.get() == twoSegmentsBaseOffsets50and100) { + return Arrays.asList(logSegmentBaseOffset50, logSegmentBaseOffset100); + } else if (localLogOffsetState.get() == oneSegmentBaseOffset100) { + return Collections.singletonList(logSegmentBaseOffset100); + } else if (localLogOffsetState.get() == oneSegmentBaseOffset101) { + return Collections.singletonList(logSegmentBaseOffset101); + } else { + throw new IllegalStateException("Unexpected LocalLogOffsetState: " + localLogOffsetState.get()); + } + }); + when(mockLog.logEndOffset()).thenReturn(300L); remoteLogManager = new RemoteLogManager(config.remoteLogManagerConfig(), brokerId, logDir, clusterId, time, partition -> Optional.of(mockLog), @@ -1745,12 +1764,12 @@ Optional lookupTimestamp(RemoteLogSegmentMetadat // Move the local-log start offset to 100L, still the read from the remote storage should be short-circuited // as the message with (timestamp + 1) exists in the local log - when(mockLog.logSegments()).thenReturn(Collections.singletonList(logSegment1)); + localLogOffsetState.set(oneSegmentBaseOffset100); assertEquals(Optional.of(expectedLocalResult), remoteLogManager.findOffsetByTimestamp(tp, timestamp + 1, 0L, cache)); // Move the local log start offset to 101L, now message with (timestamp + 1) does not exist in the local log and // the indexes needs to be fetched from the remote storage - when(logSegment1.baseOffset()).thenReturn(101L); + localLogOffsetState.set(oneSegmentBaseOffset101); assertEquals(Optional.of(expectedRemoteResult), remoteLogManager.findOffsetByTimestamp(tp, timestamp + 1, 0L, cache)); } From 3c71922b358f96448d70d1092fa606518b08c21f Mon Sep 17 00:00:00 2001 From: Will Perlichek Date: Fri, 8 Nov 2024 14:24:30 -0800 Subject: [PATCH 2/4] Simplify IllegalStateException message for localLogOffsetState --- 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 f406daaade834..a95600e0707f6 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -1735,7 +1735,7 @@ void testFetchOffsetByTimestampWithTieredStorageDoesNotFetchIndexWhenExistsLocal } else if (localLogOffsetState.get() == oneSegmentBaseOffset101) { return Collections.singletonList(logSegmentBaseOffset101); } else { - throw new IllegalStateException("Unexpected LocalLogOffsetState: " + localLogOffsetState.get()); + throw new IllegalStateException("Unexpected localLogOffsetState"); } }); From 32d5305f05ed95ffc78a0c2df14e516df9d60f5f Mon Sep 17 00:00:00 2001 From: Will Perlichek Date: Fri, 8 Nov 2024 14:46:23 -0800 Subject: [PATCH 3/4] Use AtomicInteger for test state to not confuse with long offsets; comments --- .../java/kafka/log/remote/RemoteLogManagerTest.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index a95600e0707f6..2db3ab46444b3 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -127,6 +127,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiConsumer; import java.util.function.Supplier; import java.util.stream.Collectors; @@ -1721,11 +1722,12 @@ void testFetchOffsetByTimestampWithTieredStorageDoesNotFetchIndexWhenExistsLocal LogSegment logSegmentBaseOffset100 = mockLogSegment(100L, timestamp + 1, expectedLocalResult); LogSegment logSegmentBaseOffset101 = mockLogSegment(101L, timestamp + 1, expectedLocalResult); - final long twoSegmentsBaseOffsets50and100 = 0L; - final long oneSegmentBaseOffset100 = 1L; - final long oneSegmentBaseOffset101 = 2L; + // Constants representing the states of local log segments + final int twoSegmentsBaseOffsets50and100 = 0; + final int oneSegmentBaseOffset100 = 1; + final int oneSegmentBaseOffset101 = 2; - AtomicLong localLogOffsetState = new AtomicLong(twoSegmentsBaseOffsets50and100); + AtomicInteger localLogOffsetState = new AtomicInteger(twoSegmentsBaseOffsets50and100); when(mockLog.logSegments()).thenAnswer(invocation -> { if (localLogOffsetState.get() == twoSegmentsBaseOffsets50and100) { From c2297be7b479a93e348fab3469e5efe7917b1e7e Mon Sep 17 00:00:00 2001 From: Will Perlichek Date: Fri, 8 Nov 2024 20:37:29 -0800 Subject: [PATCH 4/4] spotlessApply to fix imports --- 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 2db3ab46444b3..dcf73ce3d4a29 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -126,8 +126,8 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.BiConsumer; import java.util.function.Supplier; import java.util.stream.Collectors;