diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 200b304605..d0f7c20afe 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -638,7 +638,7 @@ public final class ColumnFamilyStore implements ColumnFamilyStoreMBean /* it's ok if compaction gets submitted multiple times while one is already in process. worst that happens is, compactor will count the sstable files and decide there are not enough to bother with. */ - if (ssTableCount >= MinorCompactionManager.COMPACTION_THRESHOLD) + if (ssTableCount >= MinorCompactionManager.MINCOMPACTION_THRESHOLD) { if (logger_.isDebugEnabled()) logger_.debug("Submitting " + columnFamily_ + " for compaction"); @@ -669,7 +669,7 @@ public final class ColumnFamilyStore implements ColumnFamilyStoreMBean /* * Group files of similar size into buckets. */ - static Set> getCompactionBuckets(List files, long min) + static Set> getCompactionBuckets(Iterable files, long min) { Map, Long> buckets = new ConcurrentHashMap, Long>(); for (String fname : files) @@ -711,32 +711,17 @@ public final class ColumnFamilyStore implements ColumnFamilyStoreMBean /* * Break the files into buckets and then compact. */ - int doCompaction(int threshold) throws IOException + int doCompaction(int minThreshold, int maxThreshold) throws IOException { - List files = new ArrayList(ssTables_.keySet()); int filesCompacted = 0; - Set> buckets = getCompactionBuckets(files, 50L * 1024L * 1024L); - for (List fileList : buckets) + for (List files : getCompactionBuckets(ssTables_.keySet(), 50L * 1024L * 1024L)) { - Collections.sort(fileList, new FileNameComparator(FileNameComparator.Ascending)); - if (fileList.size() < threshold) + if (files.size() < minThreshold) { continue; } - // For each bucket if it has crossed the threshhold do the compaction - // In case of range compaction merge the counting bloom filters also. - files.clear(); - int count = 0; - for (String file : fileList) - { - files.add(file); - count++; - if (count == threshold) - { - filesCompacted += doFileCompaction(files, BUFSIZE); - break; - } - } + Collections.sort(files, new FileNameComparator(FileNameComparator.Ascending)); + filesCompacted += doFileCompaction(files.subList(0, Math.min(files.size(), maxThreshold)), BUFSIZE); } return filesCompacted; } diff --git a/src/java/org/apache/cassandra/db/MinorCompactionManager.java b/src/java/org/apache/cassandra/db/MinorCompactionManager.java index c303dd9429..3fe2ed70fb 100644 --- a/src/java/org/apache/cassandra/db/MinorCompactionManager.java +++ b/src/java/org/apache/cassandra/db/MinorCompactionManager.java @@ -21,11 +21,8 @@ package org.apache.cassandra.db; import java.io.IOException; import java.util.List; import java.util.concurrent.Callable; -import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; -import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -33,7 +30,7 @@ import org.apache.cassandra.concurrent.DebuggableScheduledThreadPoolExecutor; import org.apache.cassandra.concurrent.ThreadFactoryImpl; import org.apache.cassandra.dht.Range; import org.apache.cassandra.net.EndPoint; -import org.apache.cassandra.service.StorageService; + import org.apache.log4j.Logger; class MinorCompactionManager @@ -42,7 +39,8 @@ class MinorCompactionManager private static Lock lock_ = new ReentrantLock(); private static Logger logger_ = Logger.getLogger(MinorCompactionManager.class); private static final long intervalInMins_ = 5; - static final int COMPACTION_THRESHOLD = 4; // compact this many sstables at a time + static final int MINCOMPACTION_THRESHOLD = 4; // compact this many sstables min at a time + static final int MAXCOMPACTION_THRESHOLD = 32; // compact this many sstables max at a time public static MinorCompactionManager instance() { @@ -155,16 +153,16 @@ class MinorCompactionManager public Future submit(final ColumnFamilyStore columnFamilyStore) { - return submit(columnFamilyStore, COMPACTION_THRESHOLD); + return submit(columnFamilyStore, MINCOMPACTION_THRESHOLD, MAXCOMPACTION_THRESHOLD); } - Future submit(final ColumnFamilyStore columnFamilyStore, final int threshold) + Future submit(final ColumnFamilyStore columnFamilyStore, final int minThreshold, final int maxThreshold) { Callable callable = new Callable() { public Integer call() throws IOException { - return columnFamilyStore.doCompaction(threshold); + return columnFamilyStore.doCompaction(minThreshold, maxThreshold); } }; return compactor_.submit(callable); diff --git a/test/unit/org/apache/cassandra/db/CompactionsTest.java b/test/unit/org/apache/cassandra/db/CompactionsTest.java index fa6afd6201..1ab34e912f 100644 --- a/test/unit/org/apache/cassandra/db/CompactionsTest.java +++ b/test/unit/org/apache/cassandra/db/CompactionsTest.java @@ -23,7 +23,6 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.Set; import java.util.HashSet; -import java.util.Arrays; import org.junit.Test; @@ -62,7 +61,7 @@ public class CompactionsTest extends CleanupHelper } if (store.getSSTables().size() > 1) { - store.doCompaction(store.getSSTables().size()); + store.doCompaction(2, store.getSSTables().size()); } assertEquals(table.getColumnFamilyStore("Standard1").getKeyRange("", "", 10000).keys.size(), inserted.size()); } diff --git a/test/unit/org/apache/cassandra/db/OneCompactionTest.java b/test/unit/org/apache/cassandra/db/OneCompactionTest.java index 41fd40aad1..61df9b2f5a 100644 --- a/test/unit/org/apache/cassandra/db/OneCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/OneCompactionTest.java @@ -23,7 +23,6 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.Set; import java.util.HashSet; -import java.util.Arrays; import org.junit.Test; @@ -47,7 +46,7 @@ public class OneCompactionTest store.forceBlockingFlush(); assertEquals(inserted.size(), table.getColumnFamilyStore(columnFamilyName).getKeyRange("", "", 10000).keys.size()); } - Future ft = MinorCompactionManager.instance().submit(store, 2); + Future ft = MinorCompactionManager.instance().submit(store, 2, 32); ft.get(); assertEquals(1, store.getSSTables().size()); assertEquals(table.getColumnFamilyStore(columnFamilyName).getKeyRange("", "", 10000).keys.size(), inserted.size()); diff --git a/test/unit/org/apache/cassandra/db/RemoveSuperColumnTest.java b/test/unit/org/apache/cassandra/db/RemoveSuperColumnTest.java index 9839d0eba5..e90e3e009f 100644 --- a/test/unit/org/apache/cassandra/db/RemoveSuperColumnTest.java +++ b/test/unit/org/apache/cassandra/db/RemoveSuperColumnTest.java @@ -57,7 +57,7 @@ public class RemoveSuperColumnTest store.forceBlockingFlush(); validateRemoveTwoSources(); - Future ft = MinorCompactionManager.instance().submit(store, 2); + Future ft = MinorCompactionManager.instance().submit(store, 2, 32); ft.get(); assertEquals(1, store.getSSTables().size()); validateRemoveCompacted(); @@ -109,7 +109,7 @@ public class RemoveSuperColumnTest store.forceBlockingFlush(); validateRemoveWithNewData(); - Future ft = MinorCompactionManager.instance().submit(store, 2); + Future ft = MinorCompactionManager.instance().submit(store, 2, 32); ft.get(); assertEquals(1, store.getSSTables().size()); validateRemoveWithNewData(); diff --git a/test/unit/org/apache/cassandra/db/TableTest.java b/test/unit/org/apache/cassandra/db/TableTest.java index 0ef1146a65..6abd7ad23a 100644 --- a/test/unit/org/apache/cassandra/db/TableTest.java +++ b/test/unit/org/apache/cassandra/db/TableTest.java @@ -304,7 +304,7 @@ public class TableTest extends CleanupHelper // compact so we have a big row with more than the minimum index count if (cfStore.getSSTables().size() > 1) { - cfStore.doCompaction(cfStore.getSSTables().size()); + cfStore.doCompaction(2, cfStore.getSSTables().size()); } SSTableReader sstable = cfStore.getSSTables().iterator().next(); long position = sstable.getPosition(key);