Fix flaky TimeWindowCompactionStrategyTest

patch by Andrés de la Peña; reviewed by Ekaterina Dimitrova for CASSANDRA-16680
This commit is contained in:
Andrés de la Peña 2021-06-01 16:23:13 +01:00
parent 5f09226a0d
commit 844c8b0348
1 changed files with 36 additions and 32 deletions

View File

@ -26,13 +26,12 @@ import java.util.concurrent.TimeUnit;
import com.google.common.collect.HashMultimap; import com.google.common.collect.HashMultimap;
import com.google.common.collect.Iterables; 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.BeforeClass;
import org.junit.Test; import org.junit.Test;
import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull; import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue; 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.getWindowBoundsInMillis;
import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.newestBucket; import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.newestBucket;
import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.validateOptions; import static org.apache.cassandra.db.compaction.TimeWindowCompactionStrategy.validateOptions;
import static org.apache.cassandra.utils.FBUtilities.nowInSeconds;
public class TimeWindowCompactionStrategyTest extends SchemaLoader 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 String CF_STANDARD1 = "Standard1";
private static final int TTL_SECONDS = 10;
@BeforeClass @BeforeClass
public static void defineSchema() throws ConfigurationException public static void defineSchema() throws ConfigurationException
@ -86,7 +87,10 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
validateOptions(options); validateOptions(options);
fail(String.format("%s == 0 should be rejected", TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_SIZE_KEY)); fail(String.format("%s == 0 should be rejected", TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_SIZE_KEY));
} }
catch (ConfigurationException e) {} catch (ConfigurationException e)
{
// expected exception
}
try try
{ {
@ -103,7 +107,7 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
{ {
options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY, "MONTHS"); options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY, "MONTHS");
validateOptions(options); 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) catch (ConfigurationException e)
{ {
@ -119,24 +123,21 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
@Test @Test
public void testTimeWindows() public void testTimeWindows()
{ {
Long tstamp1 = 1451001601000L; // 2015-12-25 @ 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 tstamp2 = 1451088001000L; // 2015-12-26 @ 00:00:01, in milliseconds
Long lowHour = 1451001600000L; // 2015-12-25 @ 00:00:00, 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 // 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 // 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 // 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 // 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); assertEquals(0, getWindowBoundsInMillis(TimeUnit.DAYS, 2, tstamp2).left.compareTo(lowHour));
return;
} }
@Test @Test
@ -148,8 +149,8 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
cfs.disableAutoCompaction(); cfs.disableAutoCompaction();
ByteBuffer value = ByteBuffer.wrap(new byte[100]); ByteBuffer value = ByteBuffer.wrap(new byte[100]);
Long tstamp = System.currentTimeMillis(); long tstamp = System.currentTimeMillis();
Long tstamp2 = tstamp - (2L * 3600L * 1000L); long tstamp2 = tstamp - (2L * 3600L * 1000L);
// create 5 sstables // create 5 sstables
for (int r = 0; r < 3; r++) for (int r = 0; r < 3; r++)
@ -187,7 +188,7 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
assertTrue("incoming bucket should not be accepted when it has below the min threshold SSTables", newBucket.isEmpty()); 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); 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) // And 2 into the second bucket (1 hour back)
for (int i = 3; i < 5; i++) for (int i = 3; i < 5; i++)
@ -237,9 +238,9 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
ByteBuffer value = ByteBuffer.wrap(new byte[100]); ByteBuffer value = ByteBuffer.wrap(new byte[100]);
// create 2 sstables // Create a expiring sstable with a TTL
DecoratedKey key = Util.dk(String.valueOf("expired")); DecoratedKey key = Util.dk("expired");
new RowUpdateBuilder(cfs.metadata, System.currentTimeMillis(), 1, key.getKey()) new RowUpdateBuilder(cfs.metadata, System.currentTimeMillis(), TTL_SECONDS, key.getKey())
.clustering("column") .clustering("column")
.add("val", value).build().applyUnsafe(); .add("val", value).build().applyUnsafe();
@ -247,7 +248,8 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
SSTableReader expiredSSTable = cfs.getLiveSSTables().iterator().next(); SSTableReader expiredSSTable = cfs.getLiveSSTables().iterator().next();
Thread.sleep(10); 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()) new RowUpdateBuilder(cfs.metadata, System.currentTimeMillis(), key.getKey())
.clustering("column") .clustering("column")
.add("val", value).build().applyUnsafe(); .add("val", value).build().applyUnsafe();
@ -256,7 +258,6 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
assertEquals(cfs.getLiveSSTables().size(), 2); assertEquals(cfs.getLiveSSTables().size(), 2);
Map<String, String> options = new HashMap<>(); Map<String, String> options = new HashMap<>();
options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_SIZE_KEY, "30"); options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_SIZE_KEY, "30");
options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY, "SECONDS"); options.put(TimeWindowCompactionStrategyOptions.COMPACTION_WINDOW_UNIT_KEY, "SECONDS");
options.put(TimeWindowCompactionStrategyOptions.TIMESTAMP_RESOLUTION_KEY, "MILLISECONDS"); options.put(TimeWindowCompactionStrategyOptions.TIMESTAMP_RESOLUTION_KEY, "MILLISECONDS");
@ -264,10 +265,13 @@ public class TimeWindowCompactionStrategyTest extends SchemaLoader
TimeWindowCompactionStrategy twcs = new TimeWindowCompactionStrategy(cfs, options); TimeWindowCompactionStrategy twcs = new TimeWindowCompactionStrategy(cfs, options);
for (SSTableReader sstable : cfs.getLiveSSTables()) for (SSTableReader sstable : cfs.getLiveSSTables())
twcs.addSSTable(sstable); twcs.addSSTable(sstable);
twcs.startup(); twcs.startup();
assertNull(twcs.getNextBackgroundTask((int) (System.currentTimeMillis() / 1000))); assertNull(twcs.getNextBackgroundTask(nowInSeconds()));
Thread.sleep(2000);
AbstractCompactionTask t = twcs.getNextBackgroundTask((int) (System.currentTimeMillis()/1000)); // Wait for the expiration of the first sstable
Thread.sleep(TimeUnit.SECONDS.toMillis(TTL_SECONDS + 1));
AbstractCompactionTask t = twcs.getNextBackgroundTask(nowInSeconds());
assertNotNull(t); assertNotNull(t);
assertEquals(1, Iterables.size(t.transaction.originals())); assertEquals(1, Iterables.size(t.transaction.originals()));
SSTableReader sstable = t.transaction.originals().iterator().next(); SSTableReader sstable = t.transaction.originals().iterator().next();