diff --git a/CHANGES.txt b/CHANGES.txt index 19d7f7e19a..11f6573e7f 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -21,6 +21,7 @@ * ReadResponseResolver check digests against each other (CASSANDRA-1830) * change exception for read requests during bootstrap from InvalidRequest to Unavailable (CASSANDRA-1862) + * make compaction buckets deterministic (CASSANDRA-1265) 0.6.8 diff --git a/src/java/org/apache/cassandra/db/CompactionManager.java b/src/java/org/apache/cassandra/db/CompactionManager.java index ec122e12fb..48bf51c64b 100644 --- a/src/java/org/apache/cassandra/db/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/CompactionManager.java @@ -42,6 +42,7 @@ import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.service.AntiEntropyService; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.Pair; import org.apache.log4j.Logger; import org.cliffc.high_scale_lib.NonBlockingHashMap; @@ -88,7 +89,9 @@ public class CompactionManager implements CompactionManagerMBean return 0; } logger.debug("Checking to see if compaction of " + cfs.columnFamily_ + " would be useful"); - Set> buckets = getBuckets(cfs.getSSTables(), 50L * 1024L * 1024L); + + Set> buckets = getBuckets( + convertSSTablesToPairs(cfs.getSSTables()), 50L * 1024L * 1024L); updateEstimateFor(cfs, buckets); for (List sstables : buckets) @@ -467,29 +470,43 @@ public class CompactionManager implements CompactionManagerMBean /* * Group files of similar size into buckets. */ - static Set> getBuckets(Iterable files, long min) + static Set> getBuckets(Iterable> files, long min) { - Map, Long> buckets = new HashMap, Long>(); - for (SSTableReader sstable : files) + // Sort the list in order to get deterministic results during the grouping below + List> sortedFiles = new ArrayList>(); + for (Pair pair: files) + sortedFiles.add(pair); + + Collections.sort(sortedFiles, new Comparator>() { - long size = sstable.length(); + public int compare(Pair p1, Pair p2) + { + return p1.right.compareTo(p2.right); + } + }); + + Map, Long> buckets = new HashMap, Long>(); + + for (Pair pair: sortedFiles) + { + long size = pair.right; boolean bFound = false; // look for a bucket containing similar-sized files: // group in the same bucket if it's w/in 50% of the average for this bucket, // or this file and the bucket are all considered "small" (less than `min`) - for (Entry, Long> entry : buckets.entrySet()) + for (Entry, Long> entry : buckets.entrySet()) { - List bucket = entry.getKey(); + List bucket = entry.getKey(); long averageSize = entry.getValue(); - if ((size > averageSize / 2 && size < 3 * averageSize / 2) + if ((size > (averageSize / 2) && size < (3 * averageSize) / 2) || (size < min && averageSize < min)) { // remove and re-add because adding changes the hash buckets.remove(bucket); long totalSize = bucket.size() * averageSize; averageSize = (totalSize + size) / (bucket.size() + 1); - bucket.add(sstable); + bucket.add(pair.left); buckets.put(bucket, averageSize); bFound = true; break; @@ -498,8 +515,8 @@ public class CompactionManager implements CompactionManagerMBean // no similar bucket found; put it in a new one if (!bFound) { - ArrayList bucket = new ArrayList(); - bucket.add(sstable); + ArrayList bucket = new ArrayList(); + bucket.add(pair.left); buckets.put(bucket, size); } } @@ -507,6 +524,14 @@ public class CompactionManager implements CompactionManagerMBean return buckets.keySet(); } + private static Collection> convertSSTablesToPairs(Collection collection) + { + Collection> tablePairs = new HashSet>(); + for(SSTableReader table: collection) + tablePairs.add(new Pair(table, table.length())); + return tablePairs; + } + public static int getDefaultGCBefore() { return (int)(System.currentTimeMillis() / 1000) - DatabaseDescriptor.getGcGraceInSeconds(); @@ -565,7 +590,9 @@ public class CompactionManager implements CompactionManagerMBean public void run () { logger.debug("Estimating compactions for " + cfs.columnFamily_); - final Set> buckets = getBuckets(cfs.getSSTables(), 50L * 1024L * 1024L); + + final Set> buckets = + getBuckets(convertSSTablesToPairs(cfs.getSSTables()), 50L * 1024L * 1024L); updateEstimateFor(cfs, buckets); } }; diff --git a/test/unit/org/apache/cassandra/db/CompactionsTest.java b/test/unit/org/apache/cassandra/db/CompactionsTest.java index 165d29b095..5f42618e4b 100644 --- a/test/unit/org/apache/cassandra/db/CompactionsTest.java +++ b/test/unit/org/apache/cassandra/db/CompactionsTest.java @@ -24,6 +24,8 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.Set; import java.util.HashSet; +import java.util.List; +import java.util.ArrayList; import org.apache.cassandra.Util; @@ -33,6 +35,7 @@ import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.CleanupHelper; import org.apache.cassandra.db.filter.QueryPath; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.Pair; import static junit.framework.Assert.assertEquals; public class CompactionsTest extends CleanupHelper @@ -75,4 +78,57 @@ public class CompactionsTest extends CleanupHelper } assertEquals(inserted.size(), Util.getRangeSlice(store).rows.size()); } + + @Test + public void testGetBuckets() + { + List> pairs = new ArrayList>(); + String[] strings = { "a", "bbbb", "cccccccc", "cccccccc", "bbbb", "a" }; + for (int i = 0; i < strings.length; i++) { + Pair pair = new Pair(strings[i], new Long(strings[i].length())); + pairs.add(pair); + } + + Set> buckets = CompactionManager.getBuckets(pairs, 2); + assertEquals(3, buckets.size()); + + for(List bucket: buckets) + { + assertEquals(2, bucket.size()); + assertEquals(bucket.get(0).length(), bucket.get(1).length()); + assertEquals(bucket.get(0).charAt(0), bucket.get(1).charAt(0)); + } + + pairs.clear(); + buckets.clear(); + + String[] strings2 = { "aaa", "bbbbbbbb", "aaa", "bbbbbbbb", "bbbbbbbb", "aaa" }; + for (int i = 0; i < strings2.length; i++) { + Pair pair = new Pair(strings2[i], new Long(strings2[i].length())); + pairs.add(pair); + } + + buckets = CompactionManager.getBuckets(pairs, 2); + assertEquals(2, buckets.size()); + + for(List bucket: buckets) + { + assertEquals(3, bucket.size()); + assertEquals(bucket.get(0).charAt(0), bucket.get(1).charAt(0)); + assertEquals(bucket.get(1).charAt(0), bucket.get(2).charAt(0)); + } + + // Test the "min" functionality + pairs.clear(); + buckets.clear(); + + String[] strings3 = { "aaa", "bbbbbbbb", "aaa", "bbbbbbbb", "bbbbbbbb", "aaa" }; + for (int i = 0; i < strings3.length; i++) { + Pair pair = new Pair(strings3[i], new Long(strings3[i].length())); + pairs.add(pair); + } + + buckets = CompactionManager.getBuckets(pairs, 10); // notice the min is 10 + assertEquals(1, buckets.size()); + } }