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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -465,13 +465,17 @@ public synchronized boolean addMessage(long ledgerId, long entryId, long deliver
return true;
}

private synchronized List<ImmutableBucket> selectMergedBuckets(final List<ImmutableBucket> values, int mergeNum) {
checkArgument(mergeNum < values.size());
@VisibleForTesting
synchronized List<ImmutableBucket> selectMergedBuckets(final List<ImmutableBucket> values, int mergeNum) {
if (values.size() < 2 || mergeNum < 2) {
return Collections.emptyList();
}
int actualMergeNum = Math.min(mergeNum, values.size());
long minNumberMessages = Long.MAX_VALUE;
long minScheduleTimestamp = Long.MAX_VALUE;
int minIndex = -1;
for (int i = 0; i + (mergeNum - 1) < values.size(); i++) {
List<ImmutableBucket> immutableBuckets = values.subList(i, i + mergeNum);
for (int i = 0; i + (actualMergeNum - 1) < values.size(); i++) {
List<ImmutableBucket> immutableBuckets = values.subList(i, i + actualMergeNum);
if (immutableBuckets.stream().allMatch(bucket -> {
// We should skip the bucket which last segment already been load to memory,
// avoid record replicated index.
Expand All @@ -482,8 +486,11 @@ private synchronized List<ImmutableBucket> selectMergedBuckets(final List<Immuta
.sum();
if (numberMessages <= minNumberMessages) {
minNumberMessages = numberMessages;
// Snapshot segment IDs start at 1 while the timestamp list is zero-based.
// The next unloaded segment is currentSegmentEntryId + 1, at list index
// currentSegmentEntryId.
long scheduleTimestamp = immutableBuckets.stream()
.mapToLong(bucket -> bucket.firstScheduleTimestamps.get(bucket.currentSegmentEntryId + 1))
.mapToLong(bucket -> bucket.firstScheduleTimestamps.get(bucket.currentSegmentEntryId))
.min().getAsLong();
if (scheduleTimestamp < minScheduleTimestamp) {
minScheduleTimestamp = scheduleTimestamp;
Expand All @@ -494,9 +501,9 @@ private synchronized List<ImmutableBucket> selectMergedBuckets(final List<Immuta
}

if (minIndex >= 0) {
return values.subList(minIndex, minIndex + mergeNum);
} else if (mergeNum > 2){
return selectMergedBuckets(values, mergeNum - 1);
return values.subList(minIndex, minIndex + actualMergeNum);
} else if (actualMergeNum > 2) {
return selectMergedBuckets(values, actualMergeNum - 1);
} else {
return Collections.emptyList();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import java.lang.reflect.Method;
import java.nio.ByteBuffer;
import java.time.Clock;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.NavigableMap;
Expand Down Expand Up @@ -565,6 +566,18 @@ public CompletableFuture<Void> deleteBucketSnapshot(long bucketId) {
}
}

private ImmutableBucket createMergeableBucket(TrackerWithStorage trackerWithStorage, long startLedgerId,
long endLedgerId, List<Long> firstScheduleTimestamps) {
MutableBucket mutableBucket = trackerWithStorage.tracker.getLastMutableBucket();
ImmutableBucket bucket = new ImmutableBucket(mutableBucket.dispatcherName, mutableBucket.cursor,
mutableBucket.sequencer, mutableBucket.bucketSnapshotStorage, startLedgerId, endLedgerId);
bucket.setCurrentSegmentEntryId(1);
bucket.setLastSegmentEntryId(firstScheduleTimestamps.size());
bucket.setFirstScheduleTimestamps(firstScheduleTimestamps);
bucket.setNumberBucketDelayedMessages(1);
return bucket;
}

private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, int maxNumBuckets)
throws Exception {
return createTrackerWithMockLedger(firstLedgerId, maxNumBuckets, new MockBucketSnapshotStorage());
Expand Down Expand Up @@ -606,6 +619,48 @@ public Position getMarkDeletedPosition() {
return new TrackerWithStorage(tracker, storage, mockClockTime);
}

@DataProvider(name = "smallMaxNumBuckets")
private Object[][] smallMaxNumBuckets() {
return new Object[][]{{1}, {2}, {3}};
}

@Test(dataProvider = "smallMaxNumBuckets")
public void testMergeSupportsSmallMaxNumBuckets(int maxNumBuckets) throws Exception {
TrackerWithStorage ts = createTrackerWithMockLedger(0L, maxNumBuckets);
int messageCount = (maxNumBuckets + 1) * 5 + 1;
NavigableSet<Position> expectedMessages = new TreeSet<>();
try {
for (int i = 1; i <= messageCount; i++) {
assertTrue(ts.tracker.addMessage(i, i, i % 5 == 0 ? 20L : 10L));
expectedMessages.add(PositionFactory.create(i, i));
}

Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
synchronized (ts.tracker) {
List<ImmutableBucket> buckets = List.copyOf(
ts.tracker.getImmutableBuckets().asMapOfRanges().values());
assertTrue(!buckets.isEmpty());
assertTrue(buckets.size() <= maxNumBuckets);
assertTrue(buckets.stream().noneMatch(bucket -> bucket.merging
|| bucket.getSnapshotCreateFuture()
.map(future -> !future.isDone() || future.isCompletedExceptionally())
.orElse(true)));
}
});

ts.clockTime.set(20L);
List<Position> scheduledMessages = new ArrayList<>();
Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
scheduledMessages.addAll(ts.tracker.getScheduledMessages(expectedMessages.size()));
assertEquals(scheduledMessages.size(), expectedMessages.size());
});
assertEquals(new TreeSet<>(scheduledMessages), expectedMessages);
assertEquals(ts.tracker.getNumberOfDelayedMessages(), 0L);
} finally {
ts.close();
}
}

@Test
public void testTrimRemovesOrphanedBuckets() throws Exception {
long firstLedgerId = 31L;
Expand Down Expand Up @@ -637,6 +692,50 @@ public void testTrimRemovesOrphanedBuckets() throws Exception {
ts.close();
}

@Test
public void testSelectMergedBucketsSupportsTwoBuckets() throws Exception {
TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1);
try {
ImmutableBucket firstBucket = createMergeableBucket(ts, 1L, 1L, List.of(10L, 20L));
ImmutableBucket secondBucket = createMergeableBucket(ts, 2L, 2L, List.of(10L, 20L));

assertEquals(ts.tracker.selectMergedBuckets(List.of(firstBucket, secondBucket), 4),
List.of(firstBucket, secondBucket));
} finally {
ts.close();
}
}

@Test
public void testSelectMergedBucketsHandlesOneUnloadedSegment() throws Exception {
TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1);
try {
ImmutableBucket firstBucket = createMergeableBucket(ts, 1L, 1L, List.of(10L, 20L));
ImmutableBucket secondBucket = createMergeableBucket(ts, 2L, 2L, List.of(10L, 20L));
ImmutableBucket thirdBucket = createMergeableBucket(ts, 3L, 3L, List.of(10L, 20L));

assertEquals(ts.tracker.selectMergedBuckets(List.of(firstBucket, secondBucket, thirdBucket), 2),
List.of(firstBucket, secondBucket));
} finally {
ts.close();
}
}

@Test
public void testSelectMergedBucketsUsesNextUnloadedSegmentTimestamp() throws Exception {
TrackerWithStorage ts = createTrackerWithMockLedger(0L, 1);
try {
ImmutableBucket firstBucket = createMergeableBucket(ts, 1L, 1L, List.of(10L, 100L, 10L));
ImmutableBucket secondBucket = createMergeableBucket(ts, 2L, 2L, List.of(10L, 100L, 10L));
ImmutableBucket thirdBucket = createMergeableBucket(ts, 3L, 3L, List.of(10L, 50L, 1000L));

assertEquals(ts.tracker.selectMergedBuckets(List.of(firstBucket, secondBucket, thirdBucket), 2),
List.of(secondBucket, thirdBucket));
} finally {
ts.close();
}
}

@Test
public void testTrimDoesNotDeleteBucketOverlappingFirstActiveLedger() throws Exception {
RecordingDeleteStorage storage = new RecordingDeleteStorage();
Expand Down
Loading