diff --git a/CHANGES.txt b/CHANGES.txt
index 095a90d64d19..0ae7fed53035 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,4 +1,5 @@
6.0-alpha3
+ * Give each compressed scan reader its own read-ahead buffer and keep the chunk cache for partition reads while one-shot scans bypass it (CASSANDRA-21671)
* Report the number of rows replica filtering protection cached when it warns (CASSANDRA-21623)
* Fix non-printable characters in Gossiper log for TOKENS (CASSANDRA-21417)
* Fix compression dictionary training failing with "insufficient samples" on large-chunk tables (CASSANDRA-21666)
diff --git a/src/java/org/apache/cassandra/cache/ChunkCache.java b/src/java/org/apache/cassandra/cache/ChunkCache.java
index 66782ac48bae..c8c966b7d246 100644
--- a/src/java/org/apache/cassandra/cache/ChunkCache.java
+++ b/src/java/org/apache/cassandra/cache/ChunkCache.java
@@ -37,6 +37,7 @@
import org.apache.cassandra.io.util.ChannelProxy;
import org.apache.cassandra.io.util.ChunkReader;
import org.apache.cassandra.io.util.FileHandle;
+import org.apache.cassandra.io.util.ReadPattern;
import org.apache.cassandra.io.util.Rebufferer;
import org.apache.cassandra.io.util.RebuffererFactory;
import org.apache.cassandra.metrics.ChunkCacheMetrics;
@@ -258,9 +259,15 @@ public void invalidate(long position)
}
@Override
- public Rebufferer instantiateRebufferer(boolean isScan)
+ public Rebufferer instantiateRebufferer(ReadPattern pattern)
{
- return this;
+ // A SCAN (compaction, cursor compaction) reads each chunk once and must not pollute the cache with
+ // one-shot chunks that evict hot data. It does not use the cache, so delegate to the source: the scan
+ // bypasses the cache and uses its own read-ahead buffer instead. For an Mmap source this also bypasses
+ // the chunk cache, which is correct: mmap data already lives in the OS page cache, so the chunk cache
+ // only duplicates it. A PARTITION_READ still uses the cache, because a repeated partition-range query
+ // re-reads hot data. See CASSANDRA-21671.
+ return pattern.usesCache() ? this : source.instantiateRebufferer(pattern);
}
@Override
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 ab3171200f17..284e3f05cfac 100644
--- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java
+++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java
@@ -97,6 +97,7 @@
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.FileUtils.DuplicateHardlinkException;
import org.apache.cassandra.io.util.RandomAccessReader;
+import org.apache.cassandra.io.util.ReadPattern;
import org.apache.cassandra.metrics.RestorableMeter;
import org.apache.cassandra.schema.SchemaConstants;
import org.apache.cassandra.schema.TableMetadataRef;
@@ -1421,43 +1422,51 @@ public StatsMetadata getSSTableMetadata()
public RandomAccessReader openDataReader()
{
- return openDataReaderInternal(null, null, false);
+ return openDataReaderInternal(null, null, ReadPattern.ROW_READ);
}
public RandomAccessReader openDataReader(RateLimiter limiter)
{
assert limiter != null;
- return openDataReaderInternal(null, limiter, false);
+ return openDataReaderInternal(null, limiter, ReadPattern.ROW_READ);
}
public RandomAccessReader openDataReader(DiskAccessMode diskAccessMode)
{
- return openDataReaderInternal(diskAccessMode, null, false);
+ return openDataReaderInternal(diskAccessMode, null, ReadPattern.ROW_READ);
}
- public RandomAccessReader openDataReaderForScan()
+ /**
+ * A reader for a query that walks a range of partitions, such as a token-range query. It reads in order, but
+ * keeps the chunk cache because a repeated partition-range query re-reads hot data. See CASSANDRA-21671.
+ */
+ public RandomAccessReader openDataReaderForPartitionRead()
{
- return openDataReaderInternal(null, null, true);
+ return openDataReaderInternal(null, null, ReadPattern.PARTITION_READ);
}
+ /**
+ * A reader for a one-shot scan (compaction and similar). It reads each chunk once, so it bypasses the chunk
+ * cache and uses its own read-ahead buffer. See CASSANDRA-21671.
+ */
public RandomAccessReader openDataReaderForScan(DiskAccessMode diskAccessMode)
{
- return openDataReaderInternal(diskAccessMode, null, true);
+ return openDataReaderInternal(diskAccessMode, null, ReadPattern.SCAN);
}
private RandomAccessReader openDataReaderInternal(@Nullable DiskAccessMode diskAccessMode,
@Nullable RateLimiter limiter,
- boolean forScan)
+ ReadPattern pattern)
{
if (canReuseDfile(diskAccessMode))
- return dfile.createReader(limiter, forScan, OnReaderClose.RETAIN_FILE_OPEN);
+ return dfile.createReader(limiter, pattern, OnReaderClose.RETAIN_FILE_OPEN);
FileHandle handle = dfile.toBuilder()
.withDiskAccessMode(diskAccessMode)
.complete();
try
{
- return handle.createReader(limiter, forScan, OnReaderClose.CLOSE_FILE);
+ return handle.createReader(limiter, pattern, OnReaderClose.CLOSE_FILE);
}
catch (Throwable t)
{
diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableScanner.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableScanner.java
index 0b54fec2c5da..9089b9b8bc21 100644
--- a/src/java/org/apache/cassandra/io/sstable/format/SSTableScanner.java
+++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableScanner.java
@@ -75,7 +75,7 @@ protected SSTableScanner(S sstable,
{
assert sstable != null;
- this.dfile = sstable.openDataReaderForScan();
+ this.dfile = sstable.openDataReaderForPartitionRead();
this.sstable = sstable;
this.columns = columns;
this.dataRange = dataRange;
diff --git a/src/java/org/apache/cassandra/io/util/CompressedChunkReader.java b/src/java/org/apache/cassandra/io/util/CompressedChunkReader.java
index 034fbb71532f..45dbae729e42 100644
--- a/src/java/org/apache/cassandra/io/util/CompressedChunkReader.java
+++ b/src/java/org/apache/cassandra/io/util/CompressedChunkReader.java
@@ -39,6 +39,10 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
final CompressionMetadata metadata;
final int maxCompressedLength;
final DoubleSupplier crcCheckChanceSupplier;
+ // Read-ahead is on only when a scan buffer is configured and larger than one chunk; a smaller buffer cannot
+ // batch reads, so it adds no value. A value of 0 means "no read-ahead". A per-scan view (see forScan) never
+ // reads ahead itself, so it always reports 0.
+ final int readAheadBufferSize;
protected CompressedChunkReader(ChannelProxy channel, CompressionMetadata metadata, DoubleSupplier crcCheckChanceSupplier)
{
@@ -46,9 +50,22 @@ protected CompressedChunkReader(ChannelProxy channel, CompressionMetadata metada
this.metadata = metadata;
this.maxCompressedLength = metadata.maxCompressedLength();
this.crcCheckChanceSupplier = crcCheckChanceSupplier;
+ int size = DatabaseDescriptor.getCompressedReadAheadBufferSize();
+ this.readAheadBufferSize = (size > 0 && size > metadata.chunkLength()) ? size : 0;
assert Integer.bitCount(metadata.chunkLength()) == 1; //must be a power of two
}
+ // Copy constructor for a per-scan view. The view shares the parent's channel and metadata but never reads
+ // ahead itself, so its readAheadBufferSize is 0.
+ protected CompressedChunkReader(CompressedChunkReader parent)
+ {
+ super(parent.channel, parent.metadata.dataLength);
+ this.metadata = parent.metadata;
+ this.maxCompressedLength = parent.maxCompressedLength;
+ this.crcCheckChanceSupplier = parent.crcCheckChanceSupplier;
+ this.readAheadBufferSize = 0;
+ }
+
protected CompressedChunkReader forScan()
{
return this;
@@ -90,9 +107,11 @@ public BufferType preferredBufferType()
}
@Override
- public Rebufferer instantiateRebufferer(boolean isScan)
+ public Rebufferer instantiateRebufferer(ReadPattern pattern)
{
- return new BufferManagingRebufferer.Aligned(isScan ? forScan() : this);
+ // A read-ahead pattern (PARTITION_READ, SCAN) gets a per-scan view that owns its own read-ahead buffer.
+ // A ROW_READ reads through this shared reader with no read-ahead.
+ return new BufferManagingRebufferer.Aligned(pattern.readsAhead() ? forScan() : this);
}
protected interface CompressedReader extends Closeable
@@ -206,10 +225,10 @@ private static class ScanCompressedReader implements CompressedReader
private final ChannelProxy channel;
private final ByteBufferHolder bufferHolder;
- private final ThreadLocalReadAheadBuffer readAheadBuffer;
+ private final ReadAheadBuffer readAheadBuffer;
private ScanCompressedReader(ChannelProxy channel, ByteBufferHolder bufferHolder,
- ThreadLocalReadAheadBuffer readAheadBuffer)
+ ReadAheadBuffer readAheadBuffer)
{
this.channel = channel;
this.bufferHolder = bufferHolder;
@@ -278,19 +297,25 @@ public static class Direct extends CompressedChunkReader
private final CompressedReader reader;
private final CompressedReader scanReader;
+ private final int blockSize;
public Direct(ChannelProxy channel, CompressionMetadata metadata, DoubleSupplier crcCheckChanceSupplier)
{
super(channel, metadata, crcCheckChanceSupplier);
- int blockSize = FileUtils.getFileBlockSize(channel.file());
+ this.blockSize = FileUtils.getFileBlockSize(channel.file());
this.reader = new DirectRandomAccessReader(channel, blockSize);
+ this.scanReader = null;
+ }
- int readAheadBufferSize = DatabaseDescriptor.getCompressedReadAheadBufferSize();
- this.scanReader = (readAheadBufferSize > 0 && readAheadBufferSize > metadata.chunkLength())
- ? new ScanCompressedReader(channel,
- new DirectThreadLocalByteBufferHolder(blockSize),
- new DirectThreadLocalReadAheadBuffer(channel, readAheadBufferSize, blockSize))
- : null;
+ // Per-scan view. Each scan reader is single-threaded and owns its own read-ahead buffer, so no buffer is
+ // shared across threads. It shares the parent's random-access reader as a fallback; that reader's close() is
+ // a no-op, so the view frees only its own scan buffer.
+ private Direct(Direct parent, CompressedReader scanReader)
+ {
+ super(parent);
+ this.blockSize = parent.blockSize;
+ this.reader = parent.reader;
+ this.scanReader = scanReader;
}
@Override
@@ -337,10 +362,14 @@ public void readChunk(long position, ByteBuffer uncompressed)
@Override
protected CompressedChunkReader forScan()
{
- if (scanReader != null)
- scanReader.allocateResources();
-
- return this;
+ if (readAheadBufferSize == 0)
+ return this;
+
+ ScanCompressedReader scan = new ScanCompressedReader(channel,
+ new DirectThreadLocalByteBufferHolder(blockSize),
+ new DirectReadAheadBuffer(channel, readAheadBufferSize, blockSize));
+ scan.allocateResources();
+ return new Direct(this, scan);
}
@Override
@@ -371,21 +400,28 @@ public Standard(ChannelProxy channel, CompressionMetadata metadata, DoubleSuppli
{
super(channel, metadata, crcCheckChanceSupplier);
reader = new RandomAccessCompressedReader(channel, metadata);
+ this.scanReader = null;
+ }
- int readAheadBufferSize = DatabaseDescriptor.getCompressedReadAheadBufferSize();
- scanReader = (readAheadBufferSize > 0 && readAheadBufferSize > metadata.chunkLength())
- ? new ScanCompressedReader(channel,
- new ThreadLocalByteBufferHolder(metadata.compressor().preferredBufferType()),
- new ThreadLocalReadAheadBuffer(channel, readAheadBufferSize, metadata.compressor().preferredBufferType())) : null;
+ // Per-scan view; see Direct for the ownership rationale.
+ private Standard(Standard parent, CompressedReader scanReader)
+ {
+ super(parent);
+ this.reader = parent.reader;
+ this.scanReader = scanReader;
}
@Override
protected CompressedChunkReader forScan()
{
- if (scanReader != null)
- scanReader.allocateResources();
-
- return this;
+ if (readAheadBufferSize == 0)
+ return this;
+
+ ScanCompressedReader scan = new ScanCompressedReader(channel,
+ new ThreadLocalByteBufferHolder(metadata.compressor().preferredBufferType()),
+ new ReadAheadBuffer(channel, readAheadBufferSize, metadata.compressor().preferredBufferType()));
+ scan.allocateResources();
+ return new Standard(this, scan);
}
@Override
diff --git a/src/java/org/apache/cassandra/io/util/DirectThreadLocalReadAheadBuffer.java b/src/java/org/apache/cassandra/io/util/DirectReadAheadBuffer.java
similarity index 83%
rename from src/java/org/apache/cassandra/io/util/DirectThreadLocalReadAheadBuffer.java
rename to src/java/org/apache/cassandra/io/util/DirectReadAheadBuffer.java
index 934e3620e114..5c745de020f4 100644
--- a/src/java/org/apache/cassandra/io/util/DirectThreadLocalReadAheadBuffer.java
+++ b/src/java/org/apache/cassandra/io/util/DirectReadAheadBuffer.java
@@ -28,12 +28,12 @@
import sun.nio.ch.DirectBuffer;
-public final class DirectThreadLocalReadAheadBuffer extends ThreadLocalReadAheadBuffer
+public final class DirectReadAheadBuffer extends ReadAheadBuffer
{
private final int blockSize;
- public DirectThreadLocalReadAheadBuffer(ChannelProxy channel, int bufferSize, int blockSize)
+ public DirectReadAheadBuffer(ChannelProxy channel, int bufferSize, int blockSize)
{
super(channel, () -> BufferUtil.allocateDirectAligned(BitUtil.align(bufferSize, blockSize), blockSize));
this.blockSize = blockSize;
@@ -53,8 +53,8 @@ protected void loadBlock(ByteBuffer blockBuffer, long blockPosition, int sizeToR
@Override
protected void cleanBuffer(ByteBuffer buffer)
{
- // Aligned buffers from BufferUtil.allocateDirectAligned are slices; clean the backing buffer (attachment)
+ // BufferUtil.allocateDirectAligned returns an aligned slice with no cleaner; free the backing
+ // allocation through the attachment, matching DirectThreadLocalByteBufferHolder.
MemoryUtil.clean((ByteBuffer) ((DirectBuffer) buffer).attachment());
}
-
-}
\ No newline at end of file
+}
diff --git a/src/java/org/apache/cassandra/io/util/EmptyRebufferer.java b/src/java/org/apache/cassandra/io/util/EmptyRebufferer.java
index 7f54a6b180f2..0d3f11b33aa4 100644
--- a/src/java/org/apache/cassandra/io/util/EmptyRebufferer.java
+++ b/src/java/org/apache/cassandra/io/util/EmptyRebufferer.java
@@ -64,7 +64,7 @@ public void closeReader()
}
@Override
- public Rebufferer instantiateRebufferer(boolean isScan)
+ public Rebufferer instantiateRebufferer(ReadPattern pattern)
{
return this;
}
diff --git a/src/java/org/apache/cassandra/io/util/FileHandle.java b/src/java/org/apache/cassandra/io/util/FileHandle.java
index 7f8c913b4b0a..79cbfc2aa510 100644
--- a/src/java/org/apache/cassandra/io/util/FileHandle.java
+++ b/src/java/org/apache/cassandra/io/util/FileHandle.java
@@ -205,23 +205,23 @@ public RandomAccessReader createReader()
*/
public RandomAccessReader createReader(RateLimiter limiter)
{
- return createReader(limiter, false);
+ return createReader(limiter, ReadPattern.ROW_READ);
}
- public RandomAccessReader createReader(RateLimiter limiter, boolean forScan)
+ public RandomAccessReader createReader(RateLimiter limiter, ReadPattern pattern)
{
- return createReader(limiter, forScan, OnReaderClose.RETAIN_FILE_OPEN);
+ return createReader(limiter, pattern, OnReaderClose.RETAIN_FILE_OPEN);
}
- public RandomAccessReader createReader(RateLimiter limiter, boolean forScan, OnReaderClose onReaderClose)
+ public RandomAccessReader createReader(RateLimiter limiter, ReadPattern pattern, OnReaderClose onReaderClose)
{
if (onReaderClose == OnReaderClose.CLOSE_FILE)
{
- return new RandomAccessReader.RandomAccessReaderWithOwnFile(instantiateRebufferer(limiter, forScan), this);
+ return new RandomAccessReader.RandomAccessReaderWithOwnFile(instantiateRebufferer(limiter, pattern), this);
}
else if (onReaderClose == OnReaderClose.RETAIN_FILE_OPEN)
{
- return new RandomAccessReader(instantiateRebufferer(limiter, forScan));
+ return new RandomAccessReader(instantiateRebufferer(limiter, pattern));
}
throw new IllegalArgumentException("Unknown close policy: " + onReaderClose);
}
@@ -266,12 +266,12 @@ public void dropPageCache(long before)
public Rebufferer instantiateRebufferer(RateLimiter limiter)
{
- return instantiateRebufferer(limiter, false);
+ return instantiateRebufferer(limiter, ReadPattern.ROW_READ);
}
- public Rebufferer instantiateRebufferer(RateLimiter limiter, boolean forScan)
+ public Rebufferer instantiateRebufferer(RateLimiter limiter, ReadPattern pattern)
{
- Rebufferer rebufferer = rebuffererFactory.instantiateRebufferer(forScan);
+ Rebufferer rebufferer = rebuffererFactory.instantiateRebufferer(pattern);
if (limiter != null)
rebufferer = new LimitingRebufferer(rebufferer, limiter, DiskOptimizationStrategy.MAX_BUFFER_SIZE);
diff --git a/src/java/org/apache/cassandra/io/util/MmapRebufferer.java b/src/java/org/apache/cassandra/io/util/MmapRebufferer.java
index 884bc9718642..b653c86c98d8 100644
--- a/src/java/org/apache/cassandra/io/util/MmapRebufferer.java
+++ b/src/java/org/apache/cassandra/io/util/MmapRebufferer.java
@@ -41,8 +41,10 @@ public BufferHolder rebuffer(long position)
}
@Override
- public Rebufferer instantiateRebufferer(boolean isScan)
+ public Rebufferer instantiateRebufferer(ReadPattern pattern)
{
+ // Mmap serves reads from the OS page cache regardless of the pattern, so there is no chunk cache to
+ // bypass and no separate read-ahead buffer to allocate.
return this;
}
diff --git a/src/java/org/apache/cassandra/io/util/RandomAccessReader.java b/src/java/org/apache/cassandra/io/util/RandomAccessReader.java
index 44b2f0c2fdd5..0bedceba4c33 100644
--- a/src/java/org/apache/cassandra/io/util/RandomAccessReader.java
+++ b/src/java/org/apache/cassandra/io/util/RandomAccessReader.java
@@ -385,7 +385,7 @@ public static RandomAccessReader open(File file)
try
{
ChunkReader reader = new SimpleChunkReader(channel, -1, BufferType.OFF_HEAP, DEFAULT_BUFFER_SIZE);
- Rebufferer rebufferer = reader.instantiateRebufferer(false);
+ Rebufferer rebufferer = reader.instantiateRebufferer(ReadPattern.ROW_READ);
return new RandomAccessReaderWithOwnChannel(rebufferer);
}
catch (Throwable t)
diff --git a/src/java/org/apache/cassandra/io/util/ThreadLocalReadAheadBuffer.java b/src/java/org/apache/cassandra/io/util/ReadAheadBuffer.java
similarity index 63%
rename from src/java/org/apache/cassandra/io/util/ThreadLocalReadAheadBuffer.java
rename to src/java/org/apache/cassandra/io/util/ReadAheadBuffer.java
index ff59f7cf96df..62d521ef8960 100644
--- a/src/java/org/apache/cassandra/io/util/ThreadLocalReadAheadBuffer.java
+++ b/src/java/org/apache/cassandra/io/util/ReadAheadBuffer.java
@@ -19,49 +19,40 @@
package org.apache.cassandra.io.util;
import java.nio.ByteBuffer;
-import java.util.HashMap;
-import java.util.Map;
import java.util.function.Supplier;
+import com.google.common.annotations.VisibleForTesting;
+
import org.apache.cassandra.io.compress.BufferType;
import org.apache.cassandra.io.compress.CorruptBlockException;
import org.apache.cassandra.io.sstable.CorruptSSTableException;
import org.apache.cassandra.utils.Closeable;
import org.apache.cassandra.utils.memory.MemoryUtil;
-import io.netty.util.concurrent.FastThreadLocal;
-
-public class ThreadLocalReadAheadBuffer implements Closeable
+/**
+ * A read-ahead buffer for sequential scans of a single file.
+ *
+ * Each instance owns its buffer. An instance is used by one scan reader, which is single-threaded, so the buffer is
+ * never shared across threads. A scan reader allocates one of these on open and frees it on close. N scanners over N
+ * inputs give N buffers by construction.
+ */
+public class ReadAheadBuffer implements Closeable
{
-
- private static class Block
- {
- ByteBuffer buffer = null;
- int index = -1;
- }
-
protected final ChannelProxy channel;
private final Supplier bufferSupplier;
-
- private static final FastThreadLocal