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 @@ -377,7 +377,6 @@ private Pair<Path, Boolean> spillTableDiffIntoSstFile(List<Path> deltaFilePaths,
* @param previousSnapshotInfo information about the previous snapshot.
* @param snapshotInfo information about the current snapshot for which
* incremental defragmentation is performed.
* @param snapshotVersion the version of the snapshot to be processed.
* @param checkpointStore the dbStore instance where data
* updates are ingested after being processed.
* @param bucketPrefixInfo table prefix information associated with buckets,
Expand All @@ -387,7 +386,7 @@ private Pair<Path, Boolean> spillTableDiffIntoSstFile(List<Path> deltaFilePaths,
*/
@VisibleForTesting
void performIncrementalDefragmentation(SnapshotInfo previousSnapshotInfo, SnapshotInfo snapshotInfo,
int snapshotVersion, DBStore checkpointStore, TablePrefixInfo bucketPrefixInfo, Set<String> incrementalTables)
DBStore checkpointStore, TablePrefixInfo bucketPrefixInfo, Set<String> incrementalTables)
throws IOException {
// Map of delta files grouped on the basis of the tableName.
Collection<Pair<Path, SstFileInfo>> allTableDeltaFiles = this.deltaDiffComputer.getDeltaFiles(
Expand All @@ -411,22 +410,18 @@ void performIncrementalDefragmentation(SnapshotInfo previousSnapshotInfo, Snapsh
for (Map.Entry<String, List<Path>> entry : tableGroupedDeltaFiles.entrySet()) {
String table = entry.getKey();
List<Path> deltaFiles = entry.getValue();
Path fileToBeIngested;
if (deltaFiles.size() == 1 && snapshotVersion > 0) {
// If there is only one delta file for the table and the snapshot version is also not 0 then the same delta
// file can reingested into the checkpointStore.
fileToBeIngested = deltaFiles.get(0);
} else {
Table<String, CodecBuffer> snapshotTable = snapshot.get().getMetadataManager().getStore()
.getTable(table, StringCodec.get(), CodecBufferCodec.get(true));
Table<String, CodecBuffer> previousSnapshotTable = previousSnapshot.get().getMetadataManager().getStore()
.getTable(table, StringCodec.get(), CodecBufferCodec.get(true));
String tableBucketPrefix = bucketPrefixInfo.getTablePrefix(table);
Pair<Path, Boolean> spillResult = spillTableDiffIntoSstFile(deltaFiles, snapshotTable,
previousSnapshotTable, tableBucketPrefix);
fileToBeIngested = spillResult.getValue() ? spillResult.getLeft() : null;
filesToBeDeleted.add(spillResult.getLeft());
}
// Delta candidates are live RocksDB SSTs selected at file granularity,
// not valid external SSTs containing an exact bucket-level delta. Always
// rebuild the logical delta, even when there is only one candidate file.
Table<String, CodecBuffer> snapshotTable = snapshot.get().getMetadataManager().getStore()
.getTable(table, StringCodec.get(), CodecBufferCodec.get(true));
Table<String, CodecBuffer> previousSnapshotTable = previousSnapshot.get().getMetadataManager().getStore()
.getTable(table, StringCodec.get(), CodecBufferCodec.get(true));
String tableBucketPrefix = bucketPrefixInfo.getTablePrefix(table);
Pair<Path, Boolean> spillResult = spillTableDiffIntoSstFile(deltaFiles, snapshotTable,
previousSnapshotTable, tableBucketPrefix);
Path fileToBeIngested = spillResult.getValue() ? spillResult.getLeft() : null;
filesToBeDeleted.add(spillResult.getLeft());
if (fileToBeIngested != null) {
if (!fileToBeIngested.toFile().exists()) {
throw new IOException("Delta file does not exist: " + fileToBeIngested);
Expand Down Expand Up @@ -687,8 +682,8 @@ boolean checkAndDefragSnapshot(SnapshotChainManager chainManager, UUID snapshotI
LOG.info("Performing incremental defragmentation for snapshot: {} (ID: {})", snapshotInfo.getTableKey(),
snapshotInfo.getSnapshotId());
try {
performIncrementalDefragmentation(checkpointSnapshotInfo, snapshotInfo, needsDefragVersionPair.getValue(),
checkpointDBStore, prefixInfo, COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT);
performIncrementalDefragmentation(checkpointSnapshotInfo, snapshotInfo, checkpointDBStore, prefixInfo,
COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT);
perfMetrics.setSnapshotDefragServiceIncLatencyMs(Time.monotonicNow() - defragStart);
} catch (IOException e) {
snapshotMetrics.incNumSnapshotIncDefragFails();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,9 @@
import java.io.File;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
Expand All @@ -61,6 +63,7 @@
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Optional;
import java.util.Set;
import java.util.UUID;
Expand All @@ -82,9 +85,11 @@
import org.apache.hadoop.hdds.utils.db.CodecException;
import org.apache.hadoop.hdds.utils.db.DBCheckpoint;
import org.apache.hadoop.hdds.utils.db.DBStore;
import org.apache.hadoop.hdds.utils.db.DBStoreBuilder;
import org.apache.hadoop.hdds.utils.db.InMemoryTestTable;
import org.apache.hadoop.hdds.utils.db.ManagedRawSSTFileReader;
import org.apache.hadoop.hdds.utils.db.RDBSstFileWriter;
import org.apache.hadoop.hdds.utils.db.RDBStore;
import org.apache.hadoop.hdds.utils.db.RocksDBCheckpoint;
import org.apache.hadoop.hdds.utils.db.RocksDatabaseException;
import org.apache.hadoop.hdds.utils.db.SstFileSetReader;
Expand Down Expand Up @@ -122,6 +127,7 @@
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.InOrder;
Expand All @@ -130,6 +136,7 @@
import org.mockito.MockedStatic;
import org.mockito.Mockito;
import org.mockito.MockitoAnnotations;
import org.rocksdb.LiveFileMetaData;

/**
* Unit tests for SnapshotDefragService.
Expand Down Expand Up @@ -170,6 +177,11 @@ public class TestSnapshotDefragService {
private Map<String, CodecBuffer> dummyTableValues;
private Set<CodecBuffer> closeSet = new HashSet<>();

private enum LiveSstType {
DB_GENERATED,
PREVIOUSLY_INGESTED
}

@BeforeEach
public void setup() throws IOException {
mocks = MockitoAnnotations.openMocks(this);
Expand Down Expand Up @@ -224,6 +236,83 @@ private String getFromCodecBuffer(CodecBuffer buffer) {
return StringCodec.get().fromCodecBuffer(buffer);
}

private void putString(Table<String, CodecBuffer> table, String key,
String value) throws RocksDatabaseException, CodecException {
table.put(key, StringCodec.get().toDirectCodecBuffer(value));
}

private DBStore createDBStore(String name, String tableName)
throws RocksDatabaseException {
return DBStoreBuilder.newBuilder(configuration)
.setName(name)
.setPath(tempDir)
.addTable(tableName)
.build();
}

private Path createLiveSstDelta(DBStore sourceStore, String tableName,
String key, LiveSstType liveSstType) throws Exception {
Table<String, CodecBuffer> sourceTable = sourceStore.getTable(
tableName, StringCodec.get(), CodecBufferCodec.get(true));
putString(sourceTable, key, "source-value");
sourceStore.flushDB();
Set<String> liveFilesBeforeIngestion = ((RDBStore) sourceStore).getDb()
.getLiveFilesMetaData().stream()
.map(LiveFileMetaData::fileName)
.collect(Collectors.toSet());

if (liveSstType == LiveSstType.PREVIOUSLY_INGESTED) {
File externalFile = tempDir.resolve("external-" + UUID.randomUUID()
+ ".sst").toFile();
try (RDBSstFileWriter writer = new RDBSstFileWriter(externalFile);
CodecBuffer keyBuffer = StringCodec.get().toDirectCodecBuffer(key);
CodecBuffer valueBuffer = StringCodec.get()
.toDirectCodecBuffer("ingested-value")) {
writer.put(keyBuffer, valueBuffer);
}
sourceTable.loadFromFile(externalFile);
}

List<LiveFileMetaData> candidateFiles = ((RDBStore) sourceStore).getDb()
.getLiveFilesMetaData().stream()
.filter(file -> tableName.equals(
StringUtils.bytes2String(file.columnFamilyName())))
.filter(file -> liveSstType != LiveSstType.PREVIOUSLY_INGESTED ||
!liveFilesBeforeIngestion.contains(file.fileName()))
.collect(Collectors.toList());
if (candidateFiles.size() != 1) {
throw new IllegalStateException("Expected one live SST file, found " +
candidateFiles.size());
}
LiveFileMetaData liveFile = candidateFiles.get(0);
Path sourceFile = Paths.get(liveFile.path(), liveFile.fileName());
Path deltaFile = tempDir.resolve("delta-" + UUID.randomUUID() + ".sst");
Files.createLink(deltaFile, sourceFile);
return deltaFile;
}

private void createMockSnapshot(SnapshotInfo snapshotInfo,
DBStore snapshotStore) throws IOException {
OmSnapshot snapshot = mock(OmSnapshot.class);
UncheckedAutoCloseableSupplier<OmSnapshot> snapshotSupplier =
new UncheckedAutoCloseableSupplier<OmSnapshot>() {
@Override
public void close() {
}

@Override
public OmSnapshot get() {
return snapshot;
}
};
OMMetadataManager snapshotMetadataManager = mock(OMMetadataManager.class);
when(snapshot.getMetadataManager()).thenReturn(snapshotMetadataManager);
when(snapshotMetadataManager.getStore()).thenReturn(snapshotStore);
when(omSnapshotManager.getActiveSnapshot(
eq(snapshotInfo.getVolumeName()), eq(snapshotInfo.getBucketName()),
eq(snapshotInfo.getName()))).thenReturn(snapshotSupplier);
}

@AfterEach
public void tearDown() throws Exception {
if (defragService != null) {
Expand Down Expand Up @@ -639,12 +728,87 @@ public OmSnapshot get() {
eq(snapshotInfo.getName()))).thenReturn(snapshotSupplier);
}

@ParameterizedTest
@EnumSource(LiveSstType.class)
public void testRewritesSingleLiveSstBeforeIngestion(
LiveSstType liveSstType) throws Exception {
String tableName = "cf1";
String key = "ab001";
String previousValue = "previous-value";
String currentValue = "current-value";
SnapshotInfo previousSnapshotInfo = createMockSnapshotInfo(
UUID.randomUUID(), "vol1", "bucket1", "snap1");
SnapshotInfo snapshotInfo = createMockSnapshotInfo(
UUID.randomUUID(), "vol1", "bucket1", "snap2");
TablePrefixInfo prefixInfo = new TablePrefixInfo(
ImmutableMap.of(tableName, "ab"));

try (DBStore sourceStore = createDBStore(
"source-" + liveSstType + "-" + UUID.randomUUID(), tableName);
DBStore previousStore = createDBStore(
"previous-" + liveSstType + "-" + UUID.randomUUID(), tableName);
DBStore snapshotStore = createDBStore(
"snapshot-" + liveSstType + "-" + UUID.randomUUID(), tableName);
DBStore checkpointStore = createDBStore(
"checkpoint-" + liveSstType + "-" + UUID.randomUUID(),
tableName)) {
putString(previousStore.getTable(
tableName, StringCodec.get(), CodecBufferCodec.get(true)),
key, previousValue);
previousStore.flushDB();
putString(snapshotStore.getTable(
tableName, StringCodec.get(), CodecBufferCodec.get(true)),
key, currentValue);
snapshotStore.flushDB();
createMockSnapshot(previousSnapshotInfo, previousStore);
createMockSnapshot(snapshotInfo, snapshotStore);
Path deltaFile = createLiveSstDelta(
sourceStore, tableName, key, liveSstType);
when(deltaFileComputer.getDeltaFiles(eq(previousSnapshotInfo),
eq(snapshotInfo), eq(ImmutableSet.of(tableName))))
.thenReturn(ImmutableList.of(Pair.of(deltaFile,
new SstFileInfo(deltaFile.toFile().getName(), key, key,
tableName))));

try (MockedConstruction<SstFileSetReader> ignored = mockConstruction(
SstFileSetReader.class, (mock, context) -> when(mock
.getKeyStreamWithTombstone(anyString(), anyString()))
.thenReturn(new ClosableIterator<String>() {
private boolean hasNext = true;

@Override
public void close() {
}

@Override
public boolean hasNext() {
return hasNext;
}

@Override
public String next() {
if (!hasNext) {
throw new NoSuchElementException();
}
hasNext = false;
return key;
}
}))) {
defragService.performIncrementalDefragmentation(previousSnapshotInfo,
snapshotInfo, checkpointStore, prefixInfo,
ImmutableSet.of(tableName));
}

assertEquals(currentValue, checkpointStore.getTable(
tableName, StringCodec.get(), StringCodec.get()).get(key));
}
}

/**
* Tests the incremental defragmentation process between two snapshots.
*
* <p>This parameterized test validates the {@code performIncrementalDefragmentation} method
* across different version scenarios (0, 1, 2, 10) to ensure proper handling of snapshot
* delta files and version-specific optimizations.</p>
* <p>This test validates that all snapshot delta files are rewritten into
* fresh external SST files before ingestion.</p>
*
* <h3>Test Data Generation:</h3>
* Creates 67,600 synthetic key-value pairs (26×26×100) distributed across two snapshots
Expand All @@ -667,38 +831,26 @@ public OmSnapshot get() {
* <li>SstFileSetReader mock returns keys with indices 0-4 (i % 6 < 5)</li>
* </ul>
*
* <h3>Version-Specific Behavior:</h3>
* <h3>Expected Behavior:</h3>
* <ul>
* <li><b>currentVersion == 0 (initial version):</b>
* <ul>
* <li>All incremental tables are dumped to new SST files</li>
* <li>All dumped files are ingested into the checkpoint database</li>
* </ul>
* </li>
* <li><b>currentVersion > 0 (subsequent versions):</b>
* <ul>
* <li>Single delta file tables (cf1) are ingested directly without merging</li>
* <li>Multiple delta file tables (cf2) are merged and dumped before ingestion</li>
* <li>Optimization: avoids unnecessary file I/O for single delta files</li>
* </ul>
* </li>
* <li>Every incremental table is dumped to a fresh external SST file</li>
* <li>All dumped files are ingested into the checkpoint database</li>
* <li>Behavior is independent of snapshot version and delta-file count</li>
* </ul>
*
* <h3>Assertions:</h3>
* <ul>
* <li>Verifies correct tables are dumped and ingested based on version</li>
* <li>Verifies all incremental tables are dumped and ingested</li>
* <li>Validates that only modified keys (i % 6 < 3) appear in delta files</li>
* <li>Confirms written values match snap2's values or null for deletions</li>
* <li>Ensures all incremental tables are ultimately ingested</li>
* </ul>
*
* @param currentVersion the snapshot version being defragmented (0 for initial, >0 for subsequent)
* @throws Exception if any error occurs during the test execution
*/
@SuppressWarnings("checkstyle:MethodLength")
@ParameterizedTest
@ValueSource(ints = {0, 1, 2, 10})
public void testPerformIncrementalDefragmentation(int currentVersion) throws Exception {
@Test
public void testPerformIncrementalDefragmentation() throws Exception {
DBStore checkpointDBStore = mock(DBStore.class);
String samePrefix = "samePrefix";
String snap1Prefix = "snap1Prefix";
Expand Down Expand Up @@ -839,18 +991,10 @@ public String next() {
String tableName = i.getArgument(0, String.class);
return checkpointTables.get(tableName);
}).when(checkpointDBStore).getTable(anyString());
defragService.performIncrementalDefragmentation(snap1Info, snap2Info, currentVersion, checkpointDBStore,
defragService.performIncrementalDefragmentation(snap1Info, snap2Info, checkpointDBStore,
prefixInfo, incrementalTables);
if (currentVersion == 0) {
assertEquals(incrementalTables, new HashSet<>(dumpedFileName.values()));
assertEquals(ingestedFiles, dumpedFileName);
} else {
assertEquals(ImmutableSet.of("cf2"), new HashSet<>(dumpedFileName.values()));
assertEquals("cf1", ingestedFiles.get(deltaFiles.get(0).getLeft().toAbsolutePath().toString()));
assertEquals(ingestedFiles.entrySet().stream().filter(e -> e.getValue().equals(
"cf2")).collect(Collectors.toSet()), dumpedFileName.entrySet().stream().filter(e -> e.getValue().equals(
"cf2")).collect(Collectors.toSet()));
}
assertEquals(incrementalTables, new HashSet<>(dumpedFileName.values()));
assertEquals(ingestedFiles, dumpedFileName);
assertEquals(incrementalTables, new HashSet<>(ingestedFiles.values()));
for (Map.Entry<Path, List<Pair<String, String>>> deltaFileContent : deltaFileContents.entrySet()) {
int idx = 0;
Expand Down Expand Up @@ -938,7 +1082,7 @@ public void testCheckAndDefragSnapshotFailure(boolean previousSnapshotExists) th
IOException defragException = new IOException("Defrag failed");
if (previousSnapshotExists) {
Mockito.doThrow(defragException).when(spyDefragService).performIncrementalDefragmentation(
eq(previousSnapshotInfo), eq(snapshotInfo), eq(10), eq(checkpointDBStore), eq(prefixInfo),
eq(previousSnapshotInfo), eq(snapshotInfo), eq(checkpointDBStore), eq(prefixInfo),
eq(COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT));
} else {
Mockito.doThrow(defragException).when(spyDefragService).performFullDefragmentation(
Expand Down Expand Up @@ -1045,7 +1189,7 @@ public void testCheckAndDefragActiveSnapshot(boolean previousSnapshotExists) thr
doNothing().when(spyDefragService).performFullDefragmentation(eq(checkpointDBStore), eq(prefixInfo),
eq(COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT));
doNothing().when(spyDefragService).performIncrementalDefragmentation(eq(previousSnapshotInfo),
eq(snapshotInfo), eq(10), eq(checkpointDBStore), eq(prefixInfo),
eq(snapshotInfo), eq(checkpointDBStore), eq(prefixInfo),
eq(COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT));
AtomicInteger lockAcquired = new AtomicInteger(0);
AtomicInteger lockReleased = new AtomicInteger(0);
Expand Down Expand Up @@ -1094,7 +1238,7 @@ public void testCheckAndDefragActiveSnapshot(boolean previousSnapshotExists) thr
eq(COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT));
if (previousSnapshotExists) {
verifier.verify(spyDefragService).performIncrementalDefragmentation(eq(previousSnapshotInfo),
eq(snapshotInfo), eq(10), eq(checkpointDBStore), eq(prefixInfo),
eq(snapshotInfo), eq(checkpointDBStore), eq(prefixInfo),
eq(COLUMN_FAMILIES_TO_TRACK_IN_SNAPSHOT));
} else {
verifier.verify(spyDefragService).performFullDefragmentation(eq(checkpointDBStore), eq(prefixInfo),
Expand Down
Loading