This is an automated email from the ASF dual-hosted git repository. marcuse pushed a commit to branch trunk in repository https://gitbox.apache.org/repos/asf/cassandra.git
commit 49abedc2c30c6274339bc203e0ddd10f128dae58 Author: Jeff Jirsa <[email protected]> AuthorDate: Mon Mar 23 10:16:26 2020 +0100 Fix force compaction of wrapping ranges Patch by Jeff Jirsa; reviewed by Benjamin Lerer for CASSANDRA-15664 --- CHANGES.txt | 1 + .../cassandra/db/compaction/CompactionManager.java | 18 +++++++- .../compaction/LeveledCompactionStrategyTest.java | 52 ++++++++++++++++++++++ 3 files changed, 69 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 55243b8..65111d0 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0-alpha4 + * Fix force compaction of wrapping ranges (CASSANDRA-15664) * Expose repair streaming metrics (CASSANDRA-15656) * Set now in seconds in the future for validation repairs (CASSANDRA-15655) * Emit metric on preview repair failure (CASSANDRA-15654) diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 28db027..7924a1f 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -33,6 +33,7 @@ import com.google.common.base.Preconditions; import com.google.common.collect.*; import com.google.common.util.concurrent.*; +import org.apache.cassandra.dht.AbstractBounds; import org.apache.cassandra.locator.RangesAtEndpoint; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -929,8 +930,21 @@ public class CompactionManager implements CompactionManagerMBean for (Range<Token> tokenRange : tokenRangeCollection) { - Iterable<SSTableReader> ssTableReaders = View.sstablesInBounds(tokenRange.left.minKeyBound(), tokenRange.right.maxKeyBound(), tree); - Iterables.addAll(sstables, ssTableReaders); + if (!AbstractBounds.strictlyWrapsAround(tokenRange.left, tokenRange.right)) + { + Iterable<SSTableReader> ssTableReaders = View.sstablesInBounds(tokenRange.left.minKeyBound(), tokenRange.right.maxKeyBound(), tree); + Iterables.addAll(sstables, ssTableReaders); + } + else + { + // Searching an interval tree will not return the correct results for a wrapping range + // so we have to unwrap it first + for (Range<Token> unwrappedRange : tokenRange.unwrap()) + { + Iterable<SSTableReader> ssTableReaders = View.sstablesInBounds(unwrappedRange.left.minKeyBound(), unwrappedRange.right.maxKeyBound(), tree); + Iterables.addAll(sstables, ssTableReaders); + } + } } return sstables; } diff --git a/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java b/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java index b925bab..8a8ed13 100644 --- a/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java @@ -460,6 +460,58 @@ public class LeveledCompactionStrategyTest // the 11 tables containing key1 should all compact to 1 table assertEquals(1, cfs.getLiveSSTables().size()); + // Set it up again + cfs.truncateBlocking(); + + // create 10 sstables that contain data for both key1 and key2 + for (int i = 0; i < numIterations; i++) + { + for (DecoratedKey key : keys) + { + UpdateBuilder update = UpdateBuilder.create(cfs.metadata(), key); + for (int c = 0; c < columns; c++) + update.newRow("column" + c).add("val", value); + update.applyUnsafe(); + } + cfs.forceBlockingFlush(); + } + + // create 20 more sstables with 10 containing data for key1 and other 10 containing data for key2 + for (int i = 0; i < numIterations; i++) + { + for (DecoratedKey key : keys) + { + UpdateBuilder update = UpdateBuilder.create(cfs.metadata(), key); + for (int c = 0; c < columns; c++) + update.newRow("column" + c).add("val", value); + update.applyUnsafe(); + cfs.forceBlockingFlush(); + } + } + + // We should have a total of 30 sstables again + assertEquals(30, cfs.getLiveSSTables().size()); + + // This time, we're going to make sure the token range wraps around, to cover the full range + Range<Token> wrappingRange; + if (key1.getToken().compareTo(key2.getToken()) < 0) + { + wrappingRange = new Range<>(key2.getToken(), key1.getToken()); + } + else + { + wrappingRange = new Range<>(key1.getToken(), key2.getToken()); + } + Collection<Range<Token>> wrappingRanges = new ArrayList<>(Arrays.asList(wrappingRange)); + cfs.forceCompactionForTokenRange(wrappingRanges); + + while(CompactionManager.instance.isCompacting(Arrays.asList(cfs), (sstable) -> true)) + { + Thread.sleep(100); + } + + // should all compact to 1 table + assertEquals(1, cfs.getLiveSSTables().size()); } @Test --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
