Skip to content
Open
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 @@ -30,7 +30,6 @@
import org.apache.fluss.server.log.LogManager;
import org.apache.fluss.server.zk.ZooKeeperClient;
import org.apache.fluss.server.zk.data.TableRegistration;
import org.apache.fluss.utils.FileUtils;
import org.apache.fluss.utils.FlussPaths;
import org.apache.fluss.utils.concurrent.ExecutorThreadFactory;

Expand All @@ -40,9 +39,9 @@
import java.io.File;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.LinkOption;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.LinkedHashMap;
import java.util.List;
Expand All @@ -55,7 +54,6 @@
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Collectors;

import static org.apache.fluss.utils.FlussPaths.HISTORICAL_LOOKUP_CACHE_DIR_NAME;
import static org.apache.fluss.utils.FlussPaths.KV_TABLET_DIR_PREFIX;
Expand Down Expand Up @@ -112,37 +110,26 @@ protected Map<File, List<File>> listTabletsToLoad() {
return tabletsToLoadByDataDir;
}

/** Returns the tablet directories to be loaded from a single configured data directory. */
protected List<File> listTabletsToLoad(File dataDir) {
return listTabletsToLoad(
dataDir,
(directory, nameFilter) ->
Arrays.stream(FileUtils.listDirectories(directory))
.filter(file -> nameFilter.test(file.getName()))
.collect(Collectors.toList()));
}

/**
* Lists tablet directories using the common layout and cache exclusions, with a caller-supplied
* directory lister to control symbolic-link handling and listing failures.
* Returns the tablet directories from a single configured data directory, skipping caches,
* symbolic links, and unreadable directories.
*/
protected <E extends Exception> List<File> listTabletsToLoad(
File dataDir, DirectoryLister<E> directoryLister) throws E {
protected List<File> listTabletsToLoad(File dataDir) {
List<File> tabletsToLoad = new ArrayList<>();
for (File dbDir :
directoryLister.listDirectories(
listDirectories(
dataDir,
name ->
!name.equals(HISTORICAL_LOOKUP_CACHE_DIR_NAME)
&& !name.equals(REMOTE_LOG_INDEX_LOCAL_CACHE))) {
for (File tableDir : directoryLister.listDirectories(dbDir, name -> true)) {
for (File tableDir : listDirectories(dbDir, name -> true)) {
for (File tabletOrPartitionDir :
directoryLister.listDirectories(
listDirectories(
tableDir,
name -> isPartitionDir(name) || name.startsWith(tabletDirPrefix))) {
if (isPartitionDir(tabletOrPartitionDir.getName())) {
tabletsToLoad.addAll(
directoryLister.listDirectories(
listDirectories(
tabletOrPartitionDir,
name -> name.startsWith(tabletDirPrefix)));
} else {
Expand All @@ -154,14 +141,22 @@ protected <E extends Exception> List<File> listTabletsToLoad(
return tabletsToLoad;
}

/**
* Lists child directories whose names match a filter, optionally reporting listing failures.
*/
@FunctionalInterface
protected interface DirectoryLister<E extends Exception> {

/** Returns the matching child directories. */
List<File> listDirectories(File parent, Predicate<String> nameFilter) throws E;
private List<File> listDirectories(File parent, Predicate<String> nameFilter) {
List<File> directories = new ArrayList<>();
File[] entries = parent.listFiles();
if (entries == null) {
LOG.warn("Failed to list tablet directories in {}. Skipping this directory.", parent);
return directories;
}
for (File entry : entries) {
if (!nameFilter.test(entry.getName())) {
continue;
}
if (Files.isDirectory(entry.toPath(), LinkOption.NOFOLLOW_LINKS)) {
directories.add(entry);
}
}
return directories;
}

protected ExecutorService createThreadPool(String poolName) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,22 +67,17 @@

import java.io.File;
import java.io.IOException;
import java.nio.file.DirectoryIteratorException;
import java.nio.file.DirectoryNotEmptyException;
import java.nio.file.DirectoryStream;
import java.nio.file.Files;
import java.nio.file.LinkOption;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.nio.file.attribute.BasicFileAttributes;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Predicate;

import static org.apache.fluss.utils.Preconditions.checkState;
import static org.apache.fluss.utils.concurrent.LockUtils.inLock;
Expand Down Expand Up @@ -395,9 +390,7 @@ private void cleanupStaleKvDirectories() {
for (File dataDir : dataDirs) {
try {
Path realDataDir = dataDir.toPath().toRealPath();
List<File> staleDirs =
listTabletsToLoad(
realDataDir.toFile(), this::listCleanupDirectories);
List<File> staleDirs = listTabletsToLoad(realDataDir.toFile());
int deletedDirectories = 0;
for (File tabletDir : staleDirs) {
try {
Expand Down Expand Up @@ -449,34 +442,6 @@ private void deleteStaleKvDirectory(Path tabletDir, Path dataDir) throws IOExcep
deleteEmptyParentDirectories(deletedDir.getParent(), dataDir);
}

private List<File> listCleanupDirectories(File parent, Predicate<String> nameFilter)
throws IOException {
List<File> directories = new ArrayList<>();
try (DirectoryStream<Path> entries = Files.newDirectoryStream(parent.toPath())) {
for (Path entry : entries) {
if (!nameFilter.test(entry.getFileName().toString())) {
continue;
}
BasicFileAttributes attributes =
Files.readAttributes(
entry, BasicFileAttributes.class, LinkOption.NOFOLLOW_LINKS);
if (attributes.isSymbolicLink()) {
LOG.warn(
"Skipping symbolic link {} during stale KV cleanup; its target will "
+ "not be cleaned.",
entry);
continue;
}
if (attributes.isDirectory()) {
directories.add(entry.toFile());
}
}
} catch (DirectoryIteratorException e) {
throw e.getCause();
}
return directories;
}

private void deleteEmptyParentDirectories(Path directory, Path dataDir) throws IOException {
for (Path parent = directory;
parent != null && parent.startsWith(dataDir) && !parent.equals(dataDir);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,23 +48,10 @@ final class TabletManagerBaseTest {

@TempDir private File tempDir;

@Test
void testIgnoresHistoricalLookupCacheDirectoryWhenLoadingTablets() {
File fakeTabletDir =
new File(
new File(FlussPaths.historicalLookupRootDir(tempDir), "database"),
"kv-table-1");
assertThat(fakeTabletDir.mkdirs()).isTrue();

TestingTabletManager tabletManager = new TestingTabletManager(tempDir);

assertThat(tabletManager.tabletsToLoad(tempDir)).isEmpty();
}

@ParameterizedTest
@EnumSource(TabletManagerBase.TabletType.class)
void testListsTabletLayoutAndExcludesCaches(TabletManagerBase.TabletType tabletType)
throws Exception {
void testListsTabletLayoutAndExcludesCachesAndSymbolicLinks(
TabletManagerBase.TabletType tabletType, @TempDir Path outsideDir) throws Exception {
for (String path :
Arrays.asList(
"db/table-1/kv-0",
Expand All @@ -78,6 +65,10 @@ void testListsTabletLayoutAndExcludesCaches(TabletManagerBase.TabletType tabletT
FlussPaths.REMOTE_LOG_INDEX_LOCAL_CACHE + "/table-1/log-0")) {
Files.createDirectories(tempDir.toPath().resolve(path));
}
Files.createDirectories(outsideDir.resolve("table-3/kv-0"));
Files.createDirectories(outsideDir.resolve("table-3/log-0"));
Files.createSymbolicLink(tempDir.toPath().resolve("linked-db"), outsideDir);

String prefix =
tabletType == TabletManagerBase.TabletType.KV
? FlussPaths.KV_TABLET_DIR_PREFIX
Expand All @@ -90,18 +81,6 @@ void testListsTabletLayoutAndExcludesCaches(TabletManagerBase.TabletType tabletT
new File(tempDir, "db/table-2/partition-p2/" + prefix + "1"));
}

@Test
void testLogLoadingStillFollowsSymbolicDatabaseDirectory(@TempDir Path outsideDir)
throws Exception {
Files.createDirectories(outsideDir.resolve("table-1/log-0"));
Path dbLink = Files.createSymbolicLink(tempDir.toPath().resolve("db"), outsideDir);
TestingTabletManager tabletManager =
new TestingTabletManager(tempDir, TabletManagerBase.TabletType.LOG);

assertThat(tabletManager.tabletsToLoad(tempDir))
.containsExactly(dbLink.resolve("table-1/log-0").toFile());
}

@Test
void testCloseTabletsConcurrentlyWaitsForAllTasksAndShutsDownPoolOnFailure() throws Exception {
TestingTabletManager tabletManager = new TestingTabletManager(2);
Expand Down Expand Up @@ -155,10 +134,6 @@ private TestingTabletManager(int closingThreads) {
super(TabletType.KV, Collections.emptyList(), new Configuration(), closingThreads);
}

private TestingTabletManager(File dataDir) {
this(dataDir, TabletType.KV);
}

private TestingTabletManager(File dataDir, TabletType tabletType) {
super(tabletType, Collections.singletonList(dataDir), new Configuration(), 1);
}
Expand Down
Loading
Loading