From 2b0b61a7b046af305e41fe109ded48bc4b4a0b26 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 17 Apr 2009 01:48:58 +0000 Subject: [PATCH] waitForFlush -> forceBlockingFlush. ServerTest.cleanup now flushes and cleans out all ColumnFamilyStores and commitlog, allowing remove tests to not step on each others' toes (all tests pass now). patch by jbellis; reviewed by Sandeep Tata for #85 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@765831 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/db/ColumnFamilyStore.java | 32 ++++++++++++++++-- src/org/apache/cassandra/db/CommitLog.java | 7 +++- src/org/apache/cassandra/db/Table.java | 2 +- test/org/apache/cassandra/ServerTest.java | 21 +++++++++++- .../cassandra/db/ColumnFamilyStoreTest.java | 33 ++++--------------- 5 files changed, 63 insertions(+), 32 deletions(-) diff --git a/src/org/apache/cassandra/db/ColumnFamilyStore.java b/src/org/apache/cassandra/db/ColumnFamilyStore.java index b20af678e3..00001a66a0 100644 --- a/src/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/org/apache/cassandra/db/ColumnFamilyStore.java @@ -84,7 +84,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean private AtomicReference binaryMemtable_; /* SSTables on disk for this column family */ - Set ssTables_ = new HashSet(); + private Set ssTables_ = new HashSet(); /* Modification lock used for protecting reads from compactions. */ private ReentrantReadWriteLock lock_ = new ReentrantReadWriteLock(true); @@ -433,11 +433,23 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean void forceFlush() throws IOException { - //MemtableManager.instance().submit(getColumnFamilyName(), memtable_.get() , CommitLog.CommitLogContext.NULL); - //memtable_.get().flush(true, CommitLog.CommitLogContext.NULL); memtable_.get().forceflush(this); } + void forceBlockingFlush() throws IOException, ExecutionException, InterruptedException + { + 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() + { + public void run() + { + } + }); + f.get(); + } + void forceFlushBinary() { BinaryMemtableManager.instance().submit(getColumnFamilyName(), binaryMemtable_.get()); @@ -1407,4 +1419,18 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean { return memtableSwitchCount; } + + /** + * clears out all data associated with this ColumnFamily. + * For use in testing. + */ + public void reset() throws IOException, ExecutionException, InterruptedException + { + forceBlockingFlush(); + for (String fName : ssTables_) + { + new File(fName).delete(); + } + ssTables_.clear(); + } } diff --git a/src/org/apache/cassandra/db/CommitLog.java b/src/org/apache/cassandra/db/CommitLog.java index 4aa151b8f3..6097f764f4 100644 --- a/src/org/apache/cassandra/db/CommitLog.java +++ b/src/org/apache/cassandra/db/CommitLog.java @@ -57,7 +57,7 @@ import java.util.concurrent.locks.ReentrantLock; * * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) */ -class CommitLog +public class CommitLog { private static final int bufSize_ = 128*1024*1024; private static Map instances_ = new HashMap(); @@ -623,6 +623,11 @@ class CommitLog forcedRollOver_ = true; } + public static void reset() + { + CommitLog.instances_.clear(); + } + public static void main(String[] args) throws Throwable { LogUtil.init(); diff --git a/src/org/apache/cassandra/db/Table.java b/src/org/apache/cassandra/db/Table.java index ce1af07d71..43613ec03c 100644 --- a/src/org/apache/cassandra/db/Table.java +++ b/src/org/apache/cassandra/db/Table.java @@ -425,7 +425,7 @@ public class Table return columnFamilyStores_; } - ColumnFamilyStore getColumnFamilyStore(String cfName) + public ColumnFamilyStore getColumnFamilyStore(String cfName) { return columnFamilyStores_.get(cfName); } diff --git a/test/org/apache/cassandra/ServerTest.java b/test/org/apache/cassandra/ServerTest.java index 8e66839636..1cad2cc7da 100644 --- a/test/org/apache/cassandra/ServerTest.java +++ b/test/org/apache/cassandra/ServerTest.java @@ -3,15 +3,34 @@ package org.apache.cassandra; import org.testng.annotations.Test; import org.testng.annotations.BeforeMethod; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.Table; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.CommitLog; import java.io.File; +import java.io.IOException; @Test(groups={"serial"}) public class ServerTest { - // TODO clean up static structures too (e.g. memtables) @BeforeMethod public void cleanup() { + Table table = Table.open("Table1"); + for (String cfName : table.getColumnFamilies()) + { + ColumnFamilyStore cfs = table.getColumnFamilyStore(cfName); + try + { + cfs.reset(); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + } + + CommitLog.reset(); + String[] directoryNames = { DatabaseDescriptor.getBootstrapFileLocation(), DatabaseDescriptor.getLogFileLocation(), diff --git a/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java b/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java index 8194fe816e..d6b6fc56fd 100644 --- a/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java +++ b/test/org/apache/cassandra/db/ColumnFamilyStoreTest.java @@ -67,9 +67,8 @@ public class ColumnFamilyStoreTest extends ServerTest // validateNameSort(table); - table.getColumnFamilyStore("Standard1").forceFlush(); - table.getColumnFamilyStore("Super1").forceFlush(); - waitForFlush(); + table.getColumnFamilyStore("Standard1").forceBlockingFlush(); + table.getColumnFamilyStore("Super1").forceBlockingFlush(); validateNameSort(table); } @@ -93,8 +92,7 @@ public class ColumnFamilyStoreTest extends ServerTest validateTimeSort(table); - table.getColumnFamilyStore("StandardByTime1").forceFlush(); - waitForFlush(); + table.getColumnFamilyStore("StandardByTime1").forceBlockingFlush(); validateTimeSort(table); // interleave some new data to test memtable + sstable @@ -154,18 +152,6 @@ public class ColumnFamilyStoreTest extends ServerTest } } - private void waitForFlush() - throws InterruptedException, ExecutionException - { - Future f = MemtableManager.instance().flusher_.submit(new Runnable() - { - public void run() - { - } - }); - f.get(); - } - private void validateNameSort(Table table) throws ColumnFamilyNotDefinedException, IOException { @@ -213,8 +199,7 @@ public class ColumnFamilyStoreTest extends ServerTest rm = new RowMutation("Table1", "key1"); rm.add("Standard1:Column1", "asdf".getBytes(), 0); rm.apply(); - store.forceFlush(); - waitForFlush(); + store.forceBlockingFlush(); // remove rm = new RowMutation("Table1", "key1"); @@ -236,8 +221,7 @@ public class ColumnFamilyStoreTest extends ServerTest rm = new RowMutation("Table1", "key1"); rm.add("Super1:SC1:Column1", "asdf".getBytes(), 0); rm.apply(); - store.forceFlush(); - waitForFlush(); + store.forceBlockingFlush(); // remove rm = new RowMutation("Table1", "key1"); @@ -259,8 +243,7 @@ public class ColumnFamilyStoreTest extends ServerTest rm = new RowMutation("Table1", "key1"); rm.add("Super1:SC1:Column1", "asdf".getBytes(), 0); rm.apply(); - store.forceFlush(); - waitForFlush(); + store.forceBlockingFlush(); // remove rm = new RowMutation("Table1", "key1"); @@ -359,7 +342,6 @@ public class ColumnFamilyStoreTest extends ServerTest { Table table = Table.open("Table1"); ColumnFamilyStore store = table.getColumnFamilyStore("Standard1"); - store.ssTables_.clear(); // TODO integrate this better into test setup/teardown for (int j = 0; j < 5; j++) { for (int i = 0; i < 10; i++) { @@ -369,8 +351,7 @@ public class ColumnFamilyStoreTest extends ServerTest rm.add("Standard1:A", new byte[0], epoch); rm.apply(); } - store.forceFlush(); - waitForFlush(); + store.forceBlockingFlush(); } Future ft = MinorCompactionManager.instance().submit(store); ft.get();