From bc89bc66cb762da2be61b92d56b48154d8bd3cbf Mon Sep 17 00:00:00 2001 From: Sam Tunnicliffe Date: Fri, 16 Oct 2015 17:39:07 +0100 Subject: [PATCH] Partially revert #9839 to remove reference loop Patch by Sam Tunnicliffe; reviewed by Benedict Elliot Smith for CASSANDRA-10543 --- CHANGES.txt | 1 + .../cassandra/db/ColumnFamilyStore.java | 4 +++ .../CompressedRandomAccessReader.java | 16 ++++++----- .../io/sstable/format/SSTableReader.java | 23 ++++++++++++---- .../cassandra/io/util/IChecksummedFile.java | 27 ------------------- .../cassandra/io/util/ICompressedFile.java | 2 +- .../cassandra/io/util/SegmentedFile.java | 14 +--------- .../cassandra/schema/CompressionParams.java | 12 +++++++++ .../miscellaneous/CrcCheckChanceTest.java | 19 ++++++++++++- 9 files changed, 64 insertions(+), 54 deletions(-) delete mode 100644 src/java/org/apache/cassandra/io/util/IChecksummedFile.java diff --git a/CHANGES.txt b/CHANGES.txt index 77facc4718..cb4c2d8039 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 3.0-rc2 + * Remove circular references in SegmentedFile (CASSANDRA-10543) * Ensure validation of indexed values only occurs once per-partition (CASSANDRA-10536) * Fix handling of static columns for range tombstones in thrift (CASSANDRA-10174) * Support empty ColumnFilter for backward compatility on empty IN (CASSANDRA-10471) diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 4c9fc55f58..0b838bfd8b 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -2044,7 +2044,11 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean { TableParams.builder().crcCheckChance(crcCheckChance).build().validate(); for (ColumnFamilyStore cfs : concatWithIndexes()) + { cfs.crcCheckChance.set(crcCheckChance); + for (SSTableReader sstable : cfs.getSSTables(SSTableSet.LIVE)) + sstable.setCrcCheckChance(crcCheckChance); + } } catch (ConfigurationException e) { diff --git a/src/java/org/apache/cassandra/io/compress/CompressedRandomAccessReader.java b/src/java/org/apache/cassandra/io/compress/CompressedRandomAccessReader.java index b2759e6517..329d93256e 100644 --- a/src/java/org/apache/cassandra/io/compress/CompressedRandomAccessReader.java +++ b/src/java/org/apache/cassandra/io/compress/CompressedRandomAccessReader.java @@ -23,6 +23,7 @@ import java.util.concurrent.ThreadLocalRandom; import java.util.zip.Checksum; import java.util.function.Supplier; +import com.google.common.annotations.VisibleForTesting; import com.google.common.primitives.Ints; import org.apache.cassandra.io.FSReadError; @@ -46,14 +47,18 @@ public class CompressedRandomAccessReader extends RandomAccessReader // raw checksum bytes private ByteBuffer checksumBytes; - private final Supplier crcCheckChanceSupplier; + + @VisibleForTesting + public double getCrcCheckChance() + { + return metadata.parameters.getCrcCheckChance(); + } protected CompressedRandomAccessReader(Builder builder) { super(builder); this.metadata = builder.metadata; this.checksum = metadata.checksumType.newInstance(); - crcCheckChanceSupplier = builder.crcCheckChanceSupplier; if (regions == null) { @@ -124,7 +129,7 @@ public class CompressedRandomAccessReader extends RandomAccessReader buffer.flip(); } - if (crcCheckChanceSupplier.get() > ThreadLocalRandom.current().nextDouble()) + if (getCrcCheckChance() > ThreadLocalRandom.current().nextDouble()) { compressed.rewind(); metadata.checksumType.update( checksum, (compressed)); @@ -186,7 +191,7 @@ public class CompressedRandomAccessReader extends RandomAccessReader buffer.flip(); } - if (crcCheckChanceSupplier.get() > ThreadLocalRandom.current().nextDouble()) + if (getCrcCheckChance() > ThreadLocalRandom.current().nextDouble()) { compressedChunk.position(chunkOffset).limit(chunkOffset + chunk.length); @@ -239,21 +244,18 @@ public class CompressedRandomAccessReader extends RandomAccessReader public final static class Builder extends RandomAccessReader.Builder { private final CompressionMetadata metadata; - private final Supplier crcCheckChanceSupplier; public Builder(ICompressedFile file) { super(file.channel()); this.metadata = applyMetadata(file.getMetadata()); this.regions = file.regions(); - this.crcCheckChanceSupplier = file.getCrcCheckChanceSupplier(); } public Builder(ChannelProxy channel, CompressionMetadata metadata) { super(channel); this.metadata = applyMetadata(metadata); - this.crcCheckChanceSupplier = (() -> 1.0); //100% crc_check_chance } private CompressionMetadata applyMetadata(CompressionMetadata metadata) diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java index 8d23597151..afd0a1ed9d 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java @@ -217,6 +217,8 @@ public abstract class SSTableReader extends SSTable implements SelfRefCounted getCrcCheckChanceSupplier(); - public void setCrcCheckChanceSupplier(Supplier crcCheckChanceSupplier); -} diff --git a/src/java/org/apache/cassandra/io/util/ICompressedFile.java b/src/java/org/apache/cassandra/io/util/ICompressedFile.java index c149fd12fc..43cef8c4e3 100644 --- a/src/java/org/apache/cassandra/io/util/ICompressedFile.java +++ b/src/java/org/apache/cassandra/io/util/ICompressedFile.java @@ -19,7 +19,7 @@ package org.apache.cassandra.io.util; import org.apache.cassandra.io.compress.CompressionMetadata; -public interface ICompressedFile extends IChecksummedFile +public interface ICompressedFile { ChannelProxy channel(); CompressionMetadata getMetadata(); diff --git a/src/java/org/apache/cassandra/io/util/SegmentedFile.java b/src/java/org/apache/cassandra/io/util/SegmentedFile.java index c8272556a7..ab2d291bf1 100644 --- a/src/java/org/apache/cassandra/io/util/SegmentedFile.java +++ b/src/java/org/apache/cassandra/io/util/SegmentedFile.java @@ -49,7 +49,7 @@ import static org.apache.cassandra.utils.Throwables.maybeFail; * would need to be longer than 2GB, that segment will not be mmap'd, and a new RandomAccessFile will be created for * each access to that segment. */ -public abstract class SegmentedFile extends SharedCloseableImpl implements IChecksummedFile +public abstract class SegmentedFile extends SharedCloseableImpl { public final ChannelProxy channel; public final int bufferSize; @@ -58,8 +58,6 @@ public abstract class SegmentedFile extends SharedCloseableImpl implements IChec // This differs from length for compressed files (but we still need length for // SegmentIterator because offsets in the file are relative to the uncompressed size) public final long onDiskLength; - private Supplier crcCheckChanceSupplier = () -> 1.0; - /** * Use getBuilder to get a Builder to construct a SegmentedFile. @@ -137,16 +135,6 @@ public abstract class SegmentedFile extends SharedCloseableImpl implements IChec return reader; } - public Supplier getCrcCheckChanceSupplier() - { - return crcCheckChanceSupplier; - } - - public void setCrcCheckChanceSupplier(Supplier crcCheckChanceSupplier) - { - this.crcCheckChanceSupplier = crcCheckChanceSupplier; - } - public void dropPageCache(long before) { CLibrary.trySkipCache(channel.getFileDescriptor(), 0, before, path()); diff --git a/src/java/org/apache/cassandra/schema/CompressionParams.java b/src/java/org/apache/cassandra/schema/CompressionParams.java index 7f467181d1..cd1686f5a9 100644 --- a/src/java/org/apache/cassandra/schema/CompressionParams.java +++ b/src/java/org/apache/cassandra/schema/CompressionParams.java @@ -71,6 +71,8 @@ public final class CompressionParams private final Integer chunkLength; private final ImmutableMap otherOptions; // Unrecognized options, can be used by the compressor + private volatile double crcCheckChance = 1.0; + public static CompressionParams fromMap(Map opts) { Map options = copyOptions(opts); @@ -455,6 +457,16 @@ public final class CompressionParams return String.valueOf(chunkLength() / 1024); } + public void setCrcCheckChance(double crcCheckChance) + { + this.crcCheckChance = crcCheckChance; + } + + public double getCrcCheckChance() + { + return crcCheckChance; + } + @Override public boolean equals(Object obj) { diff --git a/test/unit/org/apache/cassandra/cql3/validation/miscellaneous/CrcCheckChanceTest.java b/test/unit/org/apache/cassandra/cql3/validation/miscellaneous/CrcCheckChanceTest.java index 3a68e4a329..d059f7d16b 100644 --- a/test/unit/org/apache/cassandra/cql3/validation/miscellaneous/CrcCheckChanceTest.java +++ b/test/unit/org/apache/cassandra/cql3/validation/miscellaneous/CrcCheckChanceTest.java @@ -30,6 +30,8 @@ import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.compaction.CompactionInterruptedException; import org.apache.cassandra.db.compaction.CompactionManager; +import org.apache.cassandra.io.compress.CompressedRandomAccessReader; +import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.utils.FBUtilities; @@ -70,6 +72,7 @@ public class CrcCheckChanceTest extends CQLTester Assert.assertEquals(0.99, cfs.getCrcCheckChance()); Assert.assertEquals(0.99, cfs.getLiveSSTables().iterator().next().getCrcCheckChance()); + Assert.assertEquals(0.99, indexCfs.getCrcCheckChance()); Assert.assertEquals(0.99, indexCfs.getLiveSSTables().iterator().next().getCrcCheckChance()); @@ -145,8 +148,22 @@ public class CrcCheckChanceTest extends CQLTester Assert.assertEquals(0.03, cfs.getLiveSSTables().iterator().next().getCrcCheckChance()); Assert.assertEquals(0.03, indexCfs.getCrcCheckChance()); Assert.assertEquals(0.03, indexCfs.getLiveSSTables().iterator().next().getCrcCheckChance()); - } + // Also check that any open readers also use the updated value + // note: only compressed files currently perform crc checks, so only the dfile reader is relevant here + SSTableReader baseSSTable = cfs.getLiveSSTables().iterator().next(); + SSTableReader idxSSTable = indexCfs.getLiveSSTables().iterator().next(); + try (CompressedRandomAccessReader baseDataReader = (CompressedRandomAccessReader)baseSSTable.openDataReader(); + CompressedRandomAccessReader idxDataReader = (CompressedRandomAccessReader)idxSSTable.openDataReader()) + { + Assert.assertEquals(0.03, baseDataReader.getCrcCheckChance()); + Assert.assertEquals(0.03, idxDataReader.getCrcCheckChance()); + + cfs.setCrcCheckChance(0.31); + Assert.assertEquals(0.31, baseDataReader.getCrcCheckChance()); + Assert.assertEquals(0.31, idxDataReader.getCrcCheckChance()); + } + } @Test public void testDropDuringCompaction() throws Throwable