From 7d6a60ddd0e156f6028c05d97fa9179f86fe02ee Mon Sep 17 00:00:00 2001 From: Benedict Elliott Smith Date: Fri, 15 May 2015 11:19:23 -0500 Subject: [PATCH] Fix canonical view returning early opened SSTables patch by benedict; reviewed by yukim for CASSANDRA-9396 --- CHANGES.txt | 1 + .../cassandra/db/ColumnFamilyStore.java | 8 ++-- .../io/sstable/SSTableRewriterTest.java | 47 ++++++++++++++++++- 3 files changed, 52 insertions(+), 4 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 80619d0721..129f6a112f 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -23,6 +23,7 @@ * Fix anticompaction blocking ANTI_ENTROPY stage (CASSANDRA-9151) * Repair waits for anticompaction to finish (CASSANDRA-9097) * Fix streaming not holding ref when stream error (CASSANDRA-9295) + * Fix canonical view returning early opened SSTables (CASSANDRA-9396) Merged from 2.0: * Clone SliceQueryFilter in AbstractReadCommand implementations (CASSANDRA-8940) * Push correct protocol notification for DROP INDEX (CASSANDRA-9310) diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 978037e448..bdc2d8bf84 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -2892,10 +2892,12 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean public List apply(DataTracker.View view) { List sstables = new ArrayList<>(); - sstables.addAll(view.compacting); + for (SSTableReader sstable : view.compacting) + if (sstable.openReason != SSTableReader.OpenReason.EARLY) + sstables.add(sstable); for (SSTableReader sstable : view.sstables) - if (!view.compacting.contains(sstable) && sstable.openReason != SSTableReader.OpenReason.EARLY) - sstables.add(sstable); + if (!view.compacting.contains(sstable) && sstable.openReason != SSTableReader.OpenReason.EARLY) + sstables.add(sstable); return sstables; } }; diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableRewriterTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableRewriterTest.java index 09937bc047..b940b7b3eb 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableRewriterTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableRewriterTest.java @@ -18,6 +18,7 @@ package org.apache.cassandra.io.sstable; import java.io.File; +import java.io.IOException; import java.nio.ByteBuffer; import java.util.*; import java.util.concurrent.ThreadLocalRandom; @@ -40,10 +41,15 @@ import org.apache.cassandra.db.compaction.SSTableSplitter; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.io.sstable.metadata.MetadataCollector; +import org.apache.cassandra.io.util.DataIntegrityMetadata; +import org.apache.cassandra.io.util.RandomAccessReader; import org.apache.cassandra.metrics.StorageMetrics; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.Pair; +import org.apache.cassandra.utils.concurrent.Ref; +import org.apache.cassandra.utils.concurrent.Refs; + import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -81,7 +87,7 @@ public class SSTableRewriterTest extends SchemaLoader } } Collection newsstables = writer.finish(); - cfs.getDataTracker().markCompactedSSTablesReplaced(sstables, newsstables , OperationType.COMPACTION); + cfs.getDataTracker().markCompactedSSTablesReplaced(sstables, newsstables, OperationType.COMPACTION); Thread.sleep(100); validateCFS(cfs); int filecounts = assertFileCounts(sstables.iterator().next().descriptor.directory.list(), 0, 0); @@ -733,6 +739,45 @@ public class SSTableRewriterTest extends SchemaLoader validateCFS(cfs); } + @Test + public void testCanonicalView() throws IOException + { + Keyspace keyspace = Keyspace.open(KEYSPACE); + ColumnFamilyStore cfs = keyspace.getColumnFamilyStore(CF); + cfs.truncateBlocking(); + + SSTableReader s = writeFile(cfs, 1000); + cfs.addSSTable(s); + Set sstables = Sets.newHashSet(cfs.markAllCompacting()); + assertEquals(1, sstables.size()); + SSTableRewriter.overrideOpenInterval(10000000); + SSTableRewriter writer = new SSTableRewriter(cfs, sstables, 1000, false); + boolean checked = false; + try (AbstractCompactionStrategy.ScannerList scanners = cfs.getCompactionStrategy().getScanners(sstables)) + { + ISSTableScanner scanner = scanners.scanners.get(0); + CompactionController controller = new CompactionController(cfs, sstables, cfs.gcBefore(System.currentTimeMillis())); + writer.switchWriter(getWriter(cfs, sstables.iterator().next().descriptor.directory)); + while (scanner.hasNext()) + { + AbstractCompactedRow row = new LazilyCompactedRow(controller, Collections.singletonList(scanner.next())); + writer.append(row); + if (!checked && writer.currentWriter().getFilePointer() > 15000000) + { + checked = true; + ColumnFamilyStore.ViewFragment viewFragment = cfs.select(ColumnFamilyStore.CANONICAL_SSTABLES); + // canonical view should have only one SSTable which is not opened early. + assertEquals(1, viewFragment.sstables.size()); + SSTableReader sstable = viewFragment.sstables.get(0); + assertEquals(s.descriptor, sstable.descriptor); + assertTrue("Found early opened SSTable in canonical view: " + sstable.getFilename(), sstable.openReason != SSTableReader.OpenReason.EARLY); + } + } + } + writer.finish(); + cfs.getDataTracker().unmarkCompacting(sstables); + } + private void validateKeys(Keyspace ks) { for (int i = 0; i < 100; i++)