diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index f4a6328d0a..ec5a1cbbfa 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -41,6 +41,7 @@ import java.net.InetAddress; public class DatabaseDescriptor { private static Logger logger_ = Logger.getLogger(DatabaseDescriptor.class); + public static final String STREAMING_SUBDIR = "stream"; // don't capitalize these; we need them to match what's in the config file for CLS.valueOf to parse public static enum CommitLogSync { @@ -599,7 +600,11 @@ public class DatabaseDescriptor FileUtils.createDirectory(dataFile + File.separator + Table.SYSTEM_TABLE); for (String table : tables_) { - FileUtils.createDirectory(dataFile + File.separator + table); + String oneDir = dataFile + File.separator + table; + FileUtils.createDirectory(oneDir); + File streamingDir = new File(oneDir, STREAMING_SUBDIR); + if (streamingDir.exists()) + FileUtils.deleteDir(streamingDir); } } } diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 83390893d9..36eb4bdec7 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -822,7 +822,7 @@ public final class ColumnFamilyStore implements ColumnFamilyStoreMBean { if (target != null) { - compactionFileLocation = compactionFileLocation + File.separator + "stream"; + compactionFileLocation = compactionFileLocation + File.separator + DatabaseDescriptor.STREAMING_SUBDIR; } FileUtils.createDirectory(compactionFileLocation); String newFilename = new File(compactionFileLocation, getTempSSTableFileName()).getAbsolutePath(); diff --git a/src/java/org/apache/cassandra/io/Streaming.java b/src/java/org/apache/cassandra/io/Streaming.java index cb22649102..ffd6899e35 100644 --- a/src/java/org/apache/cassandra/io/Streaming.java +++ b/src/java/org/apache/cassandra/io/Streaming.java @@ -94,9 +94,12 @@ public class Streaming if (logger.isDebugEnabled()) logger.debug("Waiting for transfer to " + target + " to complete"); StreamManager.instance(target).waitForStreamCompletion(); - // reference sstables one more time to make sure it doesn't get GC'd early (causing delete of its files) + for (SSTableReader sstable : sstables) + { + sstable.markCompacted(); + } if (logger.isDebugEnabled()) - logger.debug("Done with transfer to " + target + " of " + StringUtils.join(sstables, ", ")); + logger.debug("Done with transfer to " + target); } }