diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 2787b6316a..920edbc301 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -437,7 +437,8 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean void forceBlockingFlush() throws IOException, ExecutionException, InterruptedException { - forceFlush(); + Memtable oldMemtable = memtable_.get(); + oldMemtable.forceflush(); // block for flush to finish by adding a no-op action to the flush executorservice // and waiting for that to finish. (this works since flush ES is single-threaded.) Future f = MemtableManager.instance().flusher_.submit(new Runnable() @@ -447,6 +448,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean } }); f.get(); + assert oldMemtable.isFlushed() || oldMemtable.isClean(); } void forceFlushBinary() diff --git a/src/java/org/apache/cassandra/db/Memtable.java b/src/java/org/apache/cassandra/db/Memtable.java index 9eeec349fa..b710723a4e 100644 --- a/src/java/org/apache/cassandra/db/Memtable.java +++ b/src/java/org/apache/cassandra/db/Memtable.java @@ -56,7 +56,8 @@ public class Memtable implements Comparable } private MemtableThreadPoolExecutor executor_; - private boolean isFrozen_; + private volatile boolean isFrozen_; + private volatile boolean isFlushed_; // for tests, in particular forceBlockingFlush asserts this private int threshold_ = DatabaseDescriptor.getMemtableSize()*1024*1024; private int thresholdCount_ = (int)(DatabaseDescriptor.getMemtableObjectCount()*1024*1024); @@ -81,6 +82,11 @@ public class Memtable implements Comparable runningExecutorServices_.add(executor_); } + public boolean isFlushed() + { + return isFlushed_; + } + class Putter implements Runnable { private String key_; @@ -203,7 +209,7 @@ public class Memtable implements Comparable */ public void forceflush() { - if (columnFamilies_.isEmpty()) + if (isClean()) return; try @@ -355,6 +361,7 @@ public class Memtable implements Comparable cfStore.onMemtableFlush(cLogCtx); cfStore.storeLocation( ssTable.getDataFileLocation(), bf ); buffer.close(); + isFlushed_ = true; } private class MemtableThreadPoolExecutor extends DebuggableThreadPoolExecutor @@ -401,4 +408,9 @@ public class Memtable implements Comparable pq.addAll(keys); return new DestructivePQIterator(pq); } + + public boolean isClean() + { + return columnFamilies_.isEmpty() && executor_.getPendingTasks() == 0; + } } diff --git a/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java b/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java index b12891a0f7..ed5efb5b3f 100644 --- a/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java +++ b/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java @@ -257,6 +257,7 @@ public class ColumnFamilyStoreTest extends ServerTest rm.apply(); List families = store.getColumnFamilies("key1", "Super1", new IdentityFilter()); + assert families.size() == 2 : StringUtils.join(families, ", "); assert families.get(0).getAllColumns().first().getMarkedForDeleteAt() == 1; // delete marker, just added assert !families.get(1).getAllColumns().first().isMarkedForDelete(); // flushed old version ColumnFamily resolved = ColumnFamily.resolve(families);