executors = new ArrayList<>();
Collections.addAll(executors, reclaimExecutor, postFlushExecutor, flushExecutor);
- Collections.addAll(executors, perDiskflushExecutors);
+ perDiskflushExecutors.appendAllExecutors(executors);
ExecutorUtils.shutdownAndWait(timeout, unit, executors);
}
@@ -1108,9 +1097,10 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
{
// flush the memtable
flushRunnables = memtable.flushRunnables(txn);
+ ExecutorService[] executors = perDiskflushExecutors.getExecutorsFor(keyspace.getName(), name);
for (int i = 0; i < flushRunnables.size(); i++)
- futures.add(perDiskflushExecutors[i].submit(flushRunnables.get(i)));
+ futures.add(executors[i].submit(flushRunnables.get(i)));
/**
* we can flush 2is as soon as the barrier completes, as they will be consistent with (or ahead of) the
@@ -2806,4 +2796,88 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean
{
return neverPurgeTombstones;
}
-}
\ No newline at end of file
+
+ /**
+ * The thread pools used to flush memtables.
+ *
+ * Each disk has its own set of thread pools to perform memtable flushes.
+ * Based on the configuration. Local system keyspaces can have their own disk
+ * to allow for special redundancy mechanism. If it is the case the executor services returned for
+ * local system keyspaces will be different from the ones for the other keyspaces.
+ */
+ private static final class PerDiskFlushExecutors
+ {
+ /**
+ * The flush executors for non local system keyspaces.
+ */
+ private final ExecutorService[] nonLocalSystemflushExecutors;
+
+ /**
+ * The flush executors for the local system keyspaces.
+ */
+ private final ExecutorService[] localSystemDiskFlushExecutors;
+
+ /**
+ * {@code true} if local system keyspaces are stored in their own directory and use an extra flush executor,
+ * {@code false} otherwise.
+ */
+ private final boolean useSpecificExecutorForSystemKeyspaces;
+
+ public PerDiskFlushExecutors(int flushWriters,
+ String[] locationsForNonSystemKeyspaces,
+ boolean useSpecificLocationForSystemKeyspaces)
+ {
+ ExecutorService[] flushExecutors = createPerDiskFlushWriters(locationsForNonSystemKeyspaces.length, flushWriters);
+ nonLocalSystemflushExecutors = flushExecutors;
+ useSpecificExecutorForSystemKeyspaces = useSpecificLocationForSystemKeyspaces;
+ localSystemDiskFlushExecutors = useSpecificLocationForSystemKeyspaces ? new ExecutorService[] {newThreadPool("LocalSystemKeyspacesDiskMemtableFlushWriter", flushWriters)}
+ : new ExecutorService[] {flushExecutors[0]};
+ }
+
+ private static ExecutorService[] createPerDiskFlushWriters(int numberOfExecutors, int flushWriters)
+ {
+ ExecutorService[] flushExecutors = new ExecutorService[numberOfExecutors];
+
+ for (int i = 0; i < numberOfExecutors; i++)
+ {
+ flushExecutors[i] = newThreadPool("PerDiskMemtableFlushWriter_" + i, flushWriters);
+ }
+ return flushExecutors;
+ }
+
+ private static JMXEnabledThreadPoolExecutor newThreadPool(String poolName, int size)
+ {
+ return new JMXEnabledThreadPoolExecutor(size,
+ Stage.KEEP_ALIVE_SECONDS,
+ TimeUnit.SECONDS,
+ new LinkedBlockingQueue<>(),
+ new NamedThreadFactory(poolName),
+ "internal");
+ }
+
+ /**
+ * Returns the flush executors for the specified keyspace.
+ *
+ * @param keyspaceName the keyspace name
+ * @param tableName the table name
+ * @return the flush executors that should be used for flushing the memtables of the specified keyspace.
+ */
+ public ExecutorService[] getExecutorsFor(String keyspaceName, String tableName)
+ {
+ return Directories.isStoredInLocalSystemKeyspacesDataLocation(keyspaceName, tableName) ? localSystemDiskFlushExecutors
+ : nonLocalSystemflushExecutors;
+ }
+
+ /**
+ * Appends all the {@code ExecutorService} used for flushes to the collection.
+ *
+ * @param collection the collection to append to.
+ */
+ public void appendAllExecutors(Collection collection)
+ {
+ Collections.addAll(collection, nonLocalSystemflushExecutors);
+ if (useSpecificExecutorForSystemKeyspaces)
+ Collections.addAll(collection, localSystemDiskFlushExecutors);
+ }
+ }
+}
diff --git a/src/java/org/apache/cassandra/db/Directories.java b/src/java/org/apache/cassandra/db/Directories.java
index 9a620a2752..cf4238c67d 100644
--- a/src/java/org/apache/cassandra/db/Directories.java
+++ b/src/java/org/apache/cassandra/db/Directories.java
@@ -17,18 +17,11 @@
*/
package org.apache.cassandra.db;
-import java.io.File;
-import java.io.FileFilter;
-import java.io.IOError;
-import java.io.IOException;
-import java.nio.file.FileStore;
-import java.nio.file.Files;
-import java.nio.file.Path;
-import java.nio.file.Paths;
+import java.io.*;
+import java.nio.file.*;
import java.util.*;
import java.util.concurrent.ThreadLocalRandom;
import java.util.function.BiPredicate;
-import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Iterables;
import com.google.common.collect.Maps;
@@ -42,9 +35,11 @@ import org.apache.cassandra.config.*;
import org.apache.cassandra.db.lifecycle.LifecycleTransaction;
import org.apache.cassandra.io.FSDiskFullWriteError;
import org.apache.cassandra.io.FSError;
+import org.apache.cassandra.io.FSNoDiskAvailableForWriteError;
import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.sstable.*;
+import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.DirectorySizeCalculator;
import org.apache.cassandra.utils.FBUtilities;
@@ -93,15 +88,11 @@ public class Directories
public static final String TMP_SUBDIR = "tmp";
public static final String SECONDARY_INDEX_NAME_SEPARATOR = ".";
- public static final DataDirectory[] dataDirectories;
-
- static
- {
- String[] locations = DatabaseDescriptor.getAllDataFileLocations();
- dataDirectories = new DataDirectory[locations.length];
- for (int i = 0; i < locations.length; ++i)
- dataDirectories[i] = new DataDirectory(new File(locations[i]));
- }
+ /**
+ * The directories used to store keyspaces data.
+ */
+ public static final DataDirectories dataDirectories = new DataDirectories(DatabaseDescriptor.getNonLocalSystemKeyspacesDataFileLocations(),
+ DatabaseDescriptor.getLocalSystemKeyspacesDataFileLocations());
/**
* Checks whether Cassandra has RWX permissions to the specified directory. Logs an error with
@@ -184,7 +175,7 @@ public class Directories
public Directories(final TableMetadata metadata)
{
- this(metadata, dataDirectories);
+ this(metadata, dataDirectories.getDataDirectoriesFor(metadata));
}
public Directories(final TableMetadata metadata, Collection paths)
@@ -445,10 +436,12 @@ public class Directories
}
if (candidates.isEmpty())
+ {
if (tooBig)
- throw new FSDiskFullWriteError(new IOException("Insufficient disk space to write " + writeSize + " bytes"), "");
- else
- throw new FSWriteError(new IOException("All configured data directories have been disallowed as unwritable for erroring out"), "");
+ throw new FSDiskFullWriteError(metadata.keyspace, writeSize);
+
+ throw new FSNoDiskAvailableForWriteError(metadata.keyspace);
+ }
// shortcut for single data directory systems
if (candidates.size() == 1)
@@ -513,14 +506,10 @@ public class Directories
allowedDirs.add(dir);
}
- Collections.sort(allowedDirs, new Comparator()
- {
- @Override
- public int compare(DataDirectory o1, DataDirectory o2)
- {
- return o1.location.compareTo(o2.location);
- }
- });
+ if (allowedDirs.isEmpty())
+ throw new FSNoDiskAvailableForWriteError(metadata.keyspace);
+
+ allowedDirs.sort(Comparator.comparing(o -> o.location));
return allowedDirs.toArray(new DataDirectory[allowedDirs.size()]);
}
@@ -592,10 +581,35 @@ public class Directories
}
}
+ /**
+ * Checks if the specified table should be stored with local system data.
+ *
+ * To minimize the risk of failures, SSTables for local system keyspaces must be stored in a single data
+ * directory. The only exception to this are some of the system table as the server can continue operating even
+ * if those tables loose some data.
+ *
+ * @param keyspace the keyspace name
+ * @param table the table name
+ * @return {@code true} if the specified table should be stored with local system data, {@code false} otherwise.
+ */
+ public static boolean isStoredInLocalSystemKeyspacesDataLocation(String keyspace, String table)
+ {
+ String keyspaceName = keyspace.toLowerCase();
+
+ return SchemaConstants.LOCAL_SYSTEM_KEYSPACE_NAMES.contains(keyspaceName)
+ && !(SchemaConstants.SYSTEM_KEYSPACE_NAME.equals(keyspaceName)
+ && SystemKeyspace.TABLES_SPLIT_ACROSS_MULTIPLE_DISKS.contains(table.toLowerCase()));
+ }
+
public static class DataDirectory
{
public final File location;
+ public DataDirectory(String location)
+ {
+ this(new File(location));
+ }
+
public DataDirectory(File location)
{
this.location = location;
@@ -632,6 +646,90 @@ public class Directories
}
}
+ /**
+ * Data directories used to store keyspace data.
+ */
+ public static final class DataDirectories implements Iterable
+ {
+ /**
+ * The directories for storing the local system keyspaces.
+ */
+ private final DataDirectory[] localSystemKeyspaceDataDirectories;
+
+ /**
+ * The directories where the data of the non local system keyspaces should be stored.
+ */
+ private final DataDirectory[] nonLocalSystemKeyspacesDirectories;
+
+
+ public DataDirectories(String[] locationsForNonSystemKeyspaces, String[] locationsForSystemKeyspace)
+ {
+ nonLocalSystemKeyspacesDirectories = toDataDirectories(locationsForNonSystemKeyspaces);
+ localSystemKeyspaceDataDirectories = toDataDirectories(locationsForSystemKeyspace);
+ }
+
+ private static DataDirectory[] toDataDirectories(String... locations)
+ {
+ DataDirectory[] directories = new DataDirectory[locations.length];
+ for (int i = 0; i < locations.length; ++i)
+ directories[i] = new DataDirectory(new File(locations[i]));
+ return directories;
+ }
+
+ /**
+ * Returns the data directories for the specified table.
+ *
+ * @param table the table metadata
+ * @return the data directories for the specified table
+ */
+ public DataDirectory[] getDataDirectoriesFor(TableMetadata table)
+ {
+ return isStoredInLocalSystemKeyspacesDataLocation(table.keyspace, table.name) ? localSystemKeyspaceDataDirectories
+ : nonLocalSystemKeyspacesDirectories;
+ }
+
+ @Override
+ public Iterator iterator()
+ {
+ return getAllDirectories().iterator();
+ }
+
+ public Set getAllDirectories()
+ {
+ Set directories = new LinkedHashSet<>(nonLocalSystemKeyspacesDirectories.length + localSystemKeyspaceDataDirectories.length);
+ Collections.addAll(directories, nonLocalSystemKeyspacesDirectories);
+ Collections.addAll(directories, localSystemKeyspaceDataDirectories);
+ return directories;
+ }
+
+ @Override
+ public boolean equals(Object o)
+ {
+ if (this == o) return true;
+ if (o == null || getClass() != o.getClass()) return false;
+
+ DataDirectories that = (DataDirectories) o;
+
+ return Arrays.equals(this.localSystemKeyspaceDataDirectories, that.localSystemKeyspaceDataDirectories)
+ && Arrays.equals(this.nonLocalSystemKeyspacesDirectories, that.nonLocalSystemKeyspacesDirectories);
+ }
+
+ @Override
+ public int hashCode()
+ {
+ return Objects.hash(localSystemKeyspaceDataDirectories, nonLocalSystemKeyspacesDirectories);
+ }
+
+ @Override
+ public String toString()
+ {
+ return "DataDirectories {" +
+ "systemKeyspaceDataDirectories=" + Arrays.toString(localSystemKeyspaceDataDirectories) +
+ ", nonSystemKeyspacesDirectories=" + Arrays.toString(nonLocalSystemKeyspacesDirectories) +
+ '}';
+ }
+ }
+
static final class DataDirectoryCandidate implements Comparable
{
final DataDirectory dataDirectory;
@@ -1001,17 +1099,11 @@ public class Directories
return visitor.getAllocatedSize();
}
+ // Recursively finds all the sub directories in the KS directory.
public static List getKSChildDirectories(String ksName)
- {
- return getKSChildDirectories(ksName, dataDirectories);
-
- }
-
- // Recursively finds all the sub directories in the KS directory.
- public static List getKSChildDirectories(String ksName, DataDirectory[] directories)
{
List result = new ArrayList<>();
- for (DataDirectory dataDirectory : directories)
+ for (DataDirectory dataDirectory : dataDirectories.getAllDirectories())
{
File ksDir = new File(dataDirectory.location, ksName);
File[] cfDirs = ksDir.listFiles();
@@ -1062,21 +1154,6 @@ public class Directories
return StringUtils.join(s, File.separator);
}
- @VisibleForTesting
- static void overrideDataDirectoriesForTest(String loc)
- {
- for (int i = 0; i < dataDirectories.length; ++i)
- dataDirectories[i] = new DataDirectory(new File(loc));
- }
-
- @VisibleForTesting
- static void resetDataDirectoriesAfterTest()
- {
- String[] locations = DatabaseDescriptor.getAllDataFileLocations();
- for (int i = 0; i < locations.length; ++i)
- dataDirectories[i] = new DataDirectory(new File(locations[i]));
- }
-
private class SSTableSizeSummer extends DirectorySizeCalculator
{
private final HashSet toSkip;
diff --git a/src/java/org/apache/cassandra/db/DiskBoundaryManager.java b/src/java/org/apache/cassandra/db/DiskBoundaryManager.java
index bbb6dbb6bf..cc617da702 100644
--- a/src/java/org/apache/cassandra/db/DiskBoundaryManager.java
+++ b/src/java/org/apache/cassandra/db/DiskBoundaryManager.java
@@ -32,7 +32,6 @@ import org.apache.cassandra.dht.Range;
import org.apache.cassandra.dht.Splitter;
import org.apache.cassandra.dht.Token;
import org.apache.cassandra.locator.RangesAtEndpoint;
-import org.apache.cassandra.locator.Replica;
import org.apache.cassandra.locator.TokenMetadata;
import org.apache.cassandra.service.PendingRangeCalculatorService;
import org.apache.cassandra.service.StorageService;
diff --git a/src/java/org/apache/cassandra/db/SystemKeyspace.java b/src/java/org/apache/cassandra/db/SystemKeyspace.java
index bb6ab4abe1..278541d449 100644
--- a/src/java/org/apache/cassandra/db/SystemKeyspace.java
+++ b/src/java/org/apache/cassandra/db/SystemKeyspace.java
@@ -72,8 +72,6 @@ import static java.util.Collections.singletonMap;
import static org.apache.cassandra.cql3.QueryProcessor.executeInternal;
import static org.apache.cassandra.cql3.QueryProcessor.executeOnceInternal;
-import static org.apache.cassandra.locator.Replica.fullReplica;
-import static org.apache.cassandra.locator.Replica.transientReplica;
public final class SystemKeyspace
{
@@ -110,6 +108,17 @@ public final class SystemKeyspace
public static final String PREPARED_STATEMENTS = "prepared_statements";
public static final String REPAIRS = "repairs";
+ /**
+ * By default the system keyspace tables should be stored in a single data directory to allow the server
+ * to handle more gracefully disk failures. Some tables through can be split accross multiple directories
+ * as the server can continue operating even if those tables lost some data.
+ */
+ public static final Set TABLES_SPLIT_ACROSS_MULTIPLE_DISKS = ImmutableSet.of(BATCHES,
+ PAXOS,
+ COMPACTION_HISTORY,
+ PREPARED_STATEMENTS,
+ REPAIRS);
+
@Deprecated public static final String LEGACY_PEERS = "peers";
@Deprecated public static final String LEGACY_PEER_EVENTS = "peer_events";
@Deprecated public static final String LEGACY_TRANSFERRED_RANGES = "transferred_ranges";
diff --git a/src/java/org/apache/cassandra/io/FSDiskFullWriteError.java b/src/java/org/apache/cassandra/io/FSDiskFullWriteError.java
index ca5d8da761..ebb07e2c69 100644
--- a/src/java/org/apache/cassandra/io/FSDiskFullWriteError.java
+++ b/src/java/org/apache/cassandra/io/FSDiskFullWriteError.java
@@ -18,16 +18,20 @@
package org.apache.cassandra.io;
+import java.io.IOException;
+
public class FSDiskFullWriteError extends FSWriteError
{
- public FSDiskFullWriteError(Throwable cause, String path)
+ public FSDiskFullWriteError(String keyspace, long mutationSize)
{
- super(cause, path);
+ super(new IOException(String.format("Insufficient disk space to write %d bytes into the %s keyspace",
+ mutationSize,
+ keyspace)));
}
@Override
public String toString()
{
- return "FSDiskFullWriteError in " + path;
+ return "FSDiskFullWriteError";
}
}
diff --git a/src/java/org/apache/cassandra/io/FSNoDiskAvailableForWriteError.java b/src/java/org/apache/cassandra/io/FSNoDiskAvailableForWriteError.java
new file mode 100644
index 0000000000..14dcd38f2a
--- /dev/null
+++ b/src/java/org/apache/cassandra/io/FSNoDiskAvailableForWriteError.java
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.cassandra.io;
+
+import java.io.IOException;
+
+/**
+ * Thrown when all the disks used by a given keyspace have been marked as unwriteable.
+ */
+public class FSNoDiskAvailableForWriteError extends FSWriteError
+{
+ public FSNoDiskAvailableForWriteError(String keyspace)
+ {
+ super(new IOException(String.format("The data directories for the %s keyspace have been marked as unwritable",
+ keyspace)));
+ }
+
+ @Override
+ public String toString()
+ {
+ return "FSNoDiskAvailableForWriteError";
+ }
+}
diff --git a/src/java/org/apache/cassandra/io/FSWriteError.java b/src/java/org/apache/cassandra/io/FSWriteError.java
index 6169904648..b419086be0 100644
--- a/src/java/org/apache/cassandra/io/FSWriteError.java
+++ b/src/java/org/apache/cassandra/io/FSWriteError.java
@@ -31,6 +31,11 @@ public class FSWriteError extends FSError
this(cause, new File(path));
}
+ public FSWriteError(Throwable cause)
+ {
+ this(cause, new File(""));
+ }
+
@Override
public String toString()
{
diff --git a/src/java/org/apache/cassandra/io/util/FileUtils.java b/src/java/org/apache/cassandra/io/util/FileUtils.java
index f5061404b0..e0ea436e5e 100644
--- a/src/java/org/apache/cassandra/io/util/FileUtils.java
+++ b/src/java/org/apache/cassandra/io/util/FileUtils.java
@@ -41,8 +41,12 @@ import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
import com.google.common.util.concurrent.RateLimiter;
+import com.google.common.base.Preconditions;
+
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -969,4 +973,68 @@ public final class FileUtils
return fileStore.getAttribute(attribute);
}
}
+
+ /**
+ * Moves the contents of a directory to another directory.
+ * Once a file has been copied to the target directory it will be deleted from the source directory.
+ * If a file already exists in the target directory a warning will be logged and the file will not
+ * be deleted.
+ *
+ * @param source the directory containing the files to move
+ * @param target the directory where the files must be moved
+ */
+ public static void moveRecursively(Path source, Path target) throws IOException
+ {
+ logger.info("Moving {} to {}" , source, target);
+
+ if (Files.isDirectory(source))
+ {
+ Files.createDirectories(target);
+
+ for (File f : source.toFile().listFiles())
+ {
+ String fileName = f.getName();
+ moveRecursively(source.resolve(fileName), target.resolve(fileName));
+ }
+
+ deleteDirectoryIfEmpty(source);
+ }
+ else
+ {
+ if (Files.exists(target))
+ {
+ logger.warn("Cannot move the file {} to {} as the target file already exists." , source, target);
+ }
+ else
+ {
+ Files.copy(source, target, StandardCopyOption.COPY_ATTRIBUTES);
+ Files.delete(source);
+ }
+ }
+ }
+
+ /**
+ * Deletes the specified directory if it is empty
+ *
+ * @param path the path to the directory
+ */
+ public static void deleteDirectoryIfEmpty(Path path) throws IOException
+ {
+ Preconditions.checkArgument(Files.isDirectory(path), String.format("%s is not a directory", path));
+
+ try
+ {
+ logger.info("Deleting directory {}", path);
+ Files.delete(path);
+ }
+ catch (DirectoryNotEmptyException e)
+ {
+ try (Stream paths = Files.list(path))
+ {
+ String content = paths.map(p -> p.getFileName().toString()).collect(Collectors.joining(", "));
+
+ logger.warn("Cannot delete the directory {} as it is not empty. (Content: {})", path, content);
+ }
+ }
+ }
}
diff --git a/src/java/org/apache/cassandra/service/CassandraDaemon.java b/src/java/org/apache/cassandra/service/CassandraDaemon.java
index 6b22635c2b..2cb12540c6 100644
--- a/src/java/org/apache/cassandra/service/CassandraDaemon.java
+++ b/src/java/org/apache/cassandra/service/CassandraDaemon.java
@@ -24,8 +24,14 @@ import java.lang.management.MemoryPoolMXBean;
import java.net.InetAddress;
import java.net.URL;
import java.net.UnknownHostException;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.Paths;
+import java.util.Arrays;
import java.util.List;
import java.util.concurrent.TimeUnit;
+import java.util.stream.Stream;
+
import javax.management.ObjectName;
import javax.management.StandardMBean;
import javax.management.remote.JMXConnectorServer;
@@ -44,6 +50,7 @@ import com.codahale.metrics.jvm.BufferPoolMetricSet;
import com.codahale.metrics.jvm.FileDescriptorRatioGauge;
import com.codahale.metrics.jvm.GarbageCollectorMetricSet;
import com.codahale.metrics.jvm.MemoryUsageGaugeSet;
+
import org.apache.cassandra.audit.AuditLogManager;
import org.apache.cassandra.concurrent.ScheduledExecutors;
import org.apache.cassandra.config.DatabaseDescriptor;
@@ -223,6 +230,19 @@ public class CassandraDaemon
{
FileUtils.setFSErrorHandler(new DefaultFSErrorHandler());
+ // Since CASSANDRA-14793 the local system keyspaces data are not dispatched across the data directories
+ // anymore to reduce the risks in case of disk failures. By consequence, the system need to ensure in case of
+ // upgrade that the old data files have been migrated to the new directories before we start deleting
+ // snapshots and upgrading system tables.
+ try
+ {
+ migrateSystemDataIfNeeded();
+ }
+ catch (IOException e)
+ {
+ exitOrFail(StartupException.ERR_WRONG_DISK_STATE, e.getMessage(), e);
+ }
+
// Delete any failed snapshot deletions on Windows - see CASSANDRA-9658
if (FBUtilities.isWindows)
WindowsFailedSnapshotTracker.deleteOldSnapshots();
@@ -247,7 +267,7 @@ public class CassandraDaemon
}
catch (IOException e)
{
- exitOrFail(3, e.getMessage(), e.getCause());
+ exitOrFail(StartupException.ERR_WRONG_DISK_STATE, e.getMessage(), e.getCause());
}
// We need to persist this as soon as possible after startup checks.
@@ -472,6 +492,73 @@ public class CassandraDaemon
}
}
+
+ /**
+ * Checks if the data of the local system keyspaces need to be migrated to a different location.
+ *
+ * @throws IOException
+ */
+ public void migrateSystemDataIfNeeded() throws IOException
+ {
+ // If there is only one directory and no system keyspace directory has been specified we do not need to do
+ // anything. If it is not the case we want to try to migrate the data.
+ if (!DatabaseDescriptor.useSpecificLocationForLocalSystemData()
+ && DatabaseDescriptor.getNonLocalSystemKeyspacesDataFileLocations().length <= 1)
+ return;
+
+ // We can face several cases:
+ // 1) The system data are spread accross the data file locations and need to be moved to
+ // the first data location (upgrade to 4.0)
+ // 2) The system data are spread accross the data file locations and need to be moved to
+ // the system keyspace location configured by the user (upgrade to 4.0)
+ // 3) The system data are stored in the first data location and need to be moved to
+ // the system keyspace location configured by the user (system_data_file_directory has been configured)
+ Path target = Paths.get(DatabaseDescriptor.getLocalSystemKeyspacesDataFileLocations()[0]);
+
+ String[] nonLocalSystemKeyspacesFileLocations = DatabaseDescriptor.getNonLocalSystemKeyspacesDataFileLocations();
+ String[] sources = DatabaseDescriptor.useSpecificLocationForLocalSystemData() ? nonLocalSystemKeyspacesFileLocations
+ : Arrays.copyOfRange(nonLocalSystemKeyspacesFileLocations,
+ 1,
+ nonLocalSystemKeyspacesFileLocations.length);
+
+ for (String source : sources)
+ {
+ Path dataFileLocation = Paths.get(source);
+
+ if (!Files.exists(dataFileLocation))
+ continue;
+
+ try (Stream locationChildren = Files.list(dataFileLocation))
+ {
+ Path[] keyspaceDirectories = locationChildren.filter(p -> SchemaConstants.isLocalSystemKeyspace(p.getFileName().toString()))
+ .toArray(Path[]::new);
+
+ for (Path keyspaceDirectory : keyspaceDirectories)
+ {
+ try (Stream keyspaceChildren = Files.list(keyspaceDirectory))
+ {
+ Path[] tableDirectories = keyspaceChildren.filter(Files::isDirectory)
+ .filter(p -> !SystemKeyspace.TABLES_SPLIT_ACROSS_MULTIPLE_DISKS
+ .contains(p.getFileName()
+ .toString()))
+ .toArray(Path[]::new);
+
+ for (Path tableDirectory : tableDirectories)
+ {
+ FileUtils.moveRecursively(tableDirectory,
+ target.resolve(dataFileLocation.relativize(tableDirectory)));
+ }
+
+ if (!SchemaConstants.SYSTEM_KEYSPACE_NAME.equals(keyspaceDirectory.getFileName().toString()))
+ {
+ FileUtils.deleteDirectoryIfEmpty(keyspaceDirectory);
+ }
+ }
+ }
+ }
+ }
+ }
+
public void setupVirtualKeyspaces()
{
VirtualKeyspaceRegistry.instance.register(VirtualSchemaKeyspace.instance);
diff --git a/src/java/org/apache/cassandra/service/DefaultFSErrorHandler.java b/src/java/org/apache/cassandra/service/DefaultFSErrorHandler.java
index d72b59a9b5..d5e3e531c1 100644
--- a/src/java/org/apache/cassandra/service/DefaultFSErrorHandler.java
+++ b/src/java/org/apache/cassandra/service/DefaultFSErrorHandler.java
@@ -26,9 +26,7 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.DisallowedDirectories;
import org.apache.cassandra.db.Keyspace;
-import org.apache.cassandra.io.FSError;
-import org.apache.cassandra.io.FSErrorHandler;
-import org.apache.cassandra.io.FSReadError;
+import org.apache.cassandra.io.*;
import org.apache.cassandra.io.sstable.CorruptSSTableException;
import org.apache.cassandra.utils.JVMStabilityInspector;
@@ -67,6 +65,18 @@ public class DefaultFSErrorHandler implements FSErrorHandler
StorageService.instance.stopTransports();
break;
case best_effort:
+
+ // There are a few scenarios where we know that the node will not be able to operate properly.
+ // For those scenarios we want to stop the transports and let the administrators handle the problem.
+ // Those scenarios are:
+ // * All the disks are full
+ // * All the disks for a given keyspace have been marked as unwriteable
+ if (e instanceof FSDiskFullWriteError || e instanceof FSNoDiskAvailableForWriteError)
+ {
+ logger.error("Stopping transports: " + e.getCause().getMessage());
+ StorageService.instance.stopTransports();
+ }
+
// for both read and write errors mark the path as unwritable.
DisallowedDirectories.maybeMarkUnwritable(e.path);
if (e instanceof FSReadError)
diff --git a/src/java/org/apache/cassandra/service/StartupChecks.java b/src/java/org/apache/cassandra/service/StartupChecks.java
index cf3d414d38..85b5836baf 100644
--- a/src/java/org/apache/cassandra/service/StartupChecks.java
+++ b/src/java/org/apache/cassandra/service/StartupChecks.java
@@ -340,6 +340,7 @@ public class StartupChecks
Arrays.asList(DatabaseDescriptor.getCommitLogLocation(),
DatabaseDescriptor.getSavedCachesLocation(),
DatabaseDescriptor.getHintsDirectory().getAbsolutePath()));
+
for (String dataDir : dirs)
{
logger.debug("Checking directory {}", dataDir);
diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java
index da0c8ea1b6..661b1a0b94 100644
--- a/src/java/org/apache/cassandra/service/StorageService.java
+++ b/src/java/org/apache/cassandra/service/StorageService.java
@@ -3460,14 +3460,32 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
return stringify(Gossiper.instance.getUnreachableMembers(), true);
}
+ @Override
public String[] getAllDataFileLocations()
{
- String[] locations = DatabaseDescriptor.getAllDataFileLocations();
- for (int i = 0; i < locations.length; i++)
- locations[i] = FileUtils.getCanonicalPath(locations[i]);
+ return getCanonicalPaths(DatabaseDescriptor.getAllDataFileLocations());
+ }
+
+ private String[] getCanonicalPaths(String[] paths)
+ {
+ String[] locations = new String[paths.length];
+ for (int i = 0; i < paths.length; i++)
+ locations[i] = FileUtils.getCanonicalPath(paths[i]);
return locations;
}
+ @Override
+ public String[] getLocalSystemKeyspacesDataFileLocations()
+ {
+ return getCanonicalPaths(DatabaseDescriptor.getLocalSystemKeyspacesDataFileLocations());
+ }
+
+ @Override
+ public String[] getNonLocalSystemKeyspacesDataFileLocations()
+ {
+ return getCanonicalPaths(DatabaseDescriptor.getNonLocalSystemKeyspacesDataFileLocations());
+ }
+
public String getCommitLogLocation()
{
return FileUtils.getCanonicalPath(DatabaseDescriptor.getCommitLogLocation());
diff --git a/src/java/org/apache/cassandra/service/StorageServiceMBean.java b/src/java/org/apache/cassandra/service/StorageServiceMBean.java
index 63c2d96037..cc69fec613 100644
--- a/src/java/org/apache/cassandra/service/StorageServiceMBean.java
+++ b/src/java/org/apache/cassandra/service/StorageServiceMBean.java
@@ -120,6 +120,20 @@ public interface StorageServiceMBean extends NotificationEmitter
*/
public String[] getAllDataFileLocations();
+ /**
+ * Returns the locations where the local system keyspaces data should be stored.
+ *
+ * @return the locations where the local system keyspaces data should be stored
+ */
+ public String[] getLocalSystemKeyspacesDataFileLocations();
+
+ /**
+ * Returns the locations where should be stored the non system keyspaces data.
+ *
+ * @return the locations where should be stored the non system keyspaces data
+ */
+ public String[] getNonLocalSystemKeyspacesDataFileLocations();
+
/**
* Get location of the commit log
* @return a string path
diff --git a/test/conf/system_keyspaces_directory.yaml b/test/conf/system_keyspaces_directory.yaml
new file mode 100644
index 0000000000..685d54ed0e
--- /dev/null
+++ b/test/conf/system_keyspaces_directory.yaml
@@ -0,0 +1 @@
+local_system_data_file_directory: build/test/cassandra/system_data
diff --git a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
index b0a6c3f279..0ed70003b4 100644
--- a/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
+++ b/test/distributed/org/apache/cassandra/distributed/impl/Instance.java
@@ -426,6 +426,7 @@ public class Instance extends IsolatedExecutor implements IInvokableInstance
DatabaseDescriptor.daemonInitialization();
FileUtils.setFSErrorHandler(new DefaultFSErrorHandler());
DatabaseDescriptor.createAllDirectories();
+ CassandraDaemon.getInstanceForTesting().migrateSystemDataIfNeeded();
CommitLog.instance.start();
CassandraDaemon.getInstanceForTesting().runStartupChecks();
diff --git a/test/unit/org/apache/cassandra/OffsetAwareConfigurationLoader.java b/test/unit/org/apache/cassandra/OffsetAwareConfigurationLoader.java
index 23138b096c..27246aa069 100644
--- a/test/unit/org/apache/cassandra/OffsetAwareConfigurationLoader.java
+++ b/test/unit/org/apache/cassandra/OffsetAwareConfigurationLoader.java
@@ -94,6 +94,9 @@ public class OffsetAwareConfigurationLoader extends YamlConfigurationLoader
for (int i = 0; i < config.data_file_directories.length; i++)
config.data_file_directories[i] += sep + offset;
+ if (config.local_system_data_file_directory != null)
+ config.local_system_data_file_directory += sep + offset;
+
return config;
}
}
diff --git a/test/unit/org/apache/cassandra/db/DirectoriesTest.java b/test/unit/org/apache/cassandra/db/DirectoriesTest.java
index eb2016f80e..507827e1e1 100644
--- a/test/unit/org/apache/cassandra/db/DirectoriesTest.java
+++ b/test/unit/org/apache/cassandra/db/DirectoriesTest.java
@@ -22,7 +22,6 @@ import java.io.FileOutputStream;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
-import java.nio.file.Paths;
import java.util.*;
import java.util.concurrent.Callable;
import java.util.concurrent.Executors;
@@ -36,10 +35,14 @@ import org.junit.Test;
import org.apache.cassandra.cql3.ColumnIdentifier;
import org.apache.cassandra.schema.Indexes;
+import org.apache.cassandra.schema.SchemaConstants;
+import org.apache.cassandra.schema.SchemaKeyspace;
import org.apache.cassandra.schema.TableMetadata;
+import org.apache.cassandra.auth.AuthKeyspace;
import org.apache.cassandra.config.Config.DiskFailurePolicy;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.cql3.statements.schema.IndexTarget;
+import org.apache.cassandra.db.Directories.DataDirectories;
import org.apache.cassandra.db.Directories.DataDirectory;
import org.apache.cassandra.db.marshal.UTF8Type;
import org.apache.cassandra.index.internal.CassandraIndex;
@@ -50,7 +53,6 @@ import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.schema.IndexMetadata;
import org.apache.cassandra.service.DefaultFSErrorHandler;
-import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.JVMStabilityInspector;
import static org.junit.Assert.assertEquals;
@@ -88,7 +90,6 @@ public class DirectoriesTest
tempDataDir.delete(); // hack to create a temp dir
tempDataDir.mkdir();
- Directories.overrideDataDirectoriesForTest(tempDataDir.getPath());
// Create two fake data dir for tests, one using CF directories, one that do not.
createTestFiles();
}
@@ -96,10 +97,14 @@ public class DirectoriesTest
@AfterClass
public static void afterClass()
{
- Directories.resetDataDirectoriesAfterTest();
FileUtils.deleteRecursive(tempDataDir);
}
+ private static DataDirectory[] toDataDirectories(File location) throws IOException
+ {
+ return new DataDirectory[] { new DataDirectory(location) };
+ }
+
private static void createTestFiles() throws IOException
{
for (TableMetadata cfm : CFM)
@@ -156,7 +161,7 @@ public class DirectoriesTest
{
for (TableMetadata cfm : CFM)
{
- Directories directories = new Directories(cfm);
+ Directories directories = new Directories(cfm, toDataDirectories(tempDataDir));
assertEquals(cfDir(cfm), directories.getDirectoryForNewSSTables());
Descriptor desc = new Descriptor(cfDir(cfm), KS, cfm.name, 1, SSTableFormat.Type.BIG);
@@ -169,7 +174,7 @@ public class DirectoriesTest
}
@Test
- public void testSecondaryIndexDirectories()
+ public void testSecondaryIndexDirectories() throws IOException
{
TableMetadata.Builder builder =
TableMetadata.builder(KS, "cf")
@@ -187,8 +192,8 @@ public class DirectoriesTest
TableMetadata PARENT_CFM = builder.build();
TableMetadata INDEX_CFM = CassandraIndex.indexCfsMetadata(PARENT_CFM, indexDef);
- Directories parentDirectories = new Directories(PARENT_CFM);
- Directories indexDirectories = new Directories(INDEX_CFM);
+ Directories parentDirectories = new Directories(PARENT_CFM, toDataDirectories(tempDataDir));
+ Directories indexDirectories = new Directories(INDEX_CFM, toDataDirectories(tempDataDir));
// secondary index has its own directory
for (File dir : indexDirectories.getCFDirectories())
{
@@ -248,11 +253,11 @@ public class DirectoriesTest
}
@Test
- public void testSSTableLister()
+ public void testSSTableLister() throws IOException
{
for (TableMetadata cfm : CFM)
{
- Directories directories = new Directories(cfm);
+ Directories directories = new Directories(cfm, toDataDirectories(tempDataDir));
checkFiles(cfm, directories);
}
}
@@ -301,7 +306,7 @@ public class DirectoriesTest
{
for (TableMetadata cfm : CFM)
{
- Directories directories = new Directories(cfm);
+ Directories directories = new Directories(cfm, toDataDirectories(tempDataDir));
File tempDir = directories.getTemporaryWriteableDirectoryAsFile(10);
tempDir.mkdir();
@@ -332,19 +337,21 @@ public class DirectoriesTest
try
{
DatabaseDescriptor.setDiskFailurePolicy(DiskFailurePolicy.best_effort);
+
+ Set directories = Directories.dataDirectories.getAllDirectories();
+ DataDirectory first = directories.iterator().next();
+
// Fake a Directory creation failure
- if (Directories.dataDirectories.length > 0)
+ if (!directories.isEmpty())
{
String[] path = new String[] {KS, "bad"};
- File dir = new File(Directories.dataDirectories[0].location, StringUtils.join(path, File.separator));
+ File dir = new File(first.location, StringUtils.join(path, File.separator));
JVMStabilityInspector.inspectThrowable(new FSWriteError(new IOException("Unable to create directory " + dir), dir));
}
- for (DataDirectory dd : Directories.dataDirectories)
- {
- File file = new File(dd.location, new File(KS, "bad").getPath());
- assertTrue(DisallowedDirectories.isUnwritable(file));
- }
+ File file = new File(first.location, new File(KS, "bad").getPath());
+ assertTrue(DisallowedDirectories.isUnwritable(file));
+
}
finally
{
@@ -357,7 +364,7 @@ public class DirectoriesTest
{
for (final TableMetadata cfm : CFM)
{
- final Directories directories = new Directories(cfm);
+ final Directories directories = new Directories(cfm, toDataDirectories(tempDataDir));
assertEquals(cfDir(cfm), directories.getDirectoryForNewSSTables());
final String n = Long.toString(System.nanoTime());
Callable directoryGetter = new Callable() {
@@ -519,9 +526,9 @@ public class DirectoriesTest
public void getDataDirectoryForFile()
{
Collection paths = new ArrayList<>();
- paths.add(new DataDirectory(new File("/tmp/a")));
- paths.add(new DataDirectory(new File("/tmp/aa")));
- paths.add(new DataDirectory(new File("/tmp/aaa")));
+ paths.add(new DataDirectory("/tmp/a"));
+ paths.add(new DataDirectory("/tmp/aa"));
+ paths.add(new DataDirectory("/tmp/aaa"));
for (TableMetadata cfm : CFM)
{
@@ -614,6 +621,52 @@ public class DirectoriesTest
}
}
+ @Test
+ public void testIsStoredInLocalSystemKeyspacesDataLocation() throws IOException
+ {
+ for (String table : SystemKeyspace.TABLES_SPLIT_ACROSS_MULTIPLE_DISKS)
+ {
+ assertFalse(Directories.isStoredInLocalSystemKeyspacesDataLocation(SchemaConstants.SYSTEM_KEYSPACE_NAME, table));
+ }
+ assertTrue(Directories.isStoredInLocalSystemKeyspacesDataLocation(SchemaConstants.SYSTEM_KEYSPACE_NAME, SystemKeyspace.PEERS_V2));
+ assertTrue(Directories.isStoredInLocalSystemKeyspacesDataLocation(SchemaConstants.SYSTEM_KEYSPACE_NAME, SystemKeyspace.TRANSFERRED_RANGES_V2));
+ assertTrue(Directories.isStoredInLocalSystemKeyspacesDataLocation(SchemaConstants.SCHEMA_KEYSPACE_NAME, SchemaKeyspace.KEYSPACES));
+ assertTrue(Directories.isStoredInLocalSystemKeyspacesDataLocation(SchemaConstants.SCHEMA_KEYSPACE_NAME, SchemaKeyspace.TABLES));
+ assertFalse(Directories.isStoredInLocalSystemKeyspacesDataLocation(SchemaConstants.AUTH_KEYSPACE_NAME, AuthKeyspace.ROLES));
+ assertFalse(Directories.isStoredInLocalSystemKeyspacesDataLocation(KS, TABLES[0]));
+ }
+
+ @Test
+ public void testDataDirectoriesIterator() throws IOException
+ {
+ Path tmpDir = Files.createTempDirectory(this.getClass().getSimpleName());
+ Path subDir_1 = Files.createDirectory(tmpDir.resolve("a"));
+ Path subDir_2 = Files.createDirectory(tmpDir.resolve("b"));
+ Path subDir_3 = Files.createDirectory(tmpDir.resolve("c"));
+
+ DataDirectories directories = new DataDirectories(new String[]{subDir_1.toString(), subDir_2.toString()},
+ new String[]{subDir_3.toString()});
+
+ Iterator iter = directories.iterator();
+ assertTrue(iter.hasNext());
+ assertEquals(new DataDirectory(subDir_1.toFile()), iter.next());
+ assertTrue(iter.hasNext());
+ assertEquals(new DataDirectory(subDir_2.toFile()), iter.next());
+ assertTrue(iter.hasNext());
+ assertEquals(new DataDirectory(subDir_3.toFile()), iter.next());
+ assertFalse(iter.hasNext());
+
+ directories = new DataDirectories(new String[]{subDir_1.toString(), subDir_2.toString()},
+ new String[]{subDir_1.toString()});
+
+ iter = directories.iterator();
+ assertTrue(iter.hasNext());
+ assertEquals(new DataDirectory(subDir_1.toFile()), iter.next());
+ assertTrue(iter.hasNext());
+ assertEquals(new DataDirectory(subDir_2.toFile()), iter.next());
+ assertFalse(iter.hasNext());
+ }
+
private String getNewFilename(TableMetadata tm, boolean oldStyle)
{
return tm.keyspace + File.separator + tm.name + (oldStyle ? "" : Component.separator + tm.id.toHexString()) + "/na-1-big-Data.db";
diff --git a/test/unit/org/apache/cassandra/io/util/FileUtilsTest.java b/test/unit/org/apache/cassandra/io/util/FileUtilsTest.java
index 373232df02..7d19f516ff 100644
--- a/test/unit/org/apache/cassandra/io/util/FileUtilsTest.java
+++ b/test/unit/org/apache/cassandra/io/util/FileUtilsTest.java
@@ -31,7 +31,9 @@ import org.junit.BeforeClass;
import org.junit.Test;
import org.apache.cassandra.config.DatabaseDescriptor;
+import org.assertj.core.api.Assertions;
+import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
@@ -128,6 +130,96 @@ public class FileUtilsTest
assertFalse(FileUtils.isContained(new File("/tmp/abc/../abc"), new File("/tmp/abcc")));
}
+ @Test
+ public void testMoveFiles() throws IOException
+ {
+ Path tmpDir = Files.createTempDirectory(this.getClass().getSimpleName());
+ Path sourceDir = Files.createDirectory(tmpDir.resolve("source"));
+ Path subDir_1 = Files.createDirectory(sourceDir.resolve("a"));
+ subDir_1.resolve("file_1.txt").toFile().createNewFile();
+ subDir_1.resolve("file_2.txt").toFile().createNewFile();
+ Path subDir_11 = Files.createDirectory(subDir_1.resolve("ab"));
+ subDir_11.resolve("file_1.txt").toFile().createNewFile();
+ subDir_11.resolve("file_2.txt").toFile().createNewFile();
+ subDir_11.resolve("file_3.txt").toFile().createNewFile();
+ Path subDir_12 = Files.createDirectory(subDir_1.resolve("ac"));
+ Path subDir_2 = Files.createDirectory(sourceDir.resolve("b"));
+ subDir_2.resolve("file_1.txt").toFile().createNewFile();
+ subDir_2.resolve("file_2.txt").toFile().createNewFile();
+
+ Path targetDir = Files.createDirectory(tmpDir.resolve("target"));
+
+ FileUtils.moveRecursively(sourceDir, targetDir);
+
+ assertThat(sourceDir).doesNotExist();
+ assertThat(targetDir.resolve("a/file_1.txt")).exists();
+ assertThat(targetDir.resolve("a/file_2.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_1.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_2.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_3.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_1.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_2.txt")).exists();
+ assertThat(targetDir.resolve("a/ac/")).exists();
+ assertThat(targetDir.resolve("b/file_1.txt")).exists();
+ assertThat(targetDir.resolve("b/file_2.txt")).exists();
+
+ // Tests that files can be moved into existing directories
+
+ sourceDir = Files.createDirectory(tmpDir.resolve("source2"));
+ subDir_1 = Files.createDirectory(sourceDir.resolve("a"));
+ subDir_1.resolve("file_3.txt").toFile().createNewFile();
+ subDir_11 = Files.createDirectory(subDir_1.resolve("ab"));
+ subDir_11.resolve("file_4.txt").toFile().createNewFile();
+
+ FileUtils.moveRecursively(sourceDir, targetDir);
+
+ assertThat(sourceDir).doesNotExist();
+ assertThat(targetDir.resolve("a/file_1.txt")).exists();
+ assertThat(targetDir.resolve("a/file_2.txt")).exists();
+ assertThat(targetDir.resolve("a/file_3.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_1.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_2.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_3.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_4.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_1.txt")).exists();
+ assertThat(targetDir.resolve("a/ab/file_2.txt")).exists();
+ assertThat(targetDir.resolve("a/ac/")).exists();
+ assertThat(targetDir.resolve("b/file_1.txt")).exists();
+ assertThat(targetDir.resolve("b/file_2.txt")).exists();
+
+ // Tests that existing files are not replaced but trigger an error.
+
+ sourceDir = Files.createDirectory(tmpDir.resolve("source3"));
+ subDir_1 = Files.createDirectory(sourceDir.resolve("a"));
+ subDir_1.resolve("file_3.txt").toFile().createNewFile();
+ FileUtils.moveRecursively(sourceDir, targetDir);
+
+ assertThat(sourceDir).exists();
+ assertThat(sourceDir.resolve("a/file_3.txt")).exists();
+ assertThat(targetDir.resolve("a/file_3.txt")).exists();
+ }
+
+ @Test
+ public void testDeleteDirectoryIfEmpty() throws IOException
+ {
+ Path tmpDir = Files.createTempDirectory(this.getClass().getSimpleName());
+ Path subDir_1 = Files.createDirectory(tmpDir.resolve("a"));
+ Path subDir_2 = Files.createDirectory(tmpDir.resolve("b"));
+ Path file_1 = subDir_2.resolve("file_1.txt");
+ file_1.toFile().createNewFile();
+
+ FileUtils.deleteDirectoryIfEmpty(subDir_1);
+ assertThat(subDir_1).doesNotExist();
+
+ FileUtils.deleteDirectoryIfEmpty(subDir_2);
+ assertThat(subDir_2).exists();
+
+ Assertions.assertThatThrownBy(() -> FileUtils.deleteDirectoryIfEmpty(file_1))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("is not a directory");
+ }
+
+
private File createFolder(Path path)
{
File folder = path.toFile();
diff --git a/test/unit/org/apache/cassandra/tools/ClearSnapshotTest.java b/test/unit/org/apache/cassandra/tools/ClearSnapshotTest.java
index b63182243a..975e45bafb 100644
--- a/test/unit/org/apache/cassandra/tools/ClearSnapshotTest.java
+++ b/test/unit/org/apache/cassandra/tools/ClearSnapshotTest.java
@@ -100,7 +100,7 @@ public class ClearSnapshotTest extends CQLTester
tool = ToolRunner.invokeNodetool("snapshot","-t","some-other-name");
tool.assertOnCleanExit();
assertTrue(!tool.getStdout().isEmpty());
-
+
Map snapshots_before = probe.getSnapshotDetails();
Assert.assertTrue(snapshots_before.size() == 2);
diff --git a/test/unit/org/apache/cassandra/tools/NodeToolTPStatsTest.java b/test/unit/org/apache/cassandra/tools/NodeToolTPStatsTest.java
index e437fc11a5..31423a407d 100644
--- a/test/unit/org/apache/cassandra/tools/NodeToolTPStatsTest.java
+++ b/test/unit/org/apache/cassandra/tools/NodeToolTPStatsTest.java
@@ -112,7 +112,7 @@ public class NodeToolTPStatsTest extends CQLTester
public void testTPStats() throws Throwable
{
ToolResult tool = ToolRunner.invokeNodetool("tpstats");
- Assertions.assertThat(tool.getStdout()).containsIgnoringCase("Pool Name Active Pending Completed Blocked All time blocked");
+ Assertions.assertThat(tool.getStdout()).containsPattern("Pool Name \\s* Active Pending Completed Blocked All time blocked");
Assertions.assertThat(tool.getStdout()).containsIgnoringCase("Latencies waiting in queue (micros) per dropped message types");
assertTrue(tool.getCleanedStderr().isEmpty());
assertEquals(0, tool.getExitCode());