Enable Direct-IO feature for CommitLog files using Java native API's.

Patch by Amit Pawar and Jacek Lewandowski; reviewed by Branimir Lambov and Maxwell Guo for CASSANDRA-18464

Co-authored-by: Amit Pawar <Amit.Pawar@amd.com>
Co-authored-by: Jacek Lewandowski <lewandowski.jacek@gmail.com>
This commit is contained in:
Amit Pawar 2023-10-06 01:23:37 +05:30 committed by Jacek Lewandowski
parent 3cdf71defe
commit 3259bea533
35 changed files with 941 additions and 136 deletions

View File

@ -1,4 +1,5 @@
5.0-beta1 5.0-beta1
* Enable Direct-IO feature for CommitLog files using Java native API's. (CASSANDRA-18464)
* SAI fixes for composite partitions, and static and non-static rows intersections (CASSANDRA-19034) * SAI fixes for composite partitions, and static and non-static rows intersections (CASSANDRA-19034)
* Improve SAI IndexContext handling of indexed and non-indexed columns in queries (CASSANDRA-18166) * Improve SAI IndexContext handling of indexed and non-indexed columns in queries (CASSANDRA-18166)
* Fixed bug where UnifiedCompactionTask constructor was calling the wrong base constructor of CompactionTask (CASSANDRA-18757) * Fixed bug where UnifiedCompactionTask constructor was calling the wrong base constructor of CompactionTask (CASSANDRA-18757)

View File

@ -187,6 +187,9 @@ New features
- Added snitch for Microsoft Azure of name AzureSnitch (CASSANDRA-18646) - Added snitch for Microsoft Azure of name AzureSnitch (CASSANDRA-18646)
- legacy command line options from cassandra-stress were removed - legacy command line options from cassandra-stress were removed
- `-mode` option in cassandra-stress has `native` and `cql3` as defaults and they do not need to be specified - `-mode` option in cassandra-stress has `native` and `cql3` as defaults and they do not need to be specified
- Allow to write the commitlog using direct I/O. Direct I/O is a new feature that minimizes cache effects and
memory-mapping overhead by using user-space buffers. This helps in transferring data from/to disk at high speed.
Java enabled support for the direct I/O feature from version 10 onwards - see JDK-8164900 for reference. (CASSANDRA-18464)
Upgrading Upgrading
--------- ---------

View File

@ -606,6 +606,16 @@ commitlog_segment_size: 32MiB
# parameters: # parameters:
# - # -
# Set the disk access mode for writing commitlog segments. The allowed values are:
# - auto: version dependent optimal setting
# - legacy: the default mode as used in Cassandra 4.x and earlier (standard I/O when the commitlog is either
# compressed or encrypted or mmap otherwise)
# - mmap: use memory mapped I/O - available only when the commitlog is neither compressed nor encrypted
# - direct: use direct I/O - available only when the commitlog is neither compressed nor encrypted
# - standard: use standard I/O - available only when the commitlog is compressed or encrypted
# The default setting is legacy when the storage compatibility is set to 4 or auto otherwise.
commitlog_disk_access_mode: legacy
# Compression to apply to SSTables as they flush for compressed tables. # Compression to apply to SSTables as they flush for compressed tables.
# Note that tables without compression enabled do not respect this flag. # Note that tables without compression enabled do not respect this flag.
# #

View File

@ -386,6 +386,7 @@ public class Config
public ParameterizedClass commitlog_compression; public ParameterizedClass commitlog_compression;
public FlushCompression flush_compression = FlushCompression.fast; public FlushCompression flush_compression = FlushCompression.fast;
public int commitlog_max_compression_buffers_in_pool = 3; public int commitlog_max_compression_buffers_in_pool = 3;
public DiskAccessMode commitlog_disk_access_mode = DiskAccessMode.legacy;
@Replaces(oldName = "periodic_commitlog_sync_lag_block_in_ms", converter = Converters.MILLIS_DURATION_INT, deprecated = true) @Replaces(oldName = "periodic_commitlog_sync_lag_block_in_ms", converter = Converters.MILLIS_DURATION_INT, deprecated = true)
public DurationSpec.IntMillisecondsBound periodic_commitlog_sync_lag_block; public DurationSpec.IntMillisecondsBound periodic_commitlog_sync_lag_block;
public TransparentDataEncryptionOptions transparent_data_encryption_options = new TransparentDataEncryptionOptions(); public TransparentDataEncryptionOptions transparent_data_encryption_options = new TransparentDataEncryptionOptions();
@ -1145,6 +1146,8 @@ public class Config
mmap, mmap,
mmap_index_only, mmap_index_only,
standard, standard,
legacy,
direct // Direct-I/O is enabled for commitlog disk only.
} }
public enum MemtableAllocationType public enum MemtableAllocationType

View File

@ -73,6 +73,7 @@ import org.apache.cassandra.auth.IInternodeAuthenticator;
import org.apache.cassandra.auth.INetworkAuthorizer; import org.apache.cassandra.auth.INetworkAuthorizer;
import org.apache.cassandra.auth.IRoleManager; import org.apache.cassandra.auth.IRoleManager;
import org.apache.cassandra.config.Config.CommitLogSync; import org.apache.cassandra.config.Config.CommitLogSync;
import org.apache.cassandra.config.Config.DiskAccessMode;
import org.apache.cassandra.config.Config.PaxosOnLinearizabilityViolation; import org.apache.cassandra.config.Config.PaxosOnLinearizabilityViolation;
import org.apache.cassandra.config.Config.PaxosStatePurging; import org.apache.cassandra.config.Config.PaxosStatePurging;
import org.apache.cassandra.db.ConsistencyLevel; import org.apache.cassandra.db.ConsistencyLevel;
@ -176,7 +177,9 @@ public class DatabaseDescriptor
private static IPartitioner partitioner; private static IPartitioner partitioner;
private static String paritionerName; private static String paritionerName;
private static Config.DiskAccessMode indexAccessMode; private static DiskAccessMode indexAccessMode;
private static DiskAccessMode commitLogWriteDiskAccessMode;
private static AbstractCryptoProvider cryptoProvider; private static AbstractCryptoProvider cryptoProvider;
private static IAuthenticator authenticator; private static IAuthenticator authenticator;
@ -512,23 +515,25 @@ public class DatabaseDescriptor
} }
/* evaluate the DiskAccessMode Config directive, which also affects indexAccessMode selection */ /* evaluate the DiskAccessMode Config directive, which also affects indexAccessMode selection */
if (conf.disk_access_mode == Config.DiskAccessMode.auto) if (conf.disk_access_mode == DiskAccessMode.auto || conf.disk_access_mode == DiskAccessMode.mmap_index_only)
{ {
conf.disk_access_mode = hasLargeAddressSpace() ? Config.DiskAccessMode.mmap : Config.DiskAccessMode.standard; conf.disk_access_mode = DiskAccessMode.standard;
indexAccessMode = conf.disk_access_mode; indexAccessMode = DiskAccessMode.mmap;
logger.info("DiskAccessMode 'auto' determined to be {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode);
} }
else if (conf.disk_access_mode == Config.DiskAccessMode.mmap_index_only) else if (conf.disk_access_mode == DiskAccessMode.legacy)
{ {
conf.disk_access_mode = Config.DiskAccessMode.standard; conf.disk_access_mode = hasLargeAddressSpace() ? DiskAccessMode.mmap : DiskAccessMode.standard;
indexAccessMode = Config.DiskAccessMode.mmap; indexAccessMode = conf.disk_access_mode;
logger.info("DiskAccessMode is {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode); }
else if (conf.disk_access_mode == DiskAccessMode.direct)
{
throw new ConfigurationException(String.format("DiskAccessMode '%s' is not supported", DiskAccessMode.direct));
} }
else else
{ {
indexAccessMode = conf.disk_access_mode; indexAccessMode = conf.disk_access_mode;
logger.info("DiskAccessMode is {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode);
} }
logger.info("DiskAccessMode is {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode);
/* phi convict threshold for FailureDetector */ /* phi convict threshold for FailureDetector */
if (conf.phi_convict_threshold < 5 || conf.phi_convict_threshold > 16) if (conf.phi_convict_threshold < 5 || conf.phi_convict_threshold > 16)
@ -617,6 +622,10 @@ public class DatabaseDescriptor
conf.commitlog_directory = storagedirFor("commitlog"); conf.commitlog_directory = storagedirFor("commitlog");
} }
initializeCommitLogDiskAccessMode();
if (commitLogWriteDiskAccessMode != conf.commitlog_disk_access_mode)
logger.info("commitlog_disk_access_mode resolved to: {}", commitLogWriteDiskAccessMode);
if (conf.hints_directory == null) if (conf.hints_directory == null)
{ {
conf.hints_directory = storagedirFor("hints"); conf.hints_directory = storagedirFor("hints");
@ -1415,6 +1424,52 @@ public class DatabaseDescriptor
paritionerName = partitioner.getClass().getCanonicalName(); paritionerName = partitioner.getClass().getCanonicalName();
} }
private static DiskAccessMode resolveCommitLogWriteDiskAccessMode(DiskAccessMode providedDiskAccessMode)
{
boolean compressOrEncrypt = getCommitLogCompression() != null || (getEncryptionContext() != null && getEncryptionContext().isEnabled());
boolean directIOSupported = false;
try
{
directIOSupported = FileUtils.getBlockSize(new File(getCommitLogLocation())) > 0;
}
catch (RuntimeException e)
{
logger.warn("Unable to determine block size for commit log directory: {}", e.getMessage());
}
if (providedDiskAccessMode == DiskAccessMode.auto)
{
if (compressOrEncrypt)
providedDiskAccessMode = DiskAccessMode.legacy;
else
{
providedDiskAccessMode = directIOSupported && conf.disk_optimization_strategy == Config.DiskOptimizationStrategy.ssd ? DiskAccessMode.direct
: DiskAccessMode.legacy;
}
}
if (providedDiskAccessMode == DiskAccessMode.legacy)
{
providedDiskAccessMode = compressOrEncrypt ? DiskAccessMode.standard : DiskAccessMode.mmap;
}
return providedDiskAccessMode;
}
private static void validateCommitLogWriteDiskAccessMode(DiskAccessMode diskAccessMode) throws ConfigurationException
{
boolean compressOrEncrypt = getCommitLogCompression() != null || (getEncryptionContext() != null && getEncryptionContext().isEnabled());
if (compressOrEncrypt && diskAccessMode != DiskAccessMode.standard)
{
throw new ConfigurationException("commitlog_disk_access_mode = " + diskAccessMode + " is not supported with compression or encryption. Please use 'auto' when unsure.", false);
}
else if (!compressOrEncrypt && diskAccessMode != DiskAccessMode.mmap && diskAccessMode != DiskAccessMode.direct)
{
throw new ConfigurationException("commitlog_disk_access_mode = " + diskAccessMode + " is not supported. Please use 'auto' when unsure.", false);
}
}
private static void validateSSTableFormatFactories(Iterable<SSTableFormat.Factory> factories) private static void validateSSTableFormatFactories(Iterable<SSTableFormat.Factory> factories)
{ {
Map<String, SSTableFormat.Factory> factoryByName = new HashMap<>(); Map<String, SSTableFormat.Factory> factoryByName = new HashMap<>();
@ -2547,6 +2602,7 @@ public class DatabaseDescriptor
return conf.commitlog_compression; return conf.commitlog_compression;
} }
@VisibleForTesting
public static void setCommitLogCompression(ParameterizedClass compressor) public static void setCommitLogCompression(ParameterizedClass compressor)
{ {
conf.commitlog_compression = compressor; conf.commitlog_compression = compressor;
@ -2642,6 +2698,28 @@ public class DatabaseDescriptor
conf.commitlog_segment_size = new DataStorageSpec.IntMebibytesBound(sizeMebibytes); conf.commitlog_segment_size = new DataStorageSpec.IntMebibytesBound(sizeMebibytes);
} }
/**
* Return commitlog disk access mode.
*/
public static DiskAccessMode getCommitLogWriteDiskAccessMode()
{
return commitLogWriteDiskAccessMode;
}
@VisibleForTesting
public static void setCommitLogWriteDiskAccessMode(DiskAccessMode diskAccessMode)
{
conf.commitlog_disk_access_mode = diskAccessMode;
}
@VisibleForTesting
public static void initializeCommitLogDiskAccessMode()
{
DiskAccessMode resolved = resolveCommitLogWriteDiskAccessMode(conf.commitlog_disk_access_mode);
validateCommitLogWriteDiskAccessMode(resolved);
commitLogWriteDiskAccessMode = resolved;
}
public static String getSavedCachesLocation() public static String getSavedCachesLocation()
{ {
return conf.saved_caches_directory; return conf.saved_caches_directory;
@ -3141,26 +3219,26 @@ public class DatabaseDescriptor
conf.commitlog_sync = sync; conf.commitlog_sync = sync;
} }
public static Config.DiskAccessMode getDiskAccessMode() public static DiskAccessMode getDiskAccessMode()
{ {
return conf.disk_access_mode; return conf.disk_access_mode;
} }
// Do not use outside unit tests. // Do not use outside unit tests.
@VisibleForTesting @VisibleForTesting
public static void setDiskAccessMode(Config.DiskAccessMode mode) public static void setDiskAccessMode(DiskAccessMode mode)
{ {
conf.disk_access_mode = mode; conf.disk_access_mode = mode;
} }
public static Config.DiskAccessMode getIndexAccessMode() public static DiskAccessMode getIndexAccessMode()
{ {
return indexAccessMode; return indexAccessMode;
} }
// Do not use outside unit tests. // Do not use outside unit tests.
@VisibleForTesting @VisibleForTesting
public static void setIndexAccessMode(Config.DiskAccessMode mode) public static void setIndexAccessMode(DiskAccessMode mode)
{ {
indexAccessMode = mode; indexAccessMode = mode;
} }

View File

@ -18,7 +18,12 @@
package org.apache.cassandra.db.commitlog; package org.apache.cassandra.db.commitlog;
import java.io.IOException; import java.io.IOException;
import java.util.*; import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicLong;
@ -32,11 +37,11 @@ import com.codahale.metrics.Timer.Context;
import net.nicoulaj.compilecommand.annotations.DontInline; import net.nicoulaj.compilecommand.annotations.DontInline;
import org.apache.cassandra.concurrent.Interruptible; import org.apache.cassandra.concurrent.Interruptible;
import org.apache.cassandra.concurrent.Interruptible.TerminateException; import org.apache.cassandra.concurrent.Interruptible.TerminateException;
import org.apache.cassandra.config.Config.DiskAccessMode;
import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.db.Mutation; import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.io.compress.BufferType;
import org.apache.cassandra.io.util.File; import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.SimpleCachedBufferPool; import org.apache.cassandra.io.util.SimpleCachedBufferPool;
@ -44,7 +49,11 @@ import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.concurrent.*; import org.apache.cassandra.utils.concurrent.Future;
import org.apache.cassandra.utils.concurrent.FutureCombiner;
import org.apache.cassandra.utils.concurrent.ImmediateFuture;
import org.apache.cassandra.utils.concurrent.UncheckedInterruptedException;
import org.apache.cassandra.utils.concurrent.WaitQueue;
import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory; import static org.apache.cassandra.concurrent.ExecutorFactory.Global.executorFactory;
import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON; import static org.apache.cassandra.concurrent.InfiniteLoopExecutor.Daemon.NON_DAEMON;
@ -95,10 +104,12 @@ public abstract class AbstractCommitLogSegmentManager
@VisibleForTesting @VisibleForTesting
Interruptible executor; Interruptible executor;
protected final CommitLog commitLog; private final CommitLog commitLog;
private final BooleanSupplier managerThreadWaitCondition = () -> (availableSegment == null && !atSegmentBufferLimit()); private final BooleanSupplier managerThreadWaitCondition = () -> (availableSegment == null && !atSegmentBufferLimit());
private final WaitQueue managerThreadWaitQueue = newWaitQueue(); private final WaitQueue managerThreadWaitQueue = newWaitQueue();
private volatile CommitLogSegment.Builder segmentBuilder;
private volatile SimpleCachedBufferPool bufferPool; private volatile SimpleCachedBufferPool bufferPool;
AbstractCommitLogSegmentManager(final CommitLog commitLog, String storageDirectory) AbstractCommitLogSegmentManager(final CommitLog commitLog, String storageDirectory)
@ -107,18 +118,41 @@ public abstract class AbstractCommitLogSegmentManager
this.storageDirectory = storageDirectory; this.storageDirectory = storageDirectory;
} }
private CommitLogSegment.Builder createSegmentBuilder(CommitLog.Configuration config)
{
if (config.useEncryption())
{
assert config.diskAccessMode == DiskAccessMode.standard;
return new EncryptedSegment.EncryptedSegmentBuilder(this);
}
else if (config.useCompression())
{
assert config.diskAccessMode == DiskAccessMode.standard;
return new CompressedSegment.CompressedSegmentBuilder(this);
}
else if (config.diskAccessMode == DiskAccessMode.direct)
{
return new DirectIOSegment.DirectIOSegmentBuilder(this);
}
else if (config.diskAccessMode == DiskAccessMode.mmap)
{
return new MemoryMappedSegment.MemoryMappedSegmentBuilder(this);
}
throw new AssertionError("Unsupported disk access mode: " + config.diskAccessMode);
}
CommitLog.Configuration getConfiguration()
{
return commitLog.configuration;
}
void start() void start()
{ {
// For encrypted segments we want to keep the compression buffers on-heap as we need those bytes for encryption, assert this.segmentBuilder == null;
// and we want to avoid copying from off-heap (compression buffer) to on-heap encryption APIs assert this.bufferPool == null;
BufferType bufferType = commitLog.configuration.useEncryption() || !commitLog.configuration.useCompression() this.segmentBuilder = createSegmentBuilder(commitLog.configuration);
? BufferType.ON_HEAP this.bufferPool = segmentBuilder.createBufferPool();
: commitLog.configuration.getCompressor().preferredBufferType();
this.bufferPool = new SimpleCachedBufferPool(DatabaseDescriptor.getCommitLogMaxCompressionBuffersInPool(),
DatabaseDescriptor.getCommitLogSegmentSize(),
bufferType);
AllocatorRunnable allocator = new AllocatorRunnable(); AllocatorRunnable allocator = new AllocatorRunnable();
executor = executorFactory().infiniteLoop("COMMIT-LOG-ALLOCATOR", allocator, SAFE, NON_DAEMON, SYNCHRONIZED); executor = executorFactory().infiniteLoop("COMMIT-LOG-ALLOCATOR", allocator, SAFE, NON_DAEMON, SYNCHRONIZED);
@ -206,7 +240,7 @@ public abstract class AbstractCommitLogSegmentManager
private boolean atSegmentBufferLimit() private boolean atSegmentBufferLimit()
{ {
return CommitLogSegment.usesBufferPool(commitLog) && bufferPool.atLimit(); return bufferPool != null && bufferPool.atLimit();
} }
private void maybeFlushToReclaim() private void maybeFlushToReclaim()
@ -238,7 +272,10 @@ public abstract class AbstractCommitLogSegmentManager
* Hook to allow segment managers to track state surrounding creation of new segments. Onl perform as task submit * Hook to allow segment managers to track state surrounding creation of new segments. Onl perform as task submit
* to segment manager so it's performed on segment management thread. * to segment manager so it's performed on segment management thread.
*/ */
abstract CommitLogSegment createSegment(); protected CommitLogSegment createSegment()
{
return this.segmentBuilder.build();
}
/** /**
* Indicates that a segment file has been flushed and is no longer needed. Only perform as task submit to segment * Indicates that a segment file has been flushed and is no longer needed. Only perform as task submit to segment
@ -556,6 +593,9 @@ public abstract class AbstractCommitLogSegmentManager
if (bufferPool != null) if (bufferPool != null)
bufferPool.emptyBufferPool(); bufferPool.emptyBufferPool();
this.segmentBuilder = null;
this.bufferPool = null;
return res; return res;
} }

View File

@ -40,6 +40,7 @@ import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import org.apache.cassandra.config.Config;
import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.config.ParameterizedClass; import org.apache.cassandra.config.ParameterizedClass;
import org.apache.cassandra.db.Mutation; import org.apache.cassandra.db.Mutation;
@ -104,7 +105,8 @@ public class CommitLog implements CommitLogMBean
CommitLog(CommitLogArchiver archiver, Function<CommitLog, AbstractCommitLogSegmentManager> segmentManagerProvider) CommitLog(CommitLogArchiver archiver, Function<CommitLog, AbstractCommitLogSegmentManager> segmentManagerProvider)
{ {
this.configuration = new Configuration(DatabaseDescriptor.getCommitLogCompression(), this.configuration = new Configuration(DatabaseDescriptor.getCommitLogCompression(),
DatabaseDescriptor.getEncryptionContext()); DatabaseDescriptor.getEncryptionContext(),
DatabaseDescriptor.getCommitLogWriteDiskAccessMode());
DatabaseDescriptor.createAllDirectories(); DatabaseDescriptor.createAllDirectories();
this.archiver = archiver; this.archiver = archiver;
@ -521,7 +523,8 @@ public class CommitLog implements CommitLogMBean
synchronized public void resetConfiguration() synchronized public void resetConfiguration()
{ {
configuration = new Configuration(DatabaseDescriptor.getCommitLogCompression(), configuration = new Configuration(DatabaseDescriptor.getCommitLogCompression(),
DatabaseDescriptor.getEncryptionContext()); DatabaseDescriptor.getEncryptionContext(),
DatabaseDescriptor.getCommitLogWriteDiskAccessMode());
} }
/** /**
@ -609,6 +612,11 @@ public class CommitLog implements CommitLogMBean
public static final class Configuration public static final class Configuration
{ {
/**
* Flag used to shows user configured Direct-IO status.
*/
public final Config.DiskAccessMode diskAccessMode;
/** /**
* The compressor class. * The compressor class.
*/ */
@ -622,17 +630,18 @@ public class CommitLog implements CommitLogMBean
/** /**
* The encryption context used to encrypt the segments. * The encryption context used to encrypt the segments.
*/ */
private EncryptionContext encryptionContext; private final EncryptionContext encryptionContext;
public Configuration(ParameterizedClass compressorClass, EncryptionContext encryptionContext) public Configuration(ParameterizedClass compressorClass, EncryptionContext encryptionContext,
Config.DiskAccessMode diskAccessMode)
{ {
this.compressorClass = compressorClass; this.compressorClass = compressorClass;
this.compressor = compressorClass != null ? CompressionParams.createCompressor(compressorClass) : null; this.compressor = compressorClass != null ? CompressionParams.createCompressor(compressorClass) : null;
this.encryptionContext = encryptionContext; this.encryptionContext = encryptionContext;
this.diskAccessMode = diskAccessMode;
} }
/** /**
* Checks if the segments must be compressed.
* @return <code>true</code> if the segments must be compressed, <code>false</code> otherwise. * @return <code>true</code> if the segments must be compressed, <code>false</code> otherwise.
*/ */
public boolean useCompression() public boolean useCompression()
@ -641,7 +650,6 @@ public class CommitLog implements CommitLogMBean
} }
/** /**
* Checks if the segments must be encrypted.
* @return <code>true</code> if the segments must be encrypted, <code>false</code> otherwise. * @return <code>true</code> if the segments must be encrypted, <code>false</code> otherwise.
*/ */
public boolean useEncryption() public boolean useEncryption()
@ -650,7 +658,6 @@ public class CommitLog implements CommitLogMBean
} }
/** /**
* Returns the compressor used to compress the segments.
* @return the compressor used to compress the segments * @return the compressor used to compress the segments
*/ */
public ICompressor getCompressor() public ICompressor getCompressor()
@ -659,7 +666,6 @@ public class CommitLog implements CommitLogMBean
} }
/** /**
* Returns the compressor class.
* @return the compressor class * @return the compressor class
*/ */
public ParameterizedClass getCompressorClass() public ParameterizedClass getCompressorClass()
@ -668,7 +674,6 @@ public class CommitLog implements CommitLogMBean
} }
/** /**
* Returns the compressor name.
* @return the compressor name. * @return the compressor name.
*/ */
public String getCompressorName() public String getCompressorName()
@ -677,12 +682,27 @@ public class CommitLog implements CommitLogMBean
} }
/** /**
* Returns the encryption context used to encrypt the segments.
* @return the encryption context used to encrypt the segments * @return the encryption context used to encrypt the segments
*/ */
public EncryptionContext getEncryptionContext() public EncryptionContext getEncryptionContext()
{ {
return encryptionContext; return encryptionContext;
} }
/**
* @return Direct-IO used for CommitLog IO
*/
public boolean isDirectIOEnabled()
{
return diskAccessMode == Config.DiskAccessMode.direct;
}
/**
* @return Standard or buffered I/O used for CommitLog IO
*/
public boolean isStandardModeEnable()
{
return diskAccessMode == Config.DiskAccessMode.standard;
}
} }
} }

View File

@ -20,8 +20,14 @@ package org.apache.cassandra.db.commitlog;
import java.io.IOException; import java.io.IOException;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.nio.channels.FileChannel; import java.nio.channels.FileChannel;
import java.nio.file.StandardOpenOption; import java.nio.file.Path;
import java.util.*; import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
@ -29,21 +35,21 @@ import java.util.concurrent.locks.LockSupport;
import java.util.zip.CRC32; import java.util.zip.CRC32;
import com.google.common.annotations.VisibleForTesting; import com.google.common.annotations.VisibleForTesting;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileWriter;
import org.cliffc.high_scale_lib.NonBlockingHashMap; import org.cliffc.high_scale_lib.NonBlockingHashMap;
import com.codahale.metrics.Timer; import com.codahale.metrics.Timer;
import org.apache.cassandra.config.*; import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.Mutation; import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.db.commitlog.CommitLog.Configuration;
import org.apache.cassandra.db.partitions.PartitionUpdate; import org.apache.cassandra.db.partitions.PartitionUpdate;
import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.FileWriter;
import org.apache.cassandra.io.util.SimpleCachedBufferPool;
import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableId;
import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.schema.TableMetadata;
import org.apache.cassandra.utils.NativeLibrary;
import org.apache.cassandra.utils.IntegerInterval; import org.apache.cassandra.utils.IntegerInterval;
import org.apache.cassandra.utils.concurrent.OpOrder; import org.apache.cassandra.utils.concurrent.OpOrder;
import org.apache.cassandra.utils.concurrent.WaitQueue; import org.apache.cassandra.utils.concurrent.WaitQueue;
@ -124,7 +130,6 @@ public abstract class CommitLogSegment
final File logFile; final File logFile;
final FileChannel channel; final FileChannel channel;
final int fd;
protected final AbstractCommitLogSegmentManager manager; protected final AbstractCommitLogSegmentManager manager;
@ -133,28 +138,6 @@ public abstract class CommitLogSegment
public final CommitLogDescriptor descriptor; public final CommitLogDescriptor descriptor;
static CommitLogSegment createSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager)
{
Configuration config = commitLog.configuration;
CommitLogSegment segment = config.useEncryption() ? new EncryptedSegment(commitLog, manager)
: config.useCompression() ? new CompressedSegment(commitLog, manager)
: new MemoryMappedSegment(commitLog, manager);
segment.writeLogHeader();
return segment;
}
/**
* Checks if the segments use a buffer pool.
*
* @param commitLog the commit log
* @return <code>true</code> if the segments use a buffer pool, <code>false</code> otherwise.
*/
static boolean usesBufferPool(CommitLog commitLog)
{
Configuration config = commitLog.configuration;
return config.useEncryption() || config.useCompression();
}
static long getNextId() static long getNextId()
{ {
return idBase + nextId.getAndIncrement(); return idBase + nextId.getAndIncrement();
@ -163,27 +146,26 @@ public abstract class CommitLogSegment
/** /**
* Constructs a new segment file. * Constructs a new segment file.
*/ */
CommitLogSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) CommitLogSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction<Path, FileChannel, IOException> channelFactory)
{ {
this.manager = manager; this.manager = manager;
id = getNextId(); id = getNextId();
descriptor = new CommitLogDescriptor(id, descriptor = new CommitLogDescriptor(id,
commitLog.configuration.getCompressorClass(), manager.getConfiguration().getCompressorClass(),
commitLog.configuration.getEncryptionContext()); manager.getConfiguration().getEncryptionContext());
logFile = new File(manager.storageDirectory, descriptor.fileName()); logFile = new File(manager.storageDirectory, descriptor.fileName());
try try
{ {
channel = FileChannel.open(logFile.toPath(), StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE); channel = channelFactory.apply(logFile.toPath());
fd = NativeLibrary.getfd(channel);
} }
catch (IOException e) catch (IOException e)
{ {
throw new FSWriteError(e, logFile); throw new FSWriteError(e, logFile);
} }
buffer = createBuffer(commitLog); this.buffer = createBuffer();
} }
/** /**
@ -207,7 +189,10 @@ public abstract class CommitLogSegment
return Collections.<String, String>emptyMap(); return Collections.<String, String>emptyMap();
} }
abstract ByteBuffer createBuffer(CommitLog commitLog); protected ByteBuffer createBuffer()
{
return manager.getBufferPool().createBuffer();
}
/** /**
* Allocate space in this buffer for the provided mutation, and return the allocated Allocation object. * Allocate space in this buffer for the provided mutation, and return the allocated Allocation object.
@ -764,4 +749,18 @@ public abstract class CommitLogSegment
return new CommitLogPosition(segment.id, buffer.limit()); return new CommitLogPosition(segment.id, buffer.limit());
} }
} }
protected abstract static class Builder
{
protected final AbstractCommitLogSegmentManager segmentManager;
public Builder(AbstractCommitLogSegmentManager segmentManager)
{
this.segmentManager = segmentManager;
}
public abstract CommitLogSegment build();
public abstract SimpleCachedBufferPool createBufferPool();
}
} }

View File

@ -236,7 +236,8 @@ public class CommitLogSegmentManagerCDC extends AbstractCommitLogSegmentManager
@Override @Override
public CommitLogSegment createSegment() public CommitLogSegment createSegment()
{ {
CommitLogSegment segment = CommitLogSegment.createSegment(commitLog, this); CommitLogSegment segment = super.createSegment();
segment.writeLogHeader();
cdcSizeTracker.processNewSegment(segment); cdcSizeTracker.processNewSegment(segment);
// After processing, the state of the segment can either be PERMITTED or FORBIDDEN // After processing, the state of the segment can either be PERMITTED or FORBIDDEN
if (segment.getCDCState() == CDCState.PERMITTED) if (segment.getCDCState() == CDCState.PERMITTED)

View File

@ -62,6 +62,8 @@ public class CommitLogSegmentManagerStandard extends AbstractCommitLogSegmentMan
@Override @Override
public CommitLogSegment createSegment() public CommitLogSegment createSegment()
{ {
return CommitLogSegment.createSegment(commitLog, this); CommitLogSegment segment = super.createSegment();
segment.writeLogHeader();
return segment;
} }
} }

View File

@ -17,10 +17,17 @@
*/ */
package org.apache.cassandra.db.commitlog; package org.apache.cassandra.db.commitlog;
import java.io.IOException;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.compress.ICompressor; import org.apache.cassandra.io.compress.ICompressor;
import org.apache.cassandra.io.util.SimpleCachedBufferPool;
/** /**
* Compressed commit log segment. Provides an in-memory buffer for the mutation threads. On sync compresses the written * Compressed commit log segment. Provides an in-memory buffer for the mutation threads. On sync compresses the written
@ -41,15 +48,10 @@ public class CompressedSegment extends FileDirectSegment
/** /**
* Constructs a new segment file. * Constructs a new segment file.
*/ */
CompressedSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) CompressedSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction<Path, FileChannel, IOException> channelFactory)
{ {
super(commitLog, manager); super(manager, channelFactory);
this.compressor = commitLog.configuration.getCompressor(); this.compressor = manager.getConfiguration().getCompressor();
}
ByteBuffer createBuffer(CommitLog commitLog)
{
return manager.getBufferPool().createBuffer();
} }
@Override @Override
@ -92,4 +94,27 @@ public class CompressedSegment extends FileDirectSegment
{ {
return lastWrittenPos; return lastWrittenPos;
} }
protected static class CompressedSegmentBuilder extends CommitLogSegment.Builder
{
public CompressedSegmentBuilder(AbstractCommitLogSegmentManager segmentManager)
{
super(segmentManager);
}
@Override
public CompressedSegment build()
{
return new CompressedSegment(segmentManager,
path -> FileChannel.open(path, StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE));
}
@Override
public SimpleCachedBufferPool createBufferPool()
{
return new SimpleCachedBufferPool(DatabaseDescriptor.getCommitLogMaxCompressionBuffersInPool(),
DatabaseDescriptor.getCommitLogSegmentSize(),
segmentManager.getConfiguration().getCompressor().preferredBufferType());
}
}
} }

View File

@ -0,0 +1,228 @@
/*
* 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.db.commitlog;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import com.google.common.annotations.VisibleForTesting;
import com.sun.nio.file.ExtendedOpenOption;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.compress.BufferType;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.SimpleCachedBufferPool;
import org.apache.cassandra.utils.ByteBufferUtil;
import sun.nio.ch.DirectBuffer;
/*
* Direct-IO segment. Allocates ByteBuffer using ByteBuffer.allocateDirect and align
* ByteBuffer.position, ByteBuffer.limit and FileChannel.position to page size (4K).
* Java-11 forces minimum page size to be written to disk with Direct-IO.
*/
public class DirectIOSegment extends CommitLogSegment
{
private final int fsBlockSize;
private final int fsBlockRemainderMask;
// Needed to track number of bytes written to disk in multiple of page size.
long lastWritten = 0;
/**
* Constructs a new segment file.
*/
DirectIOSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction<Path, FileChannel, IOException> channelFactory, int fsBlockSize)
{
super(manager, channelFactory);
assert Integer.highestOneBit(fsBlockSize) == fsBlockSize : "fsBlockSize must be a power of 2";
// mark the initial sync marker as uninitialised
int firstSync = buffer.position();
buffer.putInt(firstSync + 0, 0);
buffer.putInt(firstSync + 4, 0);
this.fsBlockSize = fsBlockSize;
this.fsBlockRemainderMask = fsBlockSize - 1;
}
@Override
void writeLogHeader()
{
super.writeLogHeader();
// Testing shows writing initial bytes takes some time for Direct I/O. During peak load,
// it is better to make "COMMIT-LOG-ALLOCATOR" thread to write these few bytes of each
// file and this helps syncer thread to speedup the flush activity.
flush(0, lastSyncedOffset);
}
@Override
void write(int startMarker, int nextMarker)
{
// if there's room in the discard section to write an empty header,
// zero out the next sync marker so replayer can cleanly exit
if (nextMarker <= buffer.capacity() - SYNC_MARKER_SIZE)
{
buffer.putInt(nextMarker, 0);
buffer.putInt(nextMarker + 4, 0);
}
// write previous sync marker to point to next sync marker
// we don't chain the crcs here to ensure this method is idempotent if it fails
writeSyncMarker(id, buffer, startMarker, startMarker, nextMarker);
}
@Override
protected void flush(int startMarker, int nextMarker)
{
try
{
// TODO move the alignment calculations to PageAware
// lastSyncedOffset is synced to disk. Align lastSyncedOffset to start of its block
// and nextMarker to end of its block to avoid write errors.
int flushPosition = lastSyncedOffset;
ByteBuffer duplicate = buffer.duplicate();
// Aligned file position if not aligned to start of a block.
if ((flushPosition & fsBlockRemainderMask) != 0)
{
flushPosition = flushPosition & -fsBlockSize;
channel.position(flushPosition);
}
duplicate.position(flushPosition);
int flushLimit = nextMarker;
// Align last byte to end of block
flushLimit = (flushLimit + fsBlockSize - 1) & -fsBlockSize;
duplicate.limit(flushLimit);
channel.write(duplicate);
// Direct I/O always writes flushes in block size and writes more than the flush size.
// File size on disk will always multiple of block size and taking this into account
// helps testcases to pass. Avoid counting same block more than once.
if (flushLimit > lastWritten)
{
manager.addSize(flushLimit - lastWritten);
lastWritten = flushLimit;
}
}
catch (IOException e)
{
throw new FSWriteError(e, getPath());
}
}
@Override
public long onDiskSize()
{
return lastWritten;
}
@Override
protected void internalClose()
{
try
{
manager.getBufferPool().releaseBuffer(buffer);
super.internalClose();
}
finally
{
manager.notifyBufferFreed();
}
}
protected static class DirectIOSegmentBuilder extends CommitLogSegment.Builder
{
public final int fsBlockSize;
public DirectIOSegmentBuilder(AbstractCommitLogSegmentManager segmentManager)
{
this(segmentManager, FileUtils.getBlockSize(new File(segmentManager.storageDirectory)));
}
@VisibleForTesting
public DirectIOSegmentBuilder(AbstractCommitLogSegmentManager segmentManager, int fsBlockSize)
{
super(segmentManager);
this.fsBlockSize = fsBlockSize;
}
@Override
public DirectIOSegment build()
{
return new DirectIOSegment(segmentManager,
path -> FileChannel.open(path, StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE, ExtendedOpenOption.DIRECT),
fsBlockSize);
}
@Override
public SimpleCachedBufferPool createBufferPool()
{
// The direct buffer must be aligned with the file system block size. We cannot enforce that during
// allocation, but we can get an aligned slice from the allocated buffer. The buffer must be oversized by the
// alignment unit to make it possible.
return new SimpleCachedBufferPool(DatabaseDescriptor.getCommitLogMaxCompressionBuffersInPool(),
DatabaseDescriptor.getCommitLogSegmentSize() + fsBlockSize,
BufferType.OFF_HEAP) {
@Override
public ByteBuffer createBuffer()
{
int segmentSize = DatabaseDescriptor.getCommitLogSegmentSize();
ByteBuffer original = super.createBuffer();
// May get previously used buffer and zero it out to now. Direct I/O writes additional bytes during
// flush operation
ByteBufferUtil.writeZeroes(original.duplicate(), original.limit());
ByteBuffer alignedBuffer;
if (original.alignmentOffset(0, fsBlockSize) > 0)
alignedBuffer = original.alignedSlice(fsBlockSize);
else
alignedBuffer = original.slice().limit(segmentSize);
assert alignedBuffer.limit() >= segmentSize : String.format("Bytebuffer slicing failed to get required buffer size (required=%d, current size=%d", segmentSize, alignedBuffer.limit());
assert alignedBuffer.alignmentOffset(0, fsBlockSize) == 0 : String.format("Index 0 should be aligned to %d page size.", fsBlockSize);
assert alignedBuffer.alignmentOffset(alignedBuffer.limit(), fsBlockSize) == 0 : String.format("Limit should be aligned to %d page size", fsBlockSize);
return alignedBuffer;
}
@Override
public void releaseBuffer(ByteBuffer buffer)
{
ByteBuffer original = (ByteBuffer) ((DirectBuffer) buffer).attachment();
assert original != null;
super.releaseBuffer(original);
}
};
}
}
}

View File

@ -19,17 +19,23 @@ package org.apache.cassandra.db.commitlog;
import java.io.IOException; import java.io.IOException;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import java.util.Map; import java.util.Map;
import javax.crypto.Cipher; import javax.crypto.Cipher;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.compress.BufferType;
import org.apache.cassandra.io.compress.ICompressor; import org.apache.cassandra.io.compress.ICompressor;
import org.apache.cassandra.security.EncryptionUtils; import org.apache.cassandra.io.util.SimpleCachedBufferPool;
import org.apache.cassandra.security.EncryptionContext; import org.apache.cassandra.security.EncryptionContext;
import org.apache.cassandra.security.EncryptionUtils;
import org.apache.cassandra.utils.Hex; import org.apache.cassandra.utils.Hex;
import static org.apache.cassandra.security.EncryptionUtils.ENCRYPTED_BLOCK_HEADER_SIZE; import static org.apache.cassandra.security.EncryptionUtils.ENCRYPTED_BLOCK_HEADER_SIZE;
@ -63,10 +69,10 @@ public class EncryptedSegment extends FileDirectSegment
private final EncryptionContext encryptionContext; private final EncryptionContext encryptionContext;
private final Cipher cipher; private final Cipher cipher;
public EncryptedSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) public EncryptedSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction<Path, FileChannel, IOException> channelFactory)
{ {
super(commitLog, manager); super(manager, channelFactory);
this.encryptionContext = commitLog.configuration.getEncryptionContext(); this.encryptionContext = manager.getConfiguration().getEncryptionContext();
try try
{ {
@ -86,12 +92,9 @@ public class EncryptedSegment extends FileDirectSegment
return map; return map;
} }
ByteBuffer createBuffer(CommitLog commitLog) // Note: we want to keep the compression buffers on-heap as we need those bytes for encryption,
{ // and we want to avoid copying from off-heap (compression buffer) to on-heap encryption APIs
// Note: we want to keep the compression buffers on-heap as we need those bytes for encryption, // (so we do not override the createBuffer method)
// and we want to avoid copying from off-heap (compression buffer) to on-heap encryption APIs
return manager.getBufferPool().createBuffer();
}
void write(int startMarker, int nextMarker) void write(int startMarker, int nextMarker)
{ {
@ -150,4 +153,28 @@ public class EncryptedSegment extends FileDirectSegment
{ {
return lastWrittenPos; return lastWrittenPos;
} }
protected static class EncryptedSegmentBuilder extends CommitLogSegment.Builder
{
public EncryptedSegmentBuilder(AbstractCommitLogSegmentManager segmentManager)
{
super(segmentManager);
}
@Override
public EncryptedSegment build()
{
return new EncryptedSegment(segmentManager,
path -> FileChannel.open(path, StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE));
}
@Override
public SimpleCachedBufferPool createBufferPool()
{
return new SimpleCachedBufferPool(DatabaseDescriptor.getCommitLogMaxCompressionBuffersInPool(),
DatabaseDescriptor.getCommitLogSegmentSize(),
BufferType.ON_HEAP);
}
}
} }

View File

@ -19,7 +19,10 @@ package org.apache.cassandra.db.commitlog;
import java.io.IOException; import java.io.IOException;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.file.Path;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.utils.SyncUtil; import org.apache.cassandra.utils.SyncUtil;
@ -31,9 +34,9 @@ public abstract class FileDirectSegment extends CommitLogSegment
{ {
volatile long lastWrittenPos = 0; volatile long lastWrittenPos = 0;
FileDirectSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) FileDirectSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction<Path, FileChannel, IOException> channelFactory)
{ {
super(commitLog, manager); super(manager, channelFactory);
} }
@Override @Override

View File

@ -21,10 +21,14 @@ import java.io.IOException;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.nio.MappedByteBuffer; import java.nio.MappedByteBuffer;
import java.nio.channels.FileChannel; import java.nio.channels.FileChannel;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.io.util.SimpleCachedBufferPool;
import org.apache.cassandra.utils.NativeLibrary; import org.apache.cassandra.utils.NativeLibrary;
import org.apache.cassandra.utils.SyncUtil; import org.apache.cassandra.utils.SyncUtil;
@ -35,21 +39,23 @@ import org.apache.cassandra.utils.SyncUtil;
*/ */
public class MemoryMappedSegment extends CommitLogSegment public class MemoryMappedSegment extends CommitLogSegment
{ {
private final int fd;
/** /**
* Constructs a new segment file. * Constructs a new segment file.
*
* @param commitLog the commit log it will be used with.
*/ */
MemoryMappedSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) MemoryMappedSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction<Path, FileChannel, IOException> channelFactory)
{ {
super(commitLog, manager); super(manager, channelFactory);
// mark the initial sync marker as uninitialised // mark the initial sync marker as uninitialised
int firstSync = buffer.position(); int firstSync = buffer.position();
buffer.putInt(firstSync + 0, 0); buffer.putInt(firstSync + 0, 0);
buffer.putInt(firstSync + 4, 0); buffer.putInt(firstSync + 4, 0);
fd = NativeLibrary.getfd(channel);
} }
ByteBuffer createBuffer(CommitLog commitLog) @Override
protected ByteBuffer createBuffer()
{ {
try try
{ {
@ -105,4 +111,25 @@ public class MemoryMappedSegment extends CommitLogSegment
FileUtils.clean(buffer); FileUtils.clean(buffer);
super.internalClose(); super.internalClose();
} }
protected static class MemoryMappedSegmentBuilder extends CommitLogSegment.Builder
{
public MemoryMappedSegmentBuilder(AbstractCommitLogSegmentManager segmentManager)
{
super(segmentManager);
}
@Override
public MemoryMappedSegment build()
{
return new MemoryMappedSegment(segmentManager,
path -> FileChannel.open(path, StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE));
}
@Override
public SimpleCachedBufferPool createBufferPool()
{
return null;
}
}
} }

View File

@ -816,4 +816,23 @@ public final class FileUtils
} }
} }
} }
public static int getBlockSize(File directory)
{
File f = FileUtils.createTempFile("block-size-test", ".tmp", directory);
try
{
long bs = Files.getFileStore(f.toPath()).getBlockSize();
assert bs >= 0 && bs <= Integer.MAX_VALUE;
return (int) bs;
}
catch (IOException e)
{
throw new RuntimeException("Failed to get file block size in " + directory, e);
}
finally
{
f.tryDelete();
}
}
} }

View File

@ -49,21 +49,20 @@ import java.util.function.Function;
import java.util.function.IntFunction; import java.util.function.IntFunction;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.Stream; import java.util.stream.Stream;
import javax.annotation.Nullable; import javax.annotation.Nullable;
import com.google.common.annotations.VisibleForTesting; import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions; import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.RateLimiter; import com.google.common.util.concurrent.RateLimiter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.CassandraRelevantProperties; import org.apache.cassandra.config.CassandraRelevantProperties;
import org.apache.cassandra.io.FSError; import org.apache.cassandra.io.FSError;
import org.apache.cassandra.io.FSReadError; import org.apache.cassandra.io.FSReadError;
import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.NoSpamLogger; import org.apache.cassandra.utils.NoSpamLogger;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import static java.nio.file.StandardOpenOption.APPEND; import static java.nio.file.StandardOpenOption.APPEND;
import static java.nio.file.StandardOpenOption.CREATE; import static java.nio.file.StandardOpenOption.CREATE;
@ -71,12 +70,9 @@ import static java.nio.file.StandardOpenOption.READ;
import static java.nio.file.StandardOpenOption.TRUNCATE_EXISTING; import static java.nio.file.StandardOpenOption.TRUNCATE_EXISTING;
import static java.nio.file.StandardOpenOption.WRITE; import static java.nio.file.StandardOpenOption.WRITE;
import static java.util.Collections.unmodifiableSet; import static java.util.Collections.unmodifiableSet;
import static org.apache.cassandra.config.CassandraRelevantProperties.USE_NIX_RECURSIVE_DELETE; import static org.apache.cassandra.config.CassandraRelevantProperties.USE_NIX_RECURSIVE_DELETE;
import static org.apache.cassandra.utils.Throwables.merge; import static org.apache.cassandra.utils.Throwables.merge;
import net.openhft.chronicle.core.util.ThrowingFunction;
/** /**
* Vernacular: tryX means return false or 0L on any failure; XIfNotY means propagate any exceptions besides those caused by Y * Vernacular: tryX means return false or 0L on any failure; XIfNotY means propagate any exceptions besides those caused by Y
* *
@ -98,12 +94,7 @@ public final class PathUtils
private static final Logger logger = LoggerFactory.getLogger(PathUtils.class); private static final Logger logger = LoggerFactory.getLogger(PathUtils.class);
private static final NoSpamLogger nospam1m = NoSpamLogger.getLogger(logger, 1, TimeUnit.MINUTES); private static final NoSpamLogger nospam1m = NoSpamLogger.getLogger(logger, 1, TimeUnit.MINUTES);
private static Consumer<Path> onDeletion = path -> { private static Consumer<Path> onDeletion = path -> {};
if (StorageService.instance.isDaemonSetupCompleted())
setDeletionListener(ignore -> {});
else
logger.trace("Deleting file during startup: {}", path);
};
public static FileChannel newReadChannel(Path path) throws NoSuchFileException public static FileChannel newReadChannel(Path path) throws NoSuchFileException
{ {

View File

@ -290,6 +290,15 @@ public class StorageService extends NotificationBroadcasterSupport implements IE
private static final boolean REQUIRE_SCHEMAS = !BOOTSTRAP_SKIP_SCHEMA_CHECK.getBoolean(); private static final boolean REQUIRE_SCHEMAS = !BOOTSTRAP_SKIP_SCHEMA_CHECK.getBoolean();
{
PathUtils.setDeletionListener(path -> {
if (isDaemonSetupCompleted())
PathUtils.setDeletionListener(ignore -> {});
else
logger.trace("Deleting file during startup: {}", path);
});
}
private final JMXProgressSupport progressSupport = new JMXProgressSupport(this); private final JMXProgressSupport progressSupport = new JMXProgressSupport(this);
private static int getRingDelay() private static int getRingDelay()

View File

@ -8,6 +8,7 @@ memtable_allocation_type: offheap_objects
commitlog_sync: batch commitlog_sync: batch
commitlog_segment_size: 5MiB commitlog_segment_size: 5MiB
commitlog_directory: build/test/cassandra/commitlog commitlog_directory: build/test/cassandra/commitlog
commitlog_disk_access_mode: legacy
# commitlog_compression: # commitlog_compression:
# - class_name: LZ4Compressor # - class_name: LZ4Compressor
cdc_raw_directory: build/test/cassandra/cdc_raw cdc_raw_directory: build/test/cassandra/cdc_raw

View File

@ -112,7 +112,8 @@ public class InstanceConfig implements IInstanceConfig
// capacities that are based on `totalMemory` that should be fixed size // capacities that are based on `totalMemory` that should be fixed size
.set("index_summary_capacity", "50MiB") .set("index_summary_capacity", "50MiB")
.set("counter_cache_size", "50MiB") .set("counter_cache_size", "50MiB")
.set("key_cache_size", "50MiB"); .set("key_cache_size", "50MiB")
.set("commitlog_disk_access_mode", "legacy");
this.featureFlags = EnumSet.noneOf(Feature.class); this.featureFlags = EnumSet.noneOf(Feature.class);
this.jmxPort = jmx_port; this.jmxPort = jmx_port;
} }

View File

@ -29,9 +29,9 @@ import org.apache.cassandra.security.EncryptionContext;
@RunWith(Parameterized.class) @RunWith(Parameterized.class)
public class BatchCommitLogStressTest extends CommitLogStressTest public class BatchCommitLogStressTest extends CommitLogStressTest
{ {
public BatchCommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext) public BatchCommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext, Config.DiskAccessMode accessMode)
{ {
super(commitLogCompression, encryptionContext); super(commitLogCompression, encryptionContext, accessMode);
DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.batch); DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.batch);
} }
} }

View File

@ -93,12 +93,15 @@ public abstract class CommitLogStressTest
private boolean randomSize = false; private boolean randomSize = false;
private boolean discardedRun = false; private boolean discardedRun = false;
private CommitLogPosition discardedPos; private CommitLogPosition discardedPos;
private long totalBytesWritten = 0;
public CommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext) public CommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext, Config.DiskAccessMode accessMode)
{ {
DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setCommitLogCompression(commitLogCompression);
DatabaseDescriptor.setEncryptionContext(encryptionContext); DatabaseDescriptor.setEncryptionContext(encryptionContext);
DatabaseDescriptor.setCommitLogSegmentSize(32); DatabaseDescriptor.setCommitLogSegmentSize(32);
DatabaseDescriptor.setCommitLogWriteDiskAccessMode(accessMode);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
@BeforeClass @BeforeClass
@ -142,11 +145,12 @@ public abstract class CommitLogStressTest
public static Collection<Object[]> buildParameterizedVariants() public static Collection<Object[]> buildParameterizedVariants()
{ {
return Arrays.asList(new Object[][]{ return Arrays.asList(new Object[][]{
{null, EncryptionContextGenerator.createDisabledContext()}, // No compression, no encryption {null, EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.legacy}, // No compression, no encryption, legacy
{null, EncryptionContextGenerator.createContext(true)}, // Encryption {null, EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.direct}, // Use Direct-I/O (non-buffered) feature.
{ new ParameterizedClass(LZ4Compressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext()}, {null, EncryptionContextGenerator.createContext(true), Config.DiskAccessMode.legacy},
{ new ParameterizedClass(SnappyCompressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext()}, { new ParameterizedClass(LZ4Compressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.legacy},
{ new ParameterizedClass(DeflateCompressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext()}}); { new ParameterizedClass(SnappyCompressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.legacy},
{ new ParameterizedClass(DeflateCompressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.legacy}});
} }
@Test @Test
@ -157,6 +161,7 @@ public abstract class CommitLogStressTest
testLog(); testLog();
} }
@Test @Test
public void testFixedSize() throws Exception public void testFixedSize() throws Exception
{ {
@ -179,6 +184,7 @@ public abstract class CommitLogStressTest
try try
{ {
DatabaseDescriptor.setCommitLogLocation(location); DatabaseDescriptor.setCommitLogLocation(location);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
CommitLog commitLog = new CommitLog(CommitLogArchiver.disabled()).start(); CommitLog commitLog = new CommitLog(CommitLogArchiver.disabled()).start();
testLog(commitLog); testLog(commitLog);
assert !failed; assert !failed;
@ -186,14 +192,17 @@ public abstract class CommitLogStressTest
finally finally
{ {
DatabaseDescriptor.setCommitLogLocation(originalDir); DatabaseDescriptor.setCommitLogLocation(originalDir);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
} }
private void testLog(CommitLog commitLog) throws IOException, InterruptedException { private void testLog(CommitLog commitLog) throws IOException, InterruptedException {
System.out.format("\nTesting commit log size %.0fmb, compressor: %s, encryption enabled: %b, sync %s%s%s\n", System.out.format("\nTesting commit log size %.0fmb, disk mode: %s, compressor: %s, encryption enabled: %b, direct I/O enabled: %b, sync %s%s%s\n",
mb(DatabaseDescriptor.getCommitLogSegmentSize()), mb(DatabaseDescriptor.getCommitLogSegmentSize()),
DatabaseDescriptor.getCommitLogWriteDiskAccessMode(),
commitLog.configuration.getCompressorName(), commitLog.configuration.getCompressorName(),
commitLog.configuration.useEncryption(), commitLog.configuration.useEncryption(),
commitLog.configuration.isDirectIOEnabled(),
commitLog.executor.getClass().getSimpleName(), commitLog.executor.getClass().getSimpleName(),
randomSize ? " random size" : "", randomSize ? " random size" : "",
discardedRun ? " with discarded run" : ""); discardedRun ? " with discarded run" : "");
@ -258,15 +267,20 @@ public abstract class CommitLogStressTest
Assert.fail("Failed to delete " + f); Assert.fail("Failed to delete " + f);
if (hash == reader.hash && cells == reader.cells) if (hash == reader.hash && cells == reader.cells)
System.out.format("Test success. compressor = %s, encryption enabled = %b; discarded = %d, skipped = %d\n", System.out.format("Test success. disk mode = %s, compressor = %s, encryption enabled = %b, direct I/O = %b; discarded = %d, skipped = %d; IO speed(total bytes=%.2fmb, rate=%.2fmb/sec)\n",
DatabaseDescriptor.getCommitLogWriteDiskAccessMode(),
commitLog.configuration.getCompressorName(), commitLog.configuration.getCompressorName(),
commitLog.configuration.useEncryption(), commitLog.configuration.useEncryption(),
reader.discarded, reader.skipped); commitLog.configuration.isDirectIOEnabled(),
reader.discarded, reader.skipped,
mb(totalBytesWritten), mb(totalBytesWritten)/runTimeMs*1000);
else else
{ {
System.out.format("Test failed (compressor = %s, encryption enabled = %b). Cells %d, expected %d, diff %d; discarded = %d, skipped = %d - hash %d expected %d.\n", System.out.format("Test failed (disk mode = %s, compressor = %s, encryption enabled = %b, direct I/O = %b). Cells %d, expected %d, diff %d; discarded = %d, skipped = %d - hash %d expected %d.\n",
DatabaseDescriptor.getCommitLogWriteDiskAccessMode(),
commitLog.configuration.getCompressorName(), commitLog.configuration.getCompressorName(),
commitLog.configuration.useEncryption(), commitLog.configuration.useEncryption(),
commitLog.configuration.isDirectIOEnabled(),
reader.cells, cells, cells - reader.cells, reader.discarded, reader.skipped, reader.cells, cells, cells - reader.cells, reader.discarded, reader.skipped,
reader.hash, hash); reader.hash, hash);
failed = true; failed = true;
@ -282,7 +296,10 @@ public abstract class CommitLogStressTest
long combinedSize = 0; long combinedSize = 0;
for (File f : new File(commitLog.segmentManager.storageDirectory).tryList()) for (File f : new File(commitLog.segmentManager.storageDirectory).tryList())
{
combinedSize += f.length(); combinedSize += f.length();
}
totalBytesWritten = combinedSize;
Assert.assertEquals(combinedSize, commitLog.getActiveOnDiskSize()); Assert.assertEquals(combinedSize, commitLog.getActiveOnDiskSize());
List<String> logFileNames = commitLog.getActiveSegmentNames(); List<String> logFileNames = commitLog.getActiveSegmentNames();

View File

@ -29,9 +29,9 @@ import org.apache.cassandra.security.EncryptionContext;
@RunWith(Parameterized.class) @RunWith(Parameterized.class)
public class GroupCommitLogStressTest extends CommitLogStressTest public class GroupCommitLogStressTest extends CommitLogStressTest
{ {
public GroupCommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext) public GroupCommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext, Config.DiskAccessMode accessMode)
{ {
super(commitLogCompression, encryptionContext); super(commitLogCompression, encryptionContext, accessMode);
DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.group); DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.group);
DatabaseDescriptor.setCommitLogSyncGroupWindow(1); DatabaseDescriptor.setCommitLogSyncGroupWindow(1);
} }

View File

@ -29,9 +29,9 @@ import org.apache.cassandra.security.EncryptionContext;
@RunWith(Parameterized.class) @RunWith(Parameterized.class)
public class PeriodicCommitLogStressTest extends CommitLogStressTest public class PeriodicCommitLogStressTest extends CommitLogStressTest
{ {
public PeriodicCommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext) public PeriodicCommitLogStressTest(ParameterizedClass commitLogCompression, EncryptionContext encryptionContext, Config.DiskAccessMode accessMode)
{ {
super(commitLogCompression, encryptionContext); super(commitLogCompression, encryptionContext, accessMode);
DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic);
DatabaseDescriptor.setCommitLogSyncPeriod(30); DatabaseDescriptor.setCommitLogSyncPeriod(30);
} }

View File

@ -18,16 +18,18 @@
*/ */
package org.apache.cassandra.config; package org.apache.cassandra.config;
import java.io.IOException;
import java.net.Inet4Address; import java.net.Inet4Address;
import java.net.Inet6Address; import java.net.Inet6Address;
import java.net.InetAddress; import java.net.InetAddress;
import java.net.NetworkInterface; import java.net.NetworkInterface;
import java.nio.file.Files;
import java.util.Arrays; import java.util.Arrays;
import java.util.Collection; import java.util.Collection;
import java.util.EnumSet;
import java.util.Enumeration; import java.util.Enumeration;
import java.util.function.Consumer; import java.util.function.Consumer;
import com.google.common.base.Throwables; import com.google.common.base.Throwables;
import org.junit.Assert; import org.junit.Assert;
import org.junit.BeforeClass; import org.junit.BeforeClass;
@ -36,6 +38,8 @@ import org.junit.Test;
import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.Keyspace;
import org.apache.cassandra.distributed.shared.WithProperties; import org.apache.cassandra.distributed.shared.WithProperties;
import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.security.EncryptionContext;
import org.apache.cassandra.security.EncryptionContextGenerator;
import org.assertj.core.api.Assertions; import org.assertj.core.api.Assertions;
import static org.apache.cassandra.config.CassandraRelevantProperties.ALLOW_UNLIMITED_CONCURRENT_VALIDATIONS; import static org.apache.cassandra.config.CassandraRelevantProperties.ALLOW_UNLIMITED_CONCURRENT_VALIDATIONS;
@ -43,6 +47,7 @@ import static org.apache.cassandra.config.CassandraRelevantProperties.CONFIG_LOA
import static org.apache.cassandra.config.CassandraRelevantProperties.PARTITIONER; import static org.apache.cassandra.config.CassandraRelevantProperties.PARTITIONER;
import static org.apache.cassandra.config.DataStorageSpec.DataStorageUnit.KIBIBYTES; import static org.apache.cassandra.config.DataStorageSpec.DataStorageUnit.KIBIBYTES;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue; import static org.junit.Assert.assertTrue;
@ -805,4 +810,106 @@ public class DatabaseDescriptorTest
{ {
DatabaseDescriptor.setDefaultKeyspaceRF(0); DatabaseDescriptor.setDefaultKeyspaceRF(0);
} }
@Test
public void testCommitLogDiskAccessMode() throws IOException
{
ParameterizedClass savedCompression = DatabaseDescriptor.getCommitLogCompression();
EncryptionContext savedEncryptionContexg = DatabaseDescriptor.getEncryptionContext();
Config.DiskAccessMode savedCommitLogDOS = DatabaseDescriptor.getCommitLogWriteDiskAccessMode();
String savedCommitLogLocation = DatabaseDescriptor.getCommitLogLocation();
try
{
// block size available
DatabaseDescriptor.setCommitLogLocation(Files.createTempDirectory("testCommitLogDiskAccessMode").toString());
// no encryption or compression
DatabaseDescriptor.setCommitLogCompression(null);
DatabaseDescriptor.setEncryptionContext(null);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.spinning;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.mmap, Config.DiskAccessMode.mmap, Config.DiskAccessMode.mmap, Config.DiskAccessMode.direct);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.ssd;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.mmap, Config.DiskAccessMode.direct, Config.DiskAccessMode.mmap, Config.DiskAccessMode.direct);
// compression enabled
DatabaseDescriptor.setCommitLogCompression(new ParameterizedClass("LZ4Compressor", null));
DatabaseDescriptor.setEncryptionContext(null);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.spinning;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.ssd;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
// encryption enabled
DatabaseDescriptor.setCommitLogCompression(null);
DatabaseDescriptor.setEncryptionContext(new EncryptionContext(EncryptionContextGenerator.createEncryptionOptions()));
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.spinning;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.ssd;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
// block size not available
DatabaseDescriptor.setCommitLogLocation(null);
// no encryption or compression
DatabaseDescriptor.setCommitLogCompression(null);
DatabaseDescriptor.setEncryptionContext(null);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.spinning;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.mmap, Config.DiskAccessMode.mmap, Config.DiskAccessMode.mmap, Config.DiskAccessMode.direct);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.ssd;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.mmap, Config.DiskAccessMode.mmap, Config.DiskAccessMode.mmap, Config.DiskAccessMode.direct);
// compression enabled
DatabaseDescriptor.setCommitLogCompression(new ParameterizedClass("LZ4Compressor", null));
DatabaseDescriptor.setEncryptionContext(null);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.spinning;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.ssd;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
// encryption enabled
DatabaseDescriptor.setCommitLogCompression(null);
DatabaseDescriptor.setEncryptionContext(new EncryptionContext(EncryptionContextGenerator.createEncryptionOptions()));
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.spinning;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
DatabaseDescriptor.getRawConfig().disk_optimization_strategy = Config.DiskOptimizationStrategy.ssd;
assertCommitLogDiskAccessModes(Config.DiskAccessMode.standard, Config.DiskAccessMode.standard, Config.DiskAccessMode.standard);
}
finally
{
DatabaseDescriptor.setCommitLogCompression(savedCompression);
DatabaseDescriptor.setEncryptionContext(savedEncryptionContexg);
DatabaseDescriptor.setCommitLogWriteDiskAccessMode(savedCommitLogDOS);
DatabaseDescriptor.setCommitLogLocation(savedCommitLogLocation);
}
}
private void assertCommitLogDiskAccessModes(Config.DiskAccessMode expectedLegacy, Config.DiskAccessMode expectedAuto, Config.DiskAccessMode... allowedModesArray)
{
EnumSet<Config.DiskAccessMode> allowedModes = EnumSet.copyOf(Arrays.asList(allowedModesArray));
allowedModes.add(Config.DiskAccessMode.legacy);
allowedModes.add(Config.DiskAccessMode.auto);
EnumSet<Config.DiskAccessMode> disallowedModes = EnumSet.complementOf(allowedModes);
for (Config.DiskAccessMode mode : disallowedModes)
{
DatabaseDescriptor.setCommitLogWriteDiskAccessMode(mode);
assertThatExceptionOfType(ConfigurationException.class).isThrownBy(DatabaseDescriptor::initializeCommitLogDiskAccessMode);
}
for (Config.DiskAccessMode mode : allowedModes)
{
DatabaseDescriptor.setCommitLogWriteDiskAccessMode(mode);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
boolean changed = DatabaseDescriptor.getCommitLogWriteDiskAccessMode() != mode;
assertThat(changed).isEqualTo(mode == Config.DiskAccessMode.legacy || mode == Config.DiskAccessMode.auto);
if (mode == Config.DiskAccessMode.legacy)
assertThat(DatabaseDescriptor.getCommitLogWriteDiskAccessMode()).isEqualTo(expectedLegacy);
else if (mode == Config.DiskAccessMode.auto)
assertThat(DatabaseDescriptor.getCommitLogWriteDiskAccessMode()).isEqualTo(expectedAuto);
else
assertThat(DatabaseDescriptor.getCommitLogWriteDiskAccessMode()).isEqualTo(mode);
}
}
} }

View File

@ -65,6 +65,7 @@ public class RecoveryManagerFlushedTest
{ {
DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setCommitLogCompression(commitLogCompression);
DatabaseDescriptor.setEncryptionContext(encryptionContext); DatabaseDescriptor.setEncryptionContext(encryptionContext);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
@Parameters() @Parameters()

View File

@ -61,6 +61,7 @@ public class RecoveryManagerMissingHeaderTest
{ {
DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setCommitLogCompression(commitLogCompression);
DatabaseDescriptor.setEncryptionContext(encryptionContext); DatabaseDescriptor.setEncryptionContext(encryptionContext);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
@Parameters() @Parameters()

View File

@ -78,6 +78,7 @@ public class RecoveryManagerTest
{ {
DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setCommitLogCompression(commitLogCompression);
DatabaseDescriptor.setEncryptionContext(encryptionContext); DatabaseDescriptor.setEncryptionContext(encryptionContext);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
@Parameters() @Parameters()

View File

@ -59,6 +59,7 @@ public class RecoveryManagerTruncateTest
{ {
DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setCommitLogCompression(commitLogCompression);
DatabaseDescriptor.setEncryptionContext(encryptionContext); DatabaseDescriptor.setEncryptionContext(encryptionContext);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
@Parameters() @Parameters()

View File

@ -23,8 +23,8 @@ import java.nio.ByteBuffer;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Random; import java.util.Random;
import org.apache.cassandra.io.util.File;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Assume;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
@ -38,6 +38,7 @@ import org.apache.cassandra.db.RowUpdateBuilder;
import org.apache.cassandra.db.compaction.CompactionManager; import org.apache.cassandra.db.compaction.CompactionManager;
import org.apache.cassandra.db.marshal.AsciiType; import org.apache.cassandra.db.marshal.AsciiType;
import org.apache.cassandra.db.marshal.BytesType; import org.apache.cassandra.db.marshal.BytesType;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.schema.KeyspaceParams; import org.apache.cassandra.schema.KeyspaceParams;
import org.jboss.byteman.contrib.bmunit.BMRule; import org.jboss.byteman.contrib.bmunit.BMRule;
import org.jboss.byteman.contrib.bmunit.BMUnitRunner; import org.jboss.byteman.contrib.bmunit.BMUnitRunner;
@ -61,6 +62,10 @@ public class CommitLogChainedMarkersTest
{ {
// this method is blend of CommitLogSegmentBackpressureTest & CommitLogReaderTest methods // this method is blend of CommitLogSegmentBackpressureTest & CommitLogReaderTest methods
DatabaseDescriptor.daemonInitialization(); DatabaseDescriptor.daemonInitialization();
Assume.assumeTrue("With direct IO used for the commitlog we cannot expect the data are on disk without flushing it",
DatabaseDescriptor.getCommitLogWriteDiskAccessMode() != Config.DiskAccessMode.direct);
DatabaseDescriptor.setCommitLogSegmentSize(5); DatabaseDescriptor.setCommitLogSegmentSize(5);
DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic);
DatabaseDescriptor.setCommitLogSyncPeriod(10000 * 1000); DatabaseDescriptor.setCommitLogSyncPeriod(10000 * 1000);

View File

@ -85,6 +85,7 @@ public class CommitLogSegmentBackpressureTest
DatabaseDescriptor.setCommitLogSync(CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSync(CommitLogSync.periodic);
DatabaseDescriptor.setCommitLogSyncPeriod(10 * 1000); DatabaseDescriptor.setCommitLogSyncPeriod(10 * 1000);
DatabaseDescriptor.setCommitLogMaxCompressionBuffersPerPool(3); DatabaseDescriptor.setCommitLogMaxCompressionBuffersPerPool(3);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
SchemaLoader.prepareServer(); SchemaLoader.prepareServer();
SchemaLoader.createKeyspace(KEYSPACE1, SchemaLoader.createKeyspace(KEYSPACE1,
KeyspaceParams.simple(1), KeyspaceParams.simple(1),

View File

@ -127,6 +127,7 @@ public abstract class CommitLogTest
{ {
DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setCommitLogCompression(commitLogCompression);
DatabaseDescriptor.setEncryptionContext(encryptionContext); DatabaseDescriptor.setEncryptionContext(encryptionContext);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
} }
@Parameters() @Parameters()

View File

@ -67,6 +67,7 @@ public class CommitlogShutdownTest
DatabaseDescriptor.setCommitLogSegmentSize(1); DatabaseDescriptor.setCommitLogSegmentSize(1);
DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic);
DatabaseDescriptor.setCommitLogSyncPeriod(10 * 1000); DatabaseDescriptor.setCommitLogSyncPeriod(10 * 1000);
DatabaseDescriptor.initializeCommitLogDiskAccessMode();
SchemaLoader.prepareServer(); SchemaLoader.prepareServer();
SchemaLoader.createKeyspace(KEYSPACE1, SchemaLoader.createKeyspace(KEYSPACE1,
KeyspaceParams.simple(1), KeyspaceParams.simple(1),

View File

@ -0,0 +1,173 @@
/*
* 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.db.commitlog;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.concurrent.atomic.AtomicLong;
import org.junit.BeforeClass;
import org.junit.Test;
import net.openhft.chronicle.core.util.ThrowingFunction;
import org.apache.cassandra.config.Config;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.SimpleCachedBufferPool;
import org.apache.cassandra.utils.Generators;
import org.mockito.ArgumentCaptor;
import org.mockito.internal.creation.MockSettingsImpl;
import sun.nio.ch.DirectBuffer;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doCallRealMethod;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.quicktheories.QuickTheory.qt;
public class DirectIOSegmentTest
{
@BeforeClass
public static void beforeClass() throws IOException
{
File commitLogDir = new File(Files.createTempDirectory("commitLogDir"));
DatabaseDescriptor.daemonInitialization(() -> {
Config config = DatabaseDescriptor.loadConfig();
config.commitlog_directory = commitLogDir.toString();
return config;
});
}
@Test
public void testFlushBuffer()
{
int fsBlockSize = 32;
int bufSize = 4 * fsBlockSize;
SimpleCachedBufferPool bufferPool = mock(SimpleCachedBufferPool.class);
AbstractCommitLogSegmentManager manager = mock(AbstractCommitLogSegmentManager.class,
new MockSettingsImpl<>().useConstructor(CommitLog.instance, DatabaseDescriptor.getCommitLogLocation()));
doReturn(bufferPool).when(manager).getBufferPool();
doCallRealMethod().when(manager).getConfiguration();
when(bufferPool.createBuffer()).thenReturn(ByteBuffer.allocate(bufSize + fsBlockSize));
doNothing().when(manager).addSize(anyLong());
qt().forAll(Generators.forwardRanges(0, bufSize))
.checkAssert(startEnd -> {
int start = startEnd.lowerEndpoint();
int end = startEnd.upperEndpoint();
FileChannel channel = mock(FileChannel.class);
ThrowingFunction<Path, FileChannel, IOException> channelFactory = path -> channel;
ArgumentCaptor<ByteBuffer> bufCap = ArgumentCaptor.forClass(ByteBuffer.class);
DirectIOSegment seg = new DirectIOSegment(manager, channelFactory, fsBlockSize);
seg.lastSyncedOffset = start;
seg.flush(start, end);
try
{
verify(channel).write(bufCap.capture());
}
catch (IOException e)
{
throw new RuntimeException(e);
}
ByteBuffer buf = bufCap.getValue();
// assert that the entire buffer is written
assertThat(buf.position()).isLessThanOrEqualTo(start);
assertThat(buf.limit()).isGreaterThanOrEqualTo(end);
// assert that the buffer is aligned to the fs block size
assertThat(buf.position() % fsBlockSize).isZero();
assertThat(buf.limit() % fsBlockSize).isZero();
// assert that the buffer is unnecessarily large
assertThat(buf.position()).isGreaterThan(start - fsBlockSize);
assertThat(buf.limit()).isLessThan(end + fsBlockSize);
assertThat(seg.lastWritten).isEqualTo(buf.limit());
});
}
@Test
public void testFlushSize()
{
int fsBlockSize = 32;
int bufSize = 4 * fsBlockSize;
SimpleCachedBufferPool bufferPool = mock(SimpleCachedBufferPool.class);
AbstractCommitLogSegmentManager manager = mock(AbstractCommitLogSegmentManager.class,
new MockSettingsImpl<>().useConstructor(CommitLog.instance, DatabaseDescriptor.getCommitLogLocation()));
doReturn(bufferPool).when(manager).getBufferPool();
doCallRealMethod().when(manager).getConfiguration();
when(bufferPool.createBuffer()).thenReturn(ByteBuffer.allocate(bufSize + fsBlockSize));
doNothing().when(manager).addSize(anyLong());
FileChannel channel = mock(FileChannel.class);
ThrowingFunction<Path, FileChannel, IOException> channelFactory = path -> channel;
ArgumentCaptor<ByteBuffer> bufCap = ArgumentCaptor.forClass(ByteBuffer.class);
DirectIOSegment seg = new DirectIOSegment(manager, channelFactory, fsBlockSize);
AtomicLong size = new AtomicLong();
doAnswer(i -> size.addAndGet(i.getArgument(0, Long.class))).when(manager).addSize(anyLong());
for (int start = 0; start < bufSize - 1; start++)
{
int end = start + 1;
seg.lastSyncedOffset = start;
seg.flush(start, end);
assertThat(size.get()).isGreaterThanOrEqualTo(end);
}
assertThat(size.get()).isEqualTo(bufSize);
}
@Test
public void testBuilder()
{
AbstractCommitLogSegmentManager manager = mock(AbstractCommitLogSegmentManager.class,
new MockSettingsImpl<>().useConstructor(CommitLog.instance, DatabaseDescriptor.getCommitLogLocation()));
DirectIOSegment.DirectIOSegmentBuilder builder = new DirectIOSegment.DirectIOSegmentBuilder(manager, 4096);
assertThat(builder.fsBlockSize).isGreaterThan(0);
int segmentSize = Math.max(5 << 20, builder.fsBlockSize * 5);
DatabaseDescriptor.setCommitLogSegmentSize(segmentSize >> 20);
SimpleCachedBufferPool pool = builder.createBufferPool();
ByteBuffer buf = pool.createBuffer();
try
{
assertThat(buf.remaining()).isEqualTo(segmentSize);
assertThat(buf.alignmentOffset(buf.position(), builder.fsBlockSize)).isEqualTo(0);
assertThat(buf).isInstanceOf(DirectBuffer.class);
assertThat(((DirectBuffer) buf).attachment()).isNotNull();
}
finally
{
pool.releaseBuffer(buf);
}
}
}

View File

@ -34,6 +34,7 @@ import java.util.UUID;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.Predicate; import java.util.function.Predicate;
import com.google.common.collect.Range;
import org.apache.commons.lang3.ArrayUtils; import org.apache.commons.lang3.ArrayUtils;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
@ -525,4 +526,11 @@ public final class Generators
throw new IllegalStateException("Gave up trying to find values matching assumptions after " + maxAttempts + " attempts"); throw new IllegalStateException("Gave up trying to find values matching assumptions after " + maxAttempts + " attempts");
} }
} }
public static Gen<Range<Integer>> forwardRanges(int min, int max)
{
return SourceDSL.integers().between(min, max)
.flatMap(start -> SourceDSL.integers().between(start, max)
.map(end -> Range.closed(start, end)));
}
} }