diff --git a/test/unit/org/apache/cassandra/db/compaction/TimeWindowCompactionStrategyTest.java b/test/unit/org/apache/cassandra/db/compaction/TimeWindowCompactionStrategyTest.java index c57e0cab42..051e7c0432 100644 --- a/test/unit/org/apache/cassandra/db/compaction/TimeWindowCompactionStrategyTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/TimeWindowCompactionStrategyTest.java @@ -26,13 +26,12 @@ import java.util.concurrent.TimeUnit; import com.google.common.collect.HashMultimap; import com.google.common.collect.Iterables; -import com.google.common.collect.Iterables; -import com.google.common.collect.Lists; import org.junit.BeforeClass; import org.junit.Test; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; @@ -52,11 +51,13 @@ import org.apache.cassandra.utils.Pair; import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.getWindowBoundsInMillis; import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.newestBucket; import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.validateOptions; +import static org.apache.cassandra.utils.FBUtilities.nowInSeconds; public class TimeWindowCompactionStrategyTest extends SchemaLoader { - public static final String KEYSPACE1 = "Keyspace1"; + private static final String KEYSPACE1 = "Keyspace1"; private static final String CF_STANDARD1 = "Standard1"; + private static final int TTL_SECONDS = 10; @BeforeClass public static void defineSchema() throws ConfigurationException @@ -86,7 +87,10 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader validateOptions(options); fail(String.format("%s == 0 should be rejected", TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_SIZE_KEY)); } - catch (ConfigurationException e) {} + catch (ConfigurationException e) + { + // expected exception + } try { @@ -103,7 +107,7 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader { options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY, "MONTHS"); validateOptions(options); - fail(String.format("Invalid time units should be rejected", TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY)); + fail(String.format("Invalid %s should be rejected", TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY)); } catch (ConfigurationException e) { @@ -119,24 +123,21 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader @Test public void testTimeWindows() { - Long tstamp1 = 1451001601000L; // 2015-12-25 @ 00:00:01, in milliseconds - Long tstamp2 = 1451088001000L; // 2015-12-26 @ 00:00:01, in milliseconds + long tstamp1 = 1451001601000L; // 2015-12-25 @ 00:00:01, in milliseconds + long tstamp2 = 1451088001000L; // 2015-12-26 @ 00:00:01, in milliseconds Long lowHour = 1451001600000L; // 2015-12-25 @ 00:00:00, in milliseconds // A 1 hour window should round down to the beginning of the hour - assertTrue(getWindowBoundsInMillis(TimeUnit.HOURS, 1, tstamp1).left.compareTo(lowHour) == 0); + assertEquals(0, getWindowBoundsInMillis(TimeUnit.HOURS, 1, tstamp1).left.compareTo(lowHour)); // A 1 minute window should round down to the beginning of the hour - assertTrue(getWindowBoundsInMillis(TimeUnit.MINUTES, 1, tstamp1).left.compareTo(lowHour) == 0); + assertEquals(0, getWindowBoundsInMillis(TimeUnit.MINUTES, 1, tstamp1).left.compareTo(lowHour)); // A 1 day window should round down to the beginning of the hour - assertTrue(getWindowBoundsInMillis(TimeUnit.DAYS, 1, tstamp1).left.compareTo(lowHour) == 0 ); + assertEquals(0, getWindowBoundsInMillis(TimeUnit.DAYS, 1, tstamp1).left.compareTo(lowHour)); // The 2 day window of 2015-12-25 + 2015-12-26 should round down to the beginning of 2015-12-25 - assertTrue(getWindowBoundsInMillis(TimeUnit.DAYS, 2, tstamp2).left.compareTo(lowHour) == 0); - - - return; + assertEquals(0, getWindowBoundsInMillis(TimeUnit.DAYS, 2, tstamp2).left.compareTo(lowHour)); } @Test @@ -148,8 +149,8 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader cfs.disableAutoCompaction(); ByteBuffer value = ByteBuffer.wrap(new byte[100]); - Long tstamp = System.currentTimeMillis(); - Long tstamp2 = tstamp - (2L * 3600L * 1000L); + long tstamp = System.currentTimeMillis(); + long tstamp2 = tstamp - (2L * 3600L * 1000L); // create 5 sstables for (int r = 0; r < 3; r++) @@ -178,21 +179,21 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader List sstrs = new ArrayList<>(cfs.getLiveSSTables()); // We'll put 3 sstables into the newest bucket - for (int i = 0 ; i < 3; i++) + for (int i = 0; i < 3; i++) { - Pair bounds = getWindowBoundsInMillis(TimeUnit.HOURS, 1, tstamp ); + Pair bounds = getWindowBoundsInMillis(TimeUnit.HOURS, 1, tstamp ); buckets.put(bounds.left, sstrs.get(i)); } List newBucket = newestBucket(buckets, 4, 32, TimeUnit.HOURS, 1, new SizeTieredCompactionStrategyOptions(), getWindowBoundsInMillis(TimeUnit.HOURS, 1, System.currentTimeMillis()).left ); assertTrue("incoming bucket should not be accepted when it has below the min threshold SSTables", newBucket.isEmpty()); newBucket = newestBucket(buckets, 2, 32, TimeUnit.HOURS, 1, new SizeTieredCompactionStrategyOptions(), getWindowBoundsInMillis(TimeUnit.HOURS, 1, System.currentTimeMillis()).left); - assertTrue("incoming bucket should be accepted when it is larger than the min threshold SSTables", !newBucket.isEmpty()); + assertFalse("incoming bucket should be accepted when it is larger than the min threshold SSTables", newBucket.isEmpty()); // And 2 into the second bucket (1 hour back) - for (int i = 3 ; i < 5; i++) + for (int i = 3; i < 5; i++) { - Pair bounds = getWindowBoundsInMillis(TimeUnit.HOURS, 1, tstamp2 ); + Pair bounds = getWindowBoundsInMillis(TimeUnit.HOURS, 1, tstamp2 ); buckets.put(bounds.left, sstrs.get(i)); } @@ -205,7 +206,7 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader for (int r = 5; r < numSSTables; r++) { DecoratedKey key = Util.dk(String.valueOf(r)); - for(int i = 0 ; i < r ; i++) + for(int i = 0; i < r; i++) { new RowUpdateBuilder(cfs.metadata, tstamp + r, key.getKey()) .clustering("column") @@ -216,9 +217,9 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader // Reset the buckets, overfill it now sstrs = new ArrayList<>(cfs.getLiveSSTables()); - for (int i = 0 ; i < 40; i++) + for (int i = 0; i < 40; i++) { - Pair bounds = getWindowBoundsInMillis(TimeUnit.HOURS, 1, sstrs.get(i).getMaxTimestamp()); + Pair bounds = getWindowBoundsInMillis(TimeUnit.HOURS, 1, sstrs.get(i).getMaxTimestamp()); buckets.put(bounds.left, sstrs.get(i)); } @@ -237,9 +238,9 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader ByteBuffer value = ByteBuffer.wrap(new byte[100]); - // create 2 sstables - DecoratedKey key = Util.dk(String.valueOf("expired")); - new RowUpdateBuilder(cfs.metadata, System.currentTimeMillis(), 1, key.getKey()) + // Create a expiring sstable with a TTL + DecoratedKey key = Util.dk("expired"); + new RowUpdateBuilder(cfs.metadata, System.currentTimeMillis(), TTL_SECONDS, key.getKey()) .clustering("column") .add("val", value).build().applyUnsafe(); @@ -247,7 +248,8 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader SSTableReader expiredSSTable = cfs.getLiveSSTables().iterator().next(); Thread.sleep(10); - key = Util.dk(String.valueOf("nonexpired")); + // Create a second sstable without TTL + key = Util.dk("nonexpired"); new RowUpdateBuilder(cfs.metadata, System.currentTimeMillis(), key.getKey()) .clustering("column") .add("val", value).build().applyUnsafe(); @@ -256,7 +258,6 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader assertEquals(cfs.getLiveSSTables().size(), 2); Map options = new HashMap<>(); - options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_SIZE_KEY, "30"); options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY, "SECONDS"); options.put(TimeWindowCompactionStrategyOptions.TIMESTAMP_RESOLUTION_KEY, "MILLISECONDS"); @@ -264,10 +265,13 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader TimeWindowCompactionStrategy twcs = new TimeWindowCompactionStrategy(cfs, options); for (SSTableReader sstable : cfs.getLiveSSTables()) twcs.addSSTable(sstable); + twcs.startup(); - assertNull(twcs.getNextBackgroundTask((int) (System.currentTimeMillis() / 1000))); - Thread.sleep(2000); - AbstractCompactionTask t = twcs.getNextBackgroundTask((int) (System.currentTimeMillis()/1000)); + assertNull(twcs.getNextBackgroundTask(nowInSeconds())); + + // Wait for the expiration of the first sstable + Thread.sleep(TimeUnit.SECONDS.toMillis(TTL_SECONDS + 1)); + AbstractCompactionTask t = twcs.getNextBackgroundTask(nowInSeconds()); assertNotNull(t); assertEquals(1, Iterables.size(t.transaction.originals())); SSTableReader sstable = t.transaction.originals().iterator().next();