mirror of https://github.com/apache/cassandra
Add read ahead buffer for scans of compressed data files
Patch by Jordan West, Jon Haddad; Reviewed by David Capwell, Caleb Rackliffe, Dmitry Konstantinov for CASSANDRA-15452
This commit is contained in:
parent
4edf79a263
commit
a09d5240ab
|
|
@ -251,7 +251,7 @@ public class ChunkCache
|
|||
}
|
||||
|
||||
@Override
|
||||
public Rebufferer instantiateRebufferer()
|
||||
public Rebufferer instantiateRebufferer(boolean isScan)
|
||||
{
|
||||
return this;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -337,6 +337,8 @@ public class Config
|
|||
@Replaces(oldName = "min_free_space_per_drive_in_mb", converter = Converters.MEBIBYTES_DATA_STORAGE_INT, deprecated = true)
|
||||
public DataStorageSpec.IntMebibytesBound min_free_space_per_drive = new DataStorageSpec.IntMebibytesBound("50MiB");
|
||||
|
||||
public DataStorageSpec.IntKibibytesBound compressed_read_ahead_buffer_size = new DataStorageSpec.IntKibibytesBound("256KiB");
|
||||
|
||||
// fraction of free disk space available for compaction after min free space is subtracted
|
||||
public volatile Double max_space_usable_for_compactions_in_percentage = .95;
|
||||
|
||||
|
|
|
|||
|
|
@ -954,6 +954,9 @@ public class DatabaseDescriptor
|
|||
break;
|
||||
}
|
||||
|
||||
if (conf.compressed_read_ahead_buffer_size.toKibibytes() > 0 && conf.compressed_read_ahead_buffer_size.toKibibytes() < 256)
|
||||
throw new ConfigurationException("compressed_read_ahead_buffer_size must be at least 256KiB (set to 0 to disable), but was " + conf.compressed_read_ahead_buffer_size, false);
|
||||
|
||||
if (conf.server_encryption_options != null)
|
||||
{
|
||||
conf.server_encryption_options.applyConfig();
|
||||
|
|
@ -2534,6 +2537,24 @@ public class DatabaseDescriptor
|
|||
conf.concurrent_materialized_view_builders = value;
|
||||
}
|
||||
|
||||
public static int getCompressedReadAheadBufferSize()
|
||||
{
|
||||
return conf.compressed_read_ahead_buffer_size.toBytes();
|
||||
}
|
||||
|
||||
public static int getCompressedReadAheadBufferSizeInKB()
|
||||
{
|
||||
return conf.compressed_read_ahead_buffer_size.toKibibytes();
|
||||
}
|
||||
|
||||
public static void setCompressedReadAheadBufferSizeInKb(int sizeInKb)
|
||||
{
|
||||
if (sizeInKb < 256)
|
||||
throw new IllegalArgumentException("compressed_read_ahead_buffer_size_in_kb must be at least 256KiB");
|
||||
|
||||
conf.compressed_read_ahead_buffer_size = createIntKibibyteBoundAndEnsureItIsValidForByteConversion(sizeInKb, "compressed_read_ahead_buffer_size");
|
||||
}
|
||||
|
||||
public static long getMinFreeSpacePerDriveInMebibytes()
|
||||
{
|
||||
return conf.min_free_space_per_drive.toMebibytes();
|
||||
|
|
|
|||
|
|
@ -1266,6 +1266,11 @@ public abstract class SSTableReader extends SSTable implements UnfilteredSource,
|
|||
return dfile.createReader();
|
||||
}
|
||||
|
||||
public RandomAccessReader openDataReaderForScan()
|
||||
{
|
||||
return dfile.createReaderForScan();
|
||||
}
|
||||
|
||||
public void trySkipFileCacheBefore(DecoratedKey key)
|
||||
{
|
||||
long position = getPosition(key, SSTableReader.Operator.GE);
|
||||
|
|
|
|||
|
|
@ -78,7 +78,7 @@ implements ISSTableScanner
|
|||
{
|
||||
assert sstable != null;
|
||||
|
||||
this.dfile = sstable.openDataReader();
|
||||
this.dfile = sstable.openDataReaderForScan();
|
||||
this.sstable = sstable;
|
||||
this.columns = columns;
|
||||
this.dataRange = dataRange;
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ public abstract class BufferManagingRebufferer implements Rebufferer, Rebufferer
|
|||
public void closeReader()
|
||||
{
|
||||
BufferPools.forChunkCache().put(buffer);
|
||||
source.releaseUnderlyingResources();
|
||||
offset = -1;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -48,4 +48,6 @@ public interface ChunkReader extends RebuffererFactory
|
|||
* This is not guaranteed to be fulfilled.
|
||||
*/
|
||||
BufferType preferredBufferType();
|
||||
|
||||
default void releaseUnderlyingResources() {}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,11 +26,13 @@ import java.util.function.Supplier;
|
|||
import com.google.common.annotations.VisibleForTesting;
|
||||
import com.google.common.primitives.Ints;
|
||||
|
||||
import org.apache.cassandra.config.DatabaseDescriptor;
|
||||
import org.apache.cassandra.io.compress.BufferType;
|
||||
import org.apache.cassandra.io.compress.CompressionMetadata;
|
||||
import org.apache.cassandra.io.compress.CorruptBlockException;
|
||||
import org.apache.cassandra.io.sstable.CorruptSSTableException;
|
||||
import org.apache.cassandra.utils.ChecksumType;
|
||||
import org.apache.cassandra.utils.Closeable;
|
||||
|
||||
public abstract class CompressedChunkReader extends AbstractReaderFileProxy implements ChunkReader
|
||||
{
|
||||
|
|
@ -47,6 +49,11 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
|
|||
assert Integer.bitCount(metadata.chunkLength()) == 1; //must be a power of two
|
||||
}
|
||||
|
||||
protected CompressedChunkReader forScan()
|
||||
{
|
||||
return this;
|
||||
}
|
||||
|
||||
@VisibleForTesting
|
||||
public double getCrcCheckChance()
|
||||
{
|
||||
|
|
@ -83,20 +90,167 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
|
|||
}
|
||||
|
||||
@Override
|
||||
public Rebufferer instantiateRebufferer()
|
||||
public Rebufferer instantiateRebufferer(boolean isScan)
|
||||
{
|
||||
return new BufferManagingRebufferer.Aligned(this);
|
||||
return new BufferManagingRebufferer.Aligned(isScan ? forScan() : this);
|
||||
}
|
||||
|
||||
protected interface CompressedReader extends Closeable
|
||||
{
|
||||
default void allocateResources()
|
||||
{
|
||||
}
|
||||
|
||||
default void deallocateResources()
|
||||
{
|
||||
}
|
||||
|
||||
default boolean allocated()
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
default void close()
|
||||
{
|
||||
|
||||
}
|
||||
|
||||
|
||||
ByteBuffer read(CompressionMetadata.Chunk chunk, boolean shouldCheckCrc) throws CorruptBlockException;
|
||||
}
|
||||
|
||||
private static class RandomAccessCompressedReader implements CompressedReader
|
||||
{
|
||||
private final ChannelProxy channel;
|
||||
private final ThreadLocalByteBufferHolder bufferHolder;
|
||||
|
||||
private RandomAccessCompressedReader(ChannelProxy channel, CompressionMetadata metadata)
|
||||
{
|
||||
this.channel = channel;
|
||||
this.bufferHolder = new ThreadLocalByteBufferHolder(metadata.compressor().preferredBufferType());
|
||||
}
|
||||
|
||||
@Override
|
||||
public ByteBuffer read(CompressionMetadata.Chunk chunk, boolean shouldCheckCrc) throws CorruptBlockException
|
||||
{
|
||||
int length = shouldCheckCrc ? chunk.length + Integer.BYTES // compressed length + checksum length
|
||||
: chunk.length;
|
||||
ByteBuffer compressed = bufferHolder.getBuffer(length);
|
||||
if (channel.read(compressed, chunk.offset) != length)
|
||||
throw new CorruptBlockException(channel.filePath(), chunk);
|
||||
compressed.flip();
|
||||
compressed.limit(chunk.length);
|
||||
|
||||
if (shouldCheckCrc)
|
||||
{
|
||||
int checksum = (int) ChecksumType.CRC32.of(compressed);
|
||||
compressed.limit(length);
|
||||
if (compressed.getInt() != checksum)
|
||||
throw new CorruptBlockException(channel.filePath(), chunk);
|
||||
compressed.position(0).limit(chunk.length);
|
||||
}
|
||||
return compressed;
|
||||
}
|
||||
}
|
||||
|
||||
private static class ScanCompressedReader implements CompressedReader
|
||||
{
|
||||
private final ChannelProxy channel;
|
||||
private final ThreadLocalByteBufferHolder bufferHolder;
|
||||
private final ThreadLocalReadAheadBuffer readAheadBuffer;
|
||||
|
||||
private ScanCompressedReader(ChannelProxy channel, CompressionMetadata metadata, int readAheadBufferSize)
|
||||
{
|
||||
this.channel = channel;
|
||||
this.bufferHolder = new ThreadLocalByteBufferHolder(metadata.compressor().preferredBufferType());
|
||||
this.readAheadBuffer = new ThreadLocalReadAheadBuffer(channel, readAheadBufferSize, metadata.compressor().preferredBufferType());
|
||||
}
|
||||
|
||||
@Override
|
||||
public ByteBuffer read(CompressionMetadata.Chunk chunk, boolean shouldCheckCrc) throws CorruptBlockException
|
||||
{
|
||||
int length = shouldCheckCrc ? chunk.length + Integer.BYTES // compressed length + checksum length
|
||||
: chunk.length;
|
||||
ByteBuffer compressed = bufferHolder.getBuffer(length);
|
||||
|
||||
int copied = 0;
|
||||
while (copied < length)
|
||||
{
|
||||
readAheadBuffer.fill(chunk.offset + copied);
|
||||
int leftToRead = length - copied;
|
||||
if (readAheadBuffer.remaining() >= leftToRead)
|
||||
copied += readAheadBuffer.read(compressed, leftToRead);
|
||||
else
|
||||
copied += readAheadBuffer.read(compressed, readAheadBuffer.remaining());
|
||||
}
|
||||
|
||||
compressed.flip();
|
||||
compressed.limit(chunk.length);
|
||||
|
||||
if (shouldCheckCrc)
|
||||
{
|
||||
int checksum = (int) ChecksumType.CRC32.of(compressed);
|
||||
compressed.limit(length);
|
||||
if (compressed.getInt() != checksum)
|
||||
throw new CorruptBlockException(channel.filePath(), chunk);
|
||||
compressed.position(0).limit(chunk.length);
|
||||
}
|
||||
return compressed;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void allocateResources()
|
||||
{
|
||||
readAheadBuffer.allocateBuffer();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deallocateResources()
|
||||
{
|
||||
readAheadBuffer.clear(true);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean allocated()
|
||||
{
|
||||
return readAheadBuffer.hasBuffer();
|
||||
}
|
||||
|
||||
public void close()
|
||||
{
|
||||
readAheadBuffer.close();
|
||||
}
|
||||
}
|
||||
|
||||
public static class Standard extends CompressedChunkReader
|
||||
{
|
||||
// we read the raw compressed bytes into this buffer, then uncompressed them into the provided one.
|
||||
private final ThreadLocalByteBufferHolder bufferHolder;
|
||||
|
||||
private final CompressedReader reader;
|
||||
private final CompressedReader scanReader;
|
||||
|
||||
public Standard(ChannelProxy channel, CompressionMetadata metadata, Supplier<Double> crcCheckChanceSupplier)
|
||||
{
|
||||
super(channel, metadata, crcCheckChanceSupplier);
|
||||
bufferHolder = new ThreadLocalByteBufferHolder(metadata.compressor().preferredBufferType());
|
||||
reader = new RandomAccessCompressedReader(channel, metadata);
|
||||
|
||||
int readAheadBufferSize = DatabaseDescriptor.getCompressedReadAheadBufferSize();
|
||||
scanReader = (readAheadBufferSize > 0 && readAheadBufferSize > metadata.chunkLength())
|
||||
? new ScanCompressedReader(channel, metadata, readAheadBufferSize) : null;
|
||||
}
|
||||
|
||||
protected CompressedChunkReader forScan()
|
||||
{
|
||||
if (scanReader != null)
|
||||
scanReader.allocateResources();
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void releaseUnderlyingResources()
|
||||
{
|
||||
if (scanReader != null)
|
||||
scanReader.deallocateResources();
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -110,31 +264,13 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
|
|||
|
||||
CompressionMetadata.Chunk chunk = metadata.chunkFor(position);
|
||||
boolean shouldCheckCrc = shouldCheckCrc();
|
||||
int length = shouldCheckCrc ? chunk.length + Integer.BYTES // compressed length + checksum length
|
||||
: chunk.length;
|
||||
|
||||
CompressedReader readFrom = (scanReader != null && scanReader.allocated()) ? scanReader : reader;
|
||||
if (chunk.length < maxCompressedLength)
|
||||
{
|
||||
ByteBuffer compressed = bufferHolder.getBuffer(length);
|
||||
|
||||
if (channel.read(compressed, chunk.offset) != length)
|
||||
throw new CorruptBlockException(channel.filePath(), chunk);
|
||||
|
||||
compressed.flip();
|
||||
compressed.limit(chunk.length);
|
||||
ByteBuffer compressed = readFrom.read(chunk, shouldCheckCrc);
|
||||
uncompressed.clear();
|
||||
|
||||
if (shouldCheckCrc)
|
||||
{
|
||||
int checksum = (int) ChecksumType.CRC32.of(compressed);
|
||||
|
||||
compressed.limit(length);
|
||||
if (compressed.getInt() != checksum)
|
||||
throw new CorruptBlockException(channel.filePath(), chunk);
|
||||
|
||||
compressed.position(0).limit(chunk.length);
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
metadata.compressor().uncompress(compressed, uncompressed);
|
||||
|
|
@ -155,10 +291,9 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
|
|||
uncompressed.flip();
|
||||
int checksum = (int) ChecksumType.CRC32.of(uncompressed);
|
||||
|
||||
ByteBuffer scratch = bufferHolder.getBuffer(Integer.BYTES);
|
||||
|
||||
ByteBuffer scratch = ByteBuffer.allocate(Integer.BYTES);
|
||||
if (channel.read(scratch, chunk.offset + chunk.length) != Integer.BYTES
|
||||
|| scratch.getInt(0) != checksum)
|
||||
|| scratch.getInt(0) != checksum)
|
||||
throw new CorruptBlockException(channel.filePath(), chunk);
|
||||
}
|
||||
}
|
||||
|
|
@ -171,6 +306,16 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
|
|||
throw new CorruptSSTableException(e, channel.filePath());
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close()
|
||||
{
|
||||
reader.close();
|
||||
if (scanReader != null)
|
||||
scanReader.close();
|
||||
|
||||
super.close();
|
||||
}
|
||||
}
|
||||
|
||||
public static class Mmap extends CompressedChunkReader
|
||||
|
|
@ -233,7 +378,6 @@ public abstract class CompressedChunkReader extends AbstractReaderFileProxy impl
|
|||
uncompressed.position(0).limit(0);
|
||||
throw new CorruptSSTableException(e, channel.filePath());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
public void close()
|
||||
|
|
|
|||
|
|
@ -64,7 +64,7 @@ public class EmptyRebufferer implements Rebufferer, RebuffererFactory
|
|||
}
|
||||
|
||||
@Override
|
||||
public Rebufferer instantiateRebufferer()
|
||||
public Rebufferer instantiateRebufferer(boolean isScan)
|
||||
{
|
||||
return this;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -134,6 +134,11 @@ public class FileHandle extends SharedCloseableImpl
|
|||
return createReader(null);
|
||||
}
|
||||
|
||||
public RandomAccessReader createReaderForScan()
|
||||
{
|
||||
return createReader(null, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* Create {@link RandomAccessReader} with configured method of reading content of the file.
|
||||
* Reading from file will be rate limited by given {@link RateLimiter}.
|
||||
|
|
@ -143,7 +148,12 @@ public class FileHandle extends SharedCloseableImpl
|
|||
*/
|
||||
public RandomAccessReader createReader(RateLimiter limiter)
|
||||
{
|
||||
return new RandomAccessReader(instantiateRebufferer(limiter));
|
||||
return createReader(limiter, false);
|
||||
}
|
||||
|
||||
public RandomAccessReader createReader(RateLimiter limiter, boolean forScan)
|
||||
{
|
||||
return new RandomAccessReader(instantiateRebufferer(limiter, forScan));
|
||||
}
|
||||
|
||||
public FileDataInput createReader(long position)
|
||||
|
|
@ -186,7 +196,12 @@ public class FileHandle extends SharedCloseableImpl
|
|||
|
||||
public Rebufferer instantiateRebufferer(RateLimiter limiter)
|
||||
{
|
||||
Rebufferer rebufferer = rebuffererFactory.instantiateRebufferer();
|
||||
return instantiateRebufferer(limiter, false);
|
||||
}
|
||||
|
||||
public Rebufferer instantiateRebufferer(RateLimiter limiter, boolean forScan)
|
||||
{
|
||||
Rebufferer rebufferer = rebuffererFactory.instantiateRebufferer(forScan);
|
||||
|
||||
if (limiter != null)
|
||||
rebufferer = new LimitingRebufferer(rebufferer, limiter, DiskOptimizationStrategy.MAX_BUFFER_SIZE);
|
||||
|
|
|
|||
|
|
@ -41,7 +41,7 @@ class MmapRebufferer extends AbstractReaderFileProxy implements Rebufferer, Rebu
|
|||
}
|
||||
|
||||
@Override
|
||||
public Rebufferer instantiateRebufferer()
|
||||
public Rebufferer instantiateRebufferer(boolean isScan)
|
||||
{
|
||||
return this;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -332,7 +332,7 @@ public class RandomAccessReader extends RebufferingInputStream implements FileDa
|
|||
try
|
||||
{
|
||||
ChunkReader reader = new SimpleChunkReader(channel, -1, BufferType.OFF_HEAP, DEFAULT_BUFFER_SIZE);
|
||||
Rebufferer rebufferer = reader.instantiateRebufferer();
|
||||
Rebufferer rebufferer = reader.instantiateRebufferer(false);
|
||||
return new RandomAccessReaderWithOwnChannel(rebufferer);
|
||||
}
|
||||
catch (Throwable t)
|
||||
|
|
|
|||
|
|
@ -28,5 +28,5 @@ package org.apache.cassandra.io.util;
|
|||
*/
|
||||
public interface RebuffererFactory extends ReaderFileProxy
|
||||
{
|
||||
Rebufferer instantiateRebufferer();
|
||||
Rebufferer instantiateRebufferer(boolean isScan);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -55,7 +55,7 @@ class SimpleChunkReader extends AbstractReaderFileProxy implements ChunkReader
|
|||
}
|
||||
|
||||
@Override
|
||||
public Rebufferer instantiateRebufferer()
|
||||
public Rebufferer instantiateRebufferer(boolean forScan)
|
||||
{
|
||||
if (Integer.bitCount(bufferSize) == 1)
|
||||
return new BufferManagingRebufferer.Aligned(this);
|
||||
|
|
|
|||
|
|
@ -0,0 +1,154 @@
|
|||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one
|
||||
* or more contributor license agreements. See the NOTICE file
|
||||
* distributed with this work for additional information
|
||||
* regarding copyright ownership. The ASF licenses this file
|
||||
* to you under the Apache License, Version 2.0 (the
|
||||
* "License"); you may not use this file except in compliance
|
||||
* with the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.cassandra.io.util;
|
||||
|
||||
import java.nio.ByteBuffer;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import io.netty.util.concurrent.FastThreadLocal;
|
||||
import org.apache.cassandra.io.compress.BufferType;
|
||||
import org.apache.cassandra.io.sstable.CorruptSSTableException;
|
||||
|
||||
public final class ThreadLocalReadAheadBuffer
|
||||
{
|
||||
private static class Block
|
||||
{
|
||||
ByteBuffer buffer = null;
|
||||
int index = -1;
|
||||
}
|
||||
|
||||
private final ChannelProxy channel;
|
||||
|
||||
private final BufferType bufferType;
|
||||
|
||||
private static final FastThreadLocal<Map<String, Block>> blockMap = new FastThreadLocal<>()
|
||||
{
|
||||
@Override
|
||||
protected Map<String, Block> initialValue()
|
||||
{
|
||||
return new HashMap<>();
|
||||
}
|
||||
};
|
||||
|
||||
private final int bufferSize;
|
||||
private final long channelSize;
|
||||
|
||||
public ThreadLocalReadAheadBuffer(ChannelProxy channel, int bufferSize, BufferType bufferType)
|
||||
{
|
||||
this.channel = channel;
|
||||
this.channelSize = channel.size();
|
||||
this.bufferSize = bufferSize;
|
||||
this.bufferType = bufferType;
|
||||
}
|
||||
|
||||
public boolean hasBuffer()
|
||||
{
|
||||
return block().buffer != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Safe to call only if {@link #hasBuffer()} is true
|
||||
*/
|
||||
public int remaining()
|
||||
{
|
||||
return getBlock().buffer.remaining();
|
||||
}
|
||||
|
||||
public void allocateBuffer()
|
||||
{
|
||||
getBlock();
|
||||
}
|
||||
|
||||
private Block getBlock()
|
||||
{
|
||||
Block block = block();
|
||||
if (block.buffer == null)
|
||||
{
|
||||
block.buffer = bufferType.allocate(bufferSize);
|
||||
block.buffer.clear();
|
||||
}
|
||||
return block;
|
||||
}
|
||||
|
||||
private Block block()
|
||||
{
|
||||
return blockMap.get().computeIfAbsent(channel.filePath(), k -> new Block());
|
||||
}
|
||||
|
||||
public void fill(long position)
|
||||
{
|
||||
Block block = getBlock();
|
||||
ByteBuffer blockBuffer = block.buffer;
|
||||
long realPosition = Math.min(channelSize, position);
|
||||
int blockNo = (int) (realPosition / bufferSize);
|
||||
long blockPosition = blockNo * (long) bufferSize;
|
||||
|
||||
long remaining = channelSize - blockPosition;
|
||||
int sizeToRead = (int) Math.min(remaining, bufferSize);
|
||||
if (block.index != blockNo)
|
||||
{
|
||||
blockBuffer.flip();
|
||||
blockBuffer.limit(sizeToRead);
|
||||
if (channel.read(blockBuffer, blockPosition) != sizeToRead)
|
||||
throw new CorruptSSTableException(null, channel.filePath());
|
||||
|
||||
block.index = blockNo;
|
||||
}
|
||||
|
||||
blockBuffer.flip();
|
||||
blockBuffer.limit(sizeToRead);
|
||||
blockBuffer.position((int) (realPosition - blockPosition));
|
||||
}
|
||||
|
||||
public int read(ByteBuffer dest, int length)
|
||||
{
|
||||
Block block = getBlock();
|
||||
ByteBuffer blockBuffer = block.buffer;
|
||||
ByteBuffer tmp = blockBuffer.duplicate();
|
||||
tmp.limit(tmp.position() + length);
|
||||
dest.put(tmp);
|
||||
blockBuffer.position(blockBuffer.position() + length);
|
||||
|
||||
return length;
|
||||
}
|
||||
|
||||
public void clear(boolean deallocate)
|
||||
{
|
||||
Block block = getBlock();
|
||||
block.index = -1;
|
||||
|
||||
ByteBuffer blockBuffer = block.buffer;
|
||||
if (blockBuffer != null)
|
||||
{
|
||||
blockBuffer.clear();
|
||||
if (deallocate)
|
||||
{
|
||||
FileUtils.clean(blockBuffer);
|
||||
block.buffer = null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public void close()
|
||||
{
|
||||
clear(true);
|
||||
blockMap.get().remove(channel.filePath());
|
||||
}
|
||||
}
|
||||
|
|
@ -1893,6 +1893,20 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
|
|||
value, oldValue);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int getCompressedReadAheadBufferInKB()
|
||||
{
|
||||
return DatabaseDescriptor.getCompressedReadAheadBufferSizeInKB();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setCompressedReadAheadBufferInKB(int sizeInKb)
|
||||
{
|
||||
DatabaseDescriptor.setCompressedReadAheadBufferSizeInKb(sizeInKb);
|
||||
logger.info("set compressed read ahead buffer size to {}KiB", sizeInKb);
|
||||
}
|
||||
|
||||
|
||||
public int getBatchlogReplayThrottleInKB()
|
||||
{
|
||||
return DatabaseDescriptor.getBatchlogReplayThrottleInKiB();
|
||||
|
|
|
|||
|
|
@ -812,6 +812,9 @@ public interface StorageServiceMBean extends NotificationEmitter
|
|||
public int getCompactionThroughputMbPerSec();
|
||||
public void setCompactionThroughputMbPerSec(int value);
|
||||
|
||||
public int getCompressedReadAheadBufferInKB();
|
||||
public void setCompressedReadAheadBufferInKB(int sizeInKb);
|
||||
|
||||
public int getBatchlogReplayThrottleInKB();
|
||||
public void setBatchlogReplayThrottleInKB(int value);
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,141 @@
|
|||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one
|
||||
* or more contributor license agreements. See the NOTICE file
|
||||
* distributed with this work for additional information
|
||||
* regarding copyright ownership. The ASF licenses this file
|
||||
* to you under the Apache License, Version 2.0 (the
|
||||
* "License"); you may not use this file except in compliance
|
||||
* with the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.cassandra.io.util;
|
||||
|
||||
import accord.utils.Gen;
|
||||
import accord.utils.Gens;
|
||||
import org.apache.cassandra.config.DatabaseDescriptor;
|
||||
import org.apache.cassandra.db.ClusteringComparator;
|
||||
import org.apache.cassandra.io.compress.CompressedSequentialWriter;
|
||||
import org.apache.cassandra.io.compress.CompressionMetadata;
|
||||
import org.apache.cassandra.io.filesystem.ListenableFileSystem;
|
||||
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
|
||||
import org.apache.cassandra.schema.CompressionParams;
|
||||
import org.assertj.core.api.Assertions;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.file.Files;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static accord.utils.Property.qt;
|
||||
|
||||
public class CompressedChunkReaderTest
|
||||
{
|
||||
static
|
||||
{
|
||||
DatabaseDescriptor.clientInitialization();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void scanReaderReadsLessThanRAReader()
|
||||
{
|
||||
var optionGen = options();
|
||||
var paramsGen = params();
|
||||
var lengthGen = Gens.longs().between(1, 1 << 16);
|
||||
qt().withSeed(-1871070464864118891L).forAll(Gens.random(), optionGen, paramsGen).check((rs, option, params) -> {
|
||||
ListenableFileSystem fs = FileSystems.newGlobalInMemoryFileSystem();
|
||||
|
||||
File f = new File("/file.db");
|
||||
AtomicInteger reads = new AtomicInteger();
|
||||
fs.onPostRead(f.path::equals, (p, c, pos, dst, r) -> {
|
||||
reads.incrementAndGet();
|
||||
});
|
||||
long length = lengthGen.nextLong(rs);
|
||||
CompressionMetadata metadata1, metadata2;
|
||||
try (CompressedSequentialWriter writer = new CompressedSequentialWriter(f, new File("/file.offset"), new File("/file.digest"), option, params, new MetadataCollector(new ClusteringComparator())))
|
||||
{
|
||||
for (long i = 0; i < length; i++)
|
||||
writer.writeLong(i);
|
||||
|
||||
writer.sync();
|
||||
metadata1 = writer.open(0);
|
||||
metadata2 = writer.open(0);
|
||||
}
|
||||
|
||||
doReads(f, metadata1, length, true);
|
||||
int scanReads = reads.getAndSet(0);
|
||||
|
||||
doReads(f, metadata2, length, false);
|
||||
int raReads = reads.getAndSet(0);
|
||||
|
||||
if (Files.size(f.toPath()) > DatabaseDescriptor.getCompressedReadAheadBufferSize())
|
||||
Assert.assertTrue(scanReads < raReads);
|
||||
});
|
||||
}
|
||||
|
||||
private void doReads(File f, CompressionMetadata metadata, long length, boolean useReadAhead)
|
||||
{
|
||||
ByteBuffer buffer = ByteBuffer.allocateDirect(metadata.chunkLength());
|
||||
|
||||
try (ChannelProxy channel = new ChannelProxy(f);
|
||||
CompressedChunkReader reader = new CompressedChunkReader.Standard(channel, metadata, () -> 1.1);
|
||||
metadata)
|
||||
{
|
||||
if (useReadAhead)
|
||||
reader.forScan();
|
||||
|
||||
long offset = 0;
|
||||
long maxOffset = length * Long.BYTES;
|
||||
do
|
||||
{
|
||||
reader.readChunk(offset, buffer);
|
||||
for (long expected = offset / Long.BYTES; buffer.hasRemaining(); expected++)
|
||||
Assertions.assertThat(buffer.getLong()).isEqualTo(expected);
|
||||
|
||||
offset += metadata.chunkLength();
|
||||
}
|
||||
while (offset < maxOffset);
|
||||
}
|
||||
finally
|
||||
{
|
||||
FileUtils.clean(buffer);
|
||||
}}
|
||||
|
||||
private static Gen<SequentialWriterOption> options()
|
||||
{
|
||||
Gen<Integer> bufferSizes = Gens.constant(1 << 10); //.pickInt(1 << 4, 1 << 10, 1 << 15);
|
||||
return rs -> SequentialWriterOption.newBuilder()
|
||||
.finishOnClose(false)
|
||||
.bufferSize(bufferSizes.next(rs))
|
||||
.build();
|
||||
}
|
||||
|
||||
private enum CompressionKind { Noop, Snappy, Deflate, Lz4, Zstd }
|
||||
|
||||
private static Gen<CompressionParams> params()
|
||||
{
|
||||
Gen<Integer> chunkLengths = Gens.constant(CompressionParams.DEFAULT_CHUNK_LENGTH);
|
||||
Gen<Double> compressionRatio = Gens.pick(1.1D);
|
||||
return rs -> {
|
||||
CompressionKind kind = rs.pick(CompressionKind.values());
|
||||
switch (kind)
|
||||
{
|
||||
case Noop: return CompressionParams.noop();
|
||||
case Snappy: return CompressionParams.snappy(chunkLengths.next(rs), compressionRatio.next(rs));
|
||||
case Deflate: return CompressionParams.deflate(chunkLengths.next(rs));
|
||||
case Lz4: return CompressionParams.lz4(chunkLengths.next(rs));
|
||||
case Zstd: return CompressionParams.zstd(chunkLengths.next(rs));
|
||||
default: throw new UnsupportedOperationException(kind.name());
|
||||
}
|
||||
};
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,204 @@
|
|||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one
|
||||
* or more contributor license agreements. See the NOTICE file
|
||||
* distributed with this work for additional information
|
||||
* regarding copyright ownership. The ASF licenses this file
|
||||
* to you under the Apache License, Version 2.0 (the
|
||||
* "License"); you may not use this file except in compliance
|
||||
* with the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.cassandra.io.util;
|
||||
|
||||
|
||||
import java.io.FileOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.file.Files;
|
||||
import java.util.List;
|
||||
import java.util.Random;
|
||||
|
||||
import org.junit.AfterClass;
|
||||
import org.junit.Assert;
|
||||
import org.junit.BeforeClass;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import org.apache.cassandra.config.DataStorageSpec;
|
||||
import org.apache.cassandra.io.compress.BufferType;
|
||||
import org.apache.cassandra.io.sstable.CorruptSSTableException;
|
||||
import org.apache.cassandra.utils.Pair;
|
||||
import org.quicktheories.WithQuickTheories;
|
||||
import org.quicktheories.core.Gen;
|
||||
|
||||
import static org.apache.cassandra.config.CassandraRelevantProperties.JAVA_IO_TMPDIR;
|
||||
|
||||
public class ThreadLocalReadAheadBufferTest implements WithQuickTheories
|
||||
{
|
||||
private static final int numFiles = 5;
|
||||
private static final File[] files = new File[numFiles];
|
||||
private static final Logger logger = LoggerFactory.getLogger(ThreadLocalReadAheadBufferTest.class);
|
||||
|
||||
@BeforeClass
|
||||
public static void setup()
|
||||
{
|
||||
int seed = new Random().nextInt();
|
||||
logger.info("Seed: {}", seed);
|
||||
|
||||
for (int i = 0; i < numFiles; i++)
|
||||
{
|
||||
int size = new Random().nextInt((Integer.MAX_VALUE - 1) / 8);
|
||||
files[i] = writeFile(seed, size);
|
||||
}
|
||||
}
|
||||
|
||||
@AfterClass
|
||||
public static void cleanup()
|
||||
{
|
||||
for (File f : files)
|
||||
{
|
||||
try
|
||||
{
|
||||
f.delete();
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
// ignore
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testLastBlockReads()
|
||||
{
|
||||
qt().forAll(lastBlockReads())
|
||||
.checkAssert(this::testReads);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testReadsLikeChannelProxy()
|
||||
{
|
||||
|
||||
qt().forAll(randomReads())
|
||||
.checkAssert(this::testReads);
|
||||
}
|
||||
|
||||
private void testReads(InputData propertyInputs)
|
||||
{
|
||||
try (ChannelProxy channel = new ChannelProxy(propertyInputs.file))
|
||||
{
|
||||
ThreadLocalReadAheadBuffer trlab = new ThreadLocalReadAheadBuffer(channel, new DataStorageSpec.IntKibibytesBound("256KiB").toBytes(), BufferType.OFF_HEAP);
|
||||
for (Pair<Long, Integer> read : propertyInputs.positionsAndLengths)
|
||||
{
|
||||
int readSize = Math.min(read.right,(int) (channel.size() - read.left));
|
||||
ByteBuffer buf1 = ByteBuffer.allocate(readSize);
|
||||
channel.read(buf1, read.left);
|
||||
|
||||
ByteBuffer buf2 = ByteBuffer.allocate(readSize);
|
||||
try
|
||||
{
|
||||
int copied = 0;
|
||||
while (copied < readSize) {
|
||||
trlab.fill(read.left + copied);
|
||||
int leftToRead = readSize - copied;
|
||||
if (trlab.remaining() >= leftToRead)
|
||||
copied += trlab.read(buf2, leftToRead);
|
||||
else
|
||||
copied += trlab.read(buf2, trlab.remaining());
|
||||
}
|
||||
}
|
||||
catch (CorruptSSTableException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
|
||||
Assert.assertEquals(buf1, buf2);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private Gen<InputData> lastBlockReads()
|
||||
{
|
||||
return arbitrary().pick(List.of(files))
|
||||
.flatMap((file) ->
|
||||
lists().of(longs().between(0, fileSize(file)).zip(integers().between(1, 100), Pair::create))
|
||||
.ofSizeBetween(5, 10)
|
||||
.map(positionsAndLengths -> new InputData(file, positionsAndLengths)));
|
||||
|
||||
}
|
||||
|
||||
private Gen<InputData> randomReads()
|
||||
{
|
||||
int blockSize = new DataStorageSpec.IntKibibytesBound("256KiB").toBytes();
|
||||
return arbitrary().pick(List.of(files))
|
||||
.flatMap((file) ->
|
||||
lists().of(longs().between(fileSize(file) - blockSize, fileSize(file)).zip(integers().between(1, 100), Pair::create))
|
||||
.ofSizeBetween(5, 10)
|
||||
.map(positionsAndLengths -> new InputData(file, positionsAndLengths)));
|
||||
|
||||
}
|
||||
|
||||
// need this becasue generators don't handle the IOException
|
||||
private long fileSize(File file)
|
||||
{
|
||||
try
|
||||
{
|
||||
return Files.size(file.toPath());
|
||||
} catch (IOException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
private static class InputData
|
||||
{
|
||||
|
||||
private final File file;
|
||||
private final List<Pair<Long, Integer>> positionsAndLengths;
|
||||
|
||||
public InputData(File file, List<Pair<Long, Integer>> positionsAndLengths)
|
||||
{
|
||||
this.file = file;
|
||||
this.positionsAndLengths = positionsAndLengths;
|
||||
}
|
||||
}
|
||||
|
||||
private static File writeFile(int seed, int length)
|
||||
{
|
||||
String fileName = JAVA_IO_TMPDIR.getString() + "data+" + length + ".bin";
|
||||
|
||||
byte[] dataChunk = new byte[4096 * 8];
|
||||
java.util.Random random = new Random(seed);
|
||||
int writtenData = 0;
|
||||
|
||||
File file = new File(fileName);
|
||||
try (FileOutputStream fos = new FileOutputStream(file.toJavaIOFile()))
|
||||
{
|
||||
while (writtenData < length)
|
||||
{
|
||||
random.nextBytes(dataChunk);
|
||||
int toWrite = Math.min((length - writtenData), dataChunk.length);
|
||||
fos.write(dataChunk, 0, toWrite);
|
||||
writtenData += toWrite;
|
||||
}
|
||||
fos.flush();
|
||||
}
|
||||
catch (IOException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
|
||||
return file;
|
||||
}
|
||||
|
||||
}
|
||||
Loading…
Reference in New Issue