diff --git a/CHANGES.txt b/CHANGES.txt index 5a4f693998..06a4e36f18 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -4,6 +4,7 @@ * Add the ability to disable bulk loading of SSTables (CASSANDRA-18781) * Clean up obsolete functions and simplify cql_version handling in cqlsh (CASSANDRA-18787) Merged from 5.0: + * 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) * 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) diff --git a/NEWS.txt b/NEWS.txt index 37c3c93928..80e2f3219b 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -261,6 +261,9 @@ New features - Added snitch for Microsoft Azure of name AzureSnitch (CASSANDRA-18646) - 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 + - 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 --------- diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 1663357bee..508a177fe1 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -606,6 +606,16 @@ commitlog_segment_size: 32MiB # 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. # Note that tables without compression enabled do not respect this flag. # diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index a898ec5cc6..014c552a84 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -395,6 +395,7 @@ public class Config public ParameterizedClass commitlog_compression; public FlushCompression flush_compression = FlushCompression.fast; 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) public DurationSpec.IntMillisecondsBound periodic_commitlog_sync_lag_block; public TransparentDataEncryptionOptions transparent_data_encryption_options = new TransparentDataEncryptionOptions(); @@ -1155,6 +1156,8 @@ public class Config mmap, mmap_index_only, standard, + legacy, + direct // Direct-I/O is enabled for commitlog disk only. } public enum MemtableAllocationType diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 9d0cdc61e8..b304100957 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -75,6 +75,7 @@ import org.apache.cassandra.auth.IInternodeAuthenticator; import org.apache.cassandra.auth.INetworkAuthorizer; import org.apache.cassandra.auth.IRoleManager; 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.PaxosStatePurging; import org.apache.cassandra.db.ConsistencyLevel; @@ -139,7 +140,12 @@ import static org.apache.cassandra.config.CassandraRelevantProperties.UNSAFE_SYS import static org.apache.cassandra.config.DataRateSpec.DataRateUnit.BYTES_PER_SECOND; import static org.apache.cassandra.config.DataRateSpec.DataRateUnit.MEBIBYTES_PER_SECOND; import static org.apache.cassandra.config.DataStorageSpec.DataStorageUnit.MEBIBYTES; -import static org.apache.cassandra.db.ConsistencyLevel.*; +import static org.apache.cassandra.db.ConsistencyLevel.ALL; +import static org.apache.cassandra.db.ConsistencyLevel.EACH_QUORUM; +import static org.apache.cassandra.db.ConsistencyLevel.LOCAL_QUORUM; +import static org.apache.cassandra.db.ConsistencyLevel.NODE_LOCAL; +import static org.apache.cassandra.db.ConsistencyLevel.ONE; +import static org.apache.cassandra.db.ConsistencyLevel.QUORUM; import static org.apache.cassandra.io.util.FileUtils.ONE_GIB; import static org.apache.cassandra.io.util.FileUtils.ONE_MIB; import static org.apache.cassandra.utils.Clock.Global.logInitializationOutcome; @@ -181,7 +187,9 @@ public class DatabaseDescriptor private static IPartitioner partitioner; private static String paritionerName; - private static Config.DiskAccessMode indexAccessMode; + private static DiskAccessMode indexAccessMode; + + private static DiskAccessMode commitLogWriteDiskAccessMode; private static AbstractCryptoProvider cryptoProvider; private static IAuthenticator authenticator; @@ -517,23 +525,25 @@ public class DatabaseDescriptor } /* 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; - indexAccessMode = conf.disk_access_mode; - logger.info("DiskAccessMode 'auto' determined to be {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode); + conf.disk_access_mode = DiskAccessMode.standard; + indexAccessMode = DiskAccessMode.mmap; } - 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; - indexAccessMode = Config.DiskAccessMode.mmap; - logger.info("DiskAccessMode is {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode); + conf.disk_access_mode = hasLargeAddressSpace() ? DiskAccessMode.mmap : DiskAccessMode.standard; + indexAccessMode = conf.disk_access_mode; + } + else if (conf.disk_access_mode == DiskAccessMode.direct) + { + throw new ConfigurationException(String.format("DiskAccessMode '%s' is not supported", DiskAccessMode.direct)); } else { 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 */ if (conf.phi_convict_threshold < 5 || conf.phi_convict_threshold > 16) @@ -622,6 +632,10 @@ public class DatabaseDescriptor 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) { conf.hints_directory = storagedirFor("hints"); @@ -1434,6 +1448,52 @@ public class DatabaseDescriptor 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 factories) { Map factoryByName = new HashMap<>(); @@ -2578,6 +2638,7 @@ public class DatabaseDescriptor return conf.commitlog_compression; } + @VisibleForTesting public static void setCommitLogCompression(ParameterizedClass compressor) { conf.commitlog_compression = compressor; @@ -2674,6 +2735,28 @@ public class DatabaseDescriptor 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() { return conf.saved_caches_directory; @@ -3173,26 +3256,26 @@ public class DatabaseDescriptor conf.commitlog_sync = sync; } - public static Config.DiskAccessMode getDiskAccessMode() + public static DiskAccessMode getDiskAccessMode() { return conf.disk_access_mode; } // Do not use outside unit tests. @VisibleForTesting - public static void setDiskAccessMode(Config.DiskAccessMode mode) + public static void setDiskAccessMode(DiskAccessMode mode) { conf.disk_access_mode = mode; } - public static Config.DiskAccessMode getIndexAccessMode() + public static DiskAccessMode getIndexAccessMode() { return indexAccessMode; } // Do not use outside unit tests. @VisibleForTesting - public static void setIndexAccessMode(Config.DiskAccessMode mode) + public static void setIndexAccessMode(DiskAccessMode mode) { indexAccessMode = mode; } diff --git a/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java b/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java index e5e82f9662..dcd791caf3 100644 --- a/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java +++ b/src/java/org/apache/cassandra/db/commitlog/AbstractCommitLogSegmentManager.java @@ -18,7 +18,12 @@ package org.apache.cassandra.db.commitlog; 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.TimeUnit; import java.util.concurrent.atomic.AtomicLong; @@ -32,11 +37,11 @@ import com.codahale.metrics.Timer.Context; import net.nicoulaj.compilecommand.annotations.DontInline; import org.apache.cassandra.concurrent.Interruptible; import org.apache.cassandra.concurrent.Interruptible.TerminateException; +import org.apache.cassandra.config.Config.DiskAccessMode; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.Keyspace; 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.FileUtils; 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.TableMetadata; 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.InfiniteLoopExecutor.Daemon.NON_DAEMON; @@ -95,10 +104,12 @@ public abstract class AbstractCommitLogSegmentManager @VisibleForTesting Interruptible executor; - protected final CommitLog commitLog; + private final CommitLog commitLog; private final BooleanSupplier managerThreadWaitCondition = () -> (availableSegment == null && !atSegmentBufferLimit()); private final WaitQueue managerThreadWaitQueue = newWaitQueue(); + private volatile CommitLogSegment.Builder segmentBuilder; + private volatile SimpleCachedBufferPool bufferPool; AbstractCommitLogSegmentManager(final CommitLog commitLog, String storageDirectory) @@ -107,18 +118,41 @@ public abstract class AbstractCommitLogSegmentManager 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() { - // For encrypted segments 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 - BufferType bufferType = commitLog.configuration.useEncryption() || !commitLog.configuration.useCompression() - ? BufferType.ON_HEAP - : commitLog.configuration.getCompressor().preferredBufferType(); - - this.bufferPool = new SimpleCachedBufferPool(DatabaseDescriptor.getCommitLogMaxCompressionBuffersInPool(), - DatabaseDescriptor.getCommitLogSegmentSize(), - bufferType); - + assert this.segmentBuilder == null; + assert this.bufferPool == null; + this.segmentBuilder = createSegmentBuilder(commitLog.configuration); + this.bufferPool = segmentBuilder.createBufferPool(); AllocatorRunnable allocator = new AllocatorRunnable(); executor = executorFactory().infiniteLoop("COMMIT-LOG-ALLOCATOR", allocator, SAFE, NON_DAEMON, SYNCHRONIZED); @@ -206,7 +240,7 @@ public abstract class AbstractCommitLogSegmentManager private boolean atSegmentBufferLimit() { - return CommitLogSegment.usesBufferPool(commitLog) && bufferPool.atLimit(); + return bufferPool != null && bufferPool.atLimit(); } 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 * 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 @@ -557,6 +594,9 @@ public abstract class AbstractCommitLogSegmentManager if (bufferPool != null) bufferPool.emptyBufferPool(); + this.segmentBuilder = null; + this.bufferPool = null; + return res; } diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLog.java b/src/java/org/apache/cassandra/db/commitlog/CommitLog.java index ca5077edd6..5a54bb5b98 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLog.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLog.java @@ -40,6 +40,7 @@ import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.Config; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.config.ParameterizedClass; import org.apache.cassandra.db.Mutation; @@ -104,7 +105,8 @@ public class CommitLog implements CommitLogMBean CommitLog(CommitLogArchiver archiver, Function segmentManagerProvider) { this.configuration = new Configuration(DatabaseDescriptor.getCommitLogCompression(), - DatabaseDescriptor.getEncryptionContext()); + DatabaseDescriptor.getEncryptionContext(), + DatabaseDescriptor.getCommitLogWriteDiskAccessMode()); DatabaseDescriptor.createAllDirectories(); this.archiver = archiver; @@ -521,7 +523,8 @@ public class CommitLog implements CommitLogMBean synchronized public void resetConfiguration() { 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 { + /** + * Flag used to shows user configured Direct-IO status. + */ + public final Config.DiskAccessMode diskAccessMode; + /** * The compressor class. */ @@ -622,17 +630,18 @@ public class CommitLog implements CommitLogMBean /** * 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.compressor = compressorClass != null ? CompressionParams.createCompressor(compressorClass) : null; this.encryptionContext = encryptionContext; + this.diskAccessMode = diskAccessMode; } /** - * Checks if the segments must be compressed. * @return true if the segments must be compressed, false otherwise. */ public boolean useCompression() @@ -641,7 +650,6 @@ public class CommitLog implements CommitLogMBean } /** - * Checks if the segments must be encrypted. * @return true if the segments must be encrypted, false otherwise. */ 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 */ public ICompressor getCompressor() @@ -659,7 +666,6 @@ public class CommitLog implements CommitLogMBean } /** - * Returns the compressor class. * @return the compressor class */ public ParameterizedClass getCompressorClass() @@ -668,7 +674,6 @@ public class CommitLog implements CommitLogMBean } /** - * Returns the compressor name. * @return the compressor name. */ 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 */ public EncryptionContext getEncryptionContext() { 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; + } } } diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java index 7a346acfd6..2a79ac4519 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegment.java @@ -20,8 +20,14 @@ package org.apache.cassandra.db.commitlog; import java.io.IOException; import java.nio.ByteBuffer; import java.nio.channels.FileChannel; -import java.nio.file.StandardOpenOption; -import java.util.*; +import java.nio.file.Path; +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.ConcurrentMap; import java.util.concurrent.atomic.AtomicInteger; @@ -29,21 +35,21 @@ import java.util.concurrent.locks.LockSupport; import java.util.zip.CRC32; 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 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.commitlog.CommitLog.Configuration; import org.apache.cassandra.db.partitions.PartitionUpdate; 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.FileWriter; +import org.apache.cassandra.io.util.SimpleCachedBufferPool; import org.apache.cassandra.schema.Schema; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.schema.TableMetadata; -import org.apache.cassandra.utils.NativeLibrary; import org.apache.cassandra.utils.IntegerInterval; import org.apache.cassandra.utils.concurrent.OpOrder; import org.apache.cassandra.utils.concurrent.WaitQueue; @@ -124,7 +130,6 @@ public abstract class CommitLogSegment final File logFile; final FileChannel channel; - final int fd; protected final AbstractCommitLogSegmentManager manager; @@ -133,28 +138,6 @@ public abstract class CommitLogSegment 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 true if the segments use a buffer pool, false otherwise. - */ - static boolean usesBufferPool(CommitLog commitLog) - { - Configuration config = commitLog.configuration; - return config.useEncryption() || config.useCompression(); - } - static long getNextId() { return idBase + nextId.getAndIncrement(); @@ -163,27 +146,26 @@ public abstract class CommitLogSegment /** * Constructs a new segment file. */ - CommitLogSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) + CommitLogSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction channelFactory) { this.manager = manager; id = getNextId(); descriptor = new CommitLogDescriptor(id, - commitLog.configuration.getCompressorClass(), - commitLog.configuration.getEncryptionContext()); + manager.getConfiguration().getCompressorClass(), + manager.getConfiguration().getEncryptionContext()); logFile = new File(manager.storageDirectory, descriptor.fileName()); try { - channel = FileChannel.open(logFile.toPath(), StandardOpenOption.WRITE, StandardOpenOption.READ, StandardOpenOption.CREATE); - fd = NativeLibrary.getfd(channel); + channel = channelFactory.apply(logFile.toPath()); } catch (IOException e) { throw new FSWriteError(e, logFile); } - buffer = createBuffer(commitLog); + this.buffer = createBuffer(); } /** @@ -207,7 +189,10 @@ public abstract class CommitLogSegment return Collections.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. @@ -764,4 +749,18 @@ public abstract class CommitLogSegment 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(); + } } diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerCDC.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerCDC.java index 349986de22..7dfe7add8f 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerCDC.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerCDC.java @@ -236,7 +236,8 @@ public class CommitLogSegmentManagerCDC extends AbstractCommitLogSegmentManager @Override public CommitLogSegment createSegment() { - CommitLogSegment segment = CommitLogSegment.createSegment(commitLog, this); + CommitLogSegment segment = super.createSegment(); + segment.writeLogHeader(); cdcSizeTracker.processNewSegment(segment); // After processing, the state of the segment can either be PERMITTED or FORBIDDEN if (segment.getCDCState() == CDCState.PERMITTED) diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerStandard.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerStandard.java index 1ae2f134b3..6ca662a3db 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerStandard.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLogSegmentManagerStandard.java @@ -62,6 +62,8 @@ public class CommitLogSegmentManagerStandard extends AbstractCommitLogSegmentMan @Override public CommitLogSegment createSegment() { - return CommitLogSegment.createSegment(commitLog, this); + CommitLogSegment segment = super.createSegment(); + segment.writeLogHeader(); + return segment; } } diff --git a/src/java/org/apache/cassandra/db/commitlog/CompressedSegment.java b/src/java/org/apache/cassandra/db/commitlog/CompressedSegment.java index 1d65fbec4e..3c7f360836 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CompressedSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/CompressedSegment.java @@ -17,10 +17,17 @@ */ 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 net.openhft.chronicle.core.util.ThrowingFunction; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.io.FSWriteError; 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 @@ -41,15 +48,10 @@ public class CompressedSegment extends FileDirectSegment /** * Constructs a new segment file. */ - CompressedSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) + CompressedSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction channelFactory) { - super(commitLog, manager); - this.compressor = commitLog.configuration.getCompressor(); - } - - ByteBuffer createBuffer(CommitLog commitLog) - { - return manager.getBufferPool().createBuffer(); + super(manager, channelFactory); + this.compressor = manager.getConfiguration().getCompressor(); } @Override @@ -92,4 +94,27 @@ public class CompressedSegment extends FileDirectSegment { 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()); + } + } } diff --git a/src/java/org/apache/cassandra/db/commitlog/DirectIOSegment.java b/src/java/org/apache/cassandra/db/commitlog/DirectIOSegment.java new file mode 100644 index 0000000000..2f350f71fb --- /dev/null +++ b/src/java/org/apache/cassandra/db/commitlog/DirectIOSegment.java @@ -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 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); + } + }; + } + } +} diff --git a/src/java/org/apache/cassandra/db/commitlog/EncryptedSegment.java b/src/java/org/apache/cassandra/db/commitlog/EncryptedSegment.java index f5036583a7..c7a8324efe 100644 --- a/src/java/org/apache/cassandra/db/commitlog/EncryptedSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/EncryptedSegment.java @@ -19,17 +19,23 @@ 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 java.util.Map; import javax.crypto.Cipher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +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.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.EncryptionUtils; import org.apache.cassandra.utils.Hex; 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 Cipher cipher; - public EncryptedSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) + public EncryptedSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction channelFactory) { - super(commitLog, manager); - this.encryptionContext = commitLog.configuration.getEncryptionContext(); + super(manager, channelFactory); + this.encryptionContext = manager.getConfiguration().getEncryptionContext(); try { @@ -86,12 +92,9 @@ public class EncryptedSegment extends FileDirectSegment 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 - return manager.getBufferPool().createBuffer(); - } + // 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 + // (so we do not override the createBuffer method) void write(int startMarker, int nextMarker) { @@ -150,4 +153,28 @@ public class EncryptedSegment extends FileDirectSegment { 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); + } + } } diff --git a/src/java/org/apache/cassandra/db/commitlog/FileDirectSegment.java b/src/java/org/apache/cassandra/db/commitlog/FileDirectSegment.java index d5431f875b..16e32ca448 100644 --- a/src/java/org/apache/cassandra/db/commitlog/FileDirectSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/FileDirectSegment.java @@ -19,7 +19,10 @@ 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 net.openhft.chronicle.core.util.ThrowingFunction; import org.apache.cassandra.io.FSWriteError; import org.apache.cassandra.utils.SyncUtil; @@ -31,9 +34,9 @@ public abstract class FileDirectSegment extends CommitLogSegment { volatile long lastWrittenPos = 0; - FileDirectSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) + FileDirectSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction channelFactory) { - super(commitLog, manager); + super(manager, channelFactory); } @Override diff --git a/src/java/org/apache/cassandra/db/commitlog/MemoryMappedSegment.java b/src/java/org/apache/cassandra/db/commitlog/MemoryMappedSegment.java index d564117d3a..fb671130a1 100644 --- a/src/java/org/apache/cassandra/db/commitlog/MemoryMappedSegment.java +++ b/src/java/org/apache/cassandra/db/commitlog/MemoryMappedSegment.java @@ -21,10 +21,14 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.nio.MappedByteBuffer; 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.util.FileUtils; +import org.apache.cassandra.io.util.SimpleCachedBufferPool; import org.apache.cassandra.utils.NativeLibrary; import org.apache.cassandra.utils.SyncUtil; @@ -35,21 +39,23 @@ import org.apache.cassandra.utils.SyncUtil; */ public class MemoryMappedSegment extends CommitLogSegment { + private final int fd; + /** * Constructs a new segment file. - * - * @param commitLog the commit log it will be used with. */ - MemoryMappedSegment(CommitLog commitLog, AbstractCommitLogSegmentManager manager) + MemoryMappedSegment(AbstractCommitLogSegmentManager manager, ThrowingFunction channelFactory) { - super(commitLog, manager); + super(manager, channelFactory); // mark the initial sync marker as uninitialised int firstSync = buffer.position(); buffer.putInt(firstSync + 0, 0); buffer.putInt(firstSync + 4, 0); + fd = NativeLibrary.getfd(channel); } - ByteBuffer createBuffer(CommitLog commitLog) + @Override + protected ByteBuffer createBuffer() { try { @@ -105,4 +111,25 @@ public class MemoryMappedSegment extends CommitLogSegment FileUtils.clean(buffer); 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; + } + } } diff --git a/src/java/org/apache/cassandra/io/util/FileUtils.java b/src/java/org/apache/cassandra/io/util/FileUtils.java index 665b6e1a7e..7027d6e114 100644 --- a/src/java/org/apache/cassandra/io/util/FileUtils.java +++ b/src/java/org/apache/cassandra/io/util/FileUtils.java @@ -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(); + } + } } \ No newline at end of file diff --git a/src/java/org/apache/cassandra/io/util/PathUtils.java b/src/java/org/apache/cassandra/io/util/PathUtils.java index 88e2d43b25..8ddd939b4c 100644 --- a/src/java/org/apache/cassandra/io/util/PathUtils.java +++ b/src/java/org/apache/cassandra/io/util/PathUtils.java @@ -49,21 +49,20 @@ import java.util.function.Function; import java.util.function.IntFunction; import java.util.stream.Collectors; import java.util.stream.Stream; - import javax.annotation.Nullable; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; 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.io.FSError; import org.apache.cassandra.io.FSReadError; import org.apache.cassandra.io.FSWriteError; -import org.apache.cassandra.service.StorageService; 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.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.WRITE; import static java.util.Collections.unmodifiableSet; - import static org.apache.cassandra.config.CassandraRelevantProperties.USE_NIX_RECURSIVE_DELETE; 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 * @@ -98,12 +94,7 @@ public final class PathUtils private static final Logger logger = LoggerFactory.getLogger(PathUtils.class); private static final NoSpamLogger nospam1m = NoSpamLogger.getLogger(logger, 1, TimeUnit.MINUTES); - private static Consumer onDeletion = path -> { - if (StorageService.instance.isDaemonSetupCompleted()) - setDeletionListener(ignore -> {}); - else - logger.trace("Deleting file during startup: {}", path); - }; + private static Consumer onDeletion = path -> {}; public static FileChannel newReadChannel(Path path) throws NoSuchFileException { diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 865ecb14b2..e491cf0777 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -271,6 +271,15 @@ public class StorageService extends NotificationBroadcasterSupport implements IE public static final int INDEFINITE = -1; public static final int RING_DELAY_MILLIS = getRingDelay(); // delay after which we assume ring has stablized + { + 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 AtomicReference ongoingBootstrap = new AtomicReference<>(); diff --git a/test/conf/cassandra.yaml b/test/conf/cassandra.yaml index bfb58b3e65..47fbde0489 100644 --- a/test/conf/cassandra.yaml +++ b/test/conf/cassandra.yaml @@ -8,6 +8,7 @@ memtable_allocation_type: offheap_objects commitlog_sync: batch commitlog_segment_size: 5MiB commitlog_directory: build/test/cassandra/commitlog +commitlog_disk_access_mode: legacy # commitlog_compression: # - class_name: LZ4Compressor cdc_raw_directory: build/test/cassandra/cdc_raw diff --git a/test/distributed/org/apache/cassandra/distributed/impl/InstanceConfig.java b/test/distributed/org/apache/cassandra/distributed/impl/InstanceConfig.java index 6f0238bd2f..fe3ab34f42 100644 --- a/test/distributed/org/apache/cassandra/distributed/impl/InstanceConfig.java +++ b/test/distributed/org/apache/cassandra/distributed/impl/InstanceConfig.java @@ -112,7 +112,8 @@ public class InstanceConfig implements IInstanceConfig // capacities that are based on `totalMemory` that should be fixed size .set("index_summary_capacity", "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.jmxPort = jmx_port; } diff --git a/test/long/org/apache/cassandra/db/commitlog/BatchCommitLogStressTest.java b/test/long/org/apache/cassandra/db/commitlog/BatchCommitLogStressTest.java index 3665882bf3..2979ae31e9 100644 --- a/test/long/org/apache/cassandra/db/commitlog/BatchCommitLogStressTest.java +++ b/test/long/org/apache/cassandra/db/commitlog/BatchCommitLogStressTest.java @@ -29,9 +29,9 @@ import org.apache.cassandra.security.EncryptionContext; @RunWith(Parameterized.class) 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); } } diff --git a/test/long/org/apache/cassandra/db/commitlog/CommitLogStressTest.java b/test/long/org/apache/cassandra/db/commitlog/CommitLogStressTest.java index 575450a93f..cf2e384769 100644 --- a/test/long/org/apache/cassandra/db/commitlog/CommitLogStressTest.java +++ b/test/long/org/apache/cassandra/db/commitlog/CommitLogStressTest.java @@ -93,12 +93,15 @@ public abstract class CommitLogStressTest private boolean randomSize = false; private boolean discardedRun = false; 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.setEncryptionContext(encryptionContext); DatabaseDescriptor.setCommitLogSegmentSize(32); + DatabaseDescriptor.setCommitLogWriteDiskAccessMode(accessMode); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } @BeforeClass @@ -142,11 +145,12 @@ public abstract class CommitLogStressTest public static Collection buildParameterizedVariants() { return Arrays.asList(new Object[][]{ - {null, EncryptionContextGenerator.createDisabledContext()}, // No compression, no encryption - {null, EncryptionContextGenerator.createContext(true)}, // Encryption - { new ParameterizedClass(LZ4Compressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext()}, - { new ParameterizedClass(SnappyCompressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext()}, - { new ParameterizedClass(DeflateCompressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext()}}); + {null, EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.legacy}, // No compression, no encryption, legacy + {null, EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.direct}, // Use Direct-I/O (non-buffered) feature. + {null, EncryptionContextGenerator.createContext(true), Config.DiskAccessMode.legacy}, + { new ParameterizedClass(LZ4Compressor.class.getName(), Collections.emptyMap()), EncryptionContextGenerator.createDisabledContext(), Config.DiskAccessMode.legacy}, + { 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 @@ -157,6 +161,7 @@ public abstract class CommitLogStressTest testLog(); } + @Test public void testFixedSize() throws Exception { @@ -179,6 +184,7 @@ public abstract class CommitLogStressTest try { DatabaseDescriptor.setCommitLogLocation(location); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); CommitLog commitLog = new CommitLog(CommitLogArchiver.disabled()).start(); testLog(commitLog); assert !failed; @@ -186,14 +192,17 @@ public abstract class CommitLogStressTest finally { DatabaseDescriptor.setCommitLogLocation(originalDir); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } } 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()), + DatabaseDescriptor.getCommitLogWriteDiskAccessMode(), commitLog.configuration.getCompressorName(), commitLog.configuration.useEncryption(), + commitLog.configuration.isDirectIOEnabled(), commitLog.executor.getClass().getSimpleName(), randomSize ? " random size" : "", discardedRun ? " with discarded run" : ""); @@ -258,15 +267,20 @@ public abstract class CommitLogStressTest Assert.fail("Failed to delete " + f); 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.useEncryption(), - reader.discarded, reader.skipped); + commitLog.configuration.isDirectIOEnabled(), + reader.discarded, reader.skipped, + mb(totalBytesWritten), mb(totalBytesWritten)/runTimeMs*1000); 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.useEncryption(), + commitLog.configuration.isDirectIOEnabled(), reader.cells, cells, cells - reader.cells, reader.discarded, reader.skipped, reader.hash, hash); failed = true; @@ -282,7 +296,10 @@ public abstract class CommitLogStressTest long combinedSize = 0; for (File f : new File(commitLog.segmentManager.storageDirectory).tryList()) + { combinedSize += f.length(); + } + totalBytesWritten = combinedSize; Assert.assertEquals(combinedSize, commitLog.getActiveOnDiskSize()); List logFileNames = commitLog.getActiveSegmentNames(); diff --git a/test/long/org/apache/cassandra/db/commitlog/GroupCommitLogStressTest.java b/test/long/org/apache/cassandra/db/commitlog/GroupCommitLogStressTest.java index e3fa961c4e..90f801b472 100644 --- a/test/long/org/apache/cassandra/db/commitlog/GroupCommitLogStressTest.java +++ b/test/long/org/apache/cassandra/db/commitlog/GroupCommitLogStressTest.java @@ -29,9 +29,9 @@ import org.apache.cassandra.security.EncryptionContext; @RunWith(Parameterized.class) 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.setCommitLogSyncGroupWindow(1); } diff --git a/test/long/org/apache/cassandra/db/commitlog/PeriodicCommitLogStressTest.java b/test/long/org/apache/cassandra/db/commitlog/PeriodicCommitLogStressTest.java index 509d46ab55..e9ea93e615 100644 --- a/test/long/org/apache/cassandra/db/commitlog/PeriodicCommitLogStressTest.java +++ b/test/long/org/apache/cassandra/db/commitlog/PeriodicCommitLogStressTest.java @@ -29,9 +29,9 @@ import org.apache.cassandra.security.EncryptionContext; @RunWith(Parameterized.class) 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.setCommitLogSyncPeriod(30); } diff --git a/test/unit/org/apache/cassandra/config/DatabaseDescriptorTest.java b/test/unit/org/apache/cassandra/config/DatabaseDescriptorTest.java index 6a83dcd4e6..bc51e73601 100644 --- a/test/unit/org/apache/cassandra/config/DatabaseDescriptorTest.java +++ b/test/unit/org/apache/cassandra/config/DatabaseDescriptorTest.java @@ -18,16 +18,18 @@ */ package org.apache.cassandra.config; +import java.io.IOException; import java.net.Inet4Address; import java.net.Inet6Address; import java.net.InetAddress; import java.net.NetworkInterface; +import java.nio.file.Files; import java.util.Arrays; import java.util.Collection; +import java.util.EnumSet; import java.util.Enumeration; import java.util.function.Consumer; - import com.google.common.base.Throwables; import org.junit.Assert; import org.junit.BeforeClass; @@ -36,6 +38,8 @@ import org.junit.Test; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.distributed.shared.WithProperties; 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 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.DataStorageSpec.DataStorageUnit.KIBIBYTES; 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.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; @@ -805,4 +810,106 @@ public class DatabaseDescriptorTest { 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 allowedModes = EnumSet.copyOf(Arrays.asList(allowedModesArray)); + allowedModes.add(Config.DiskAccessMode.legacy); + allowedModes.add(Config.DiskAccessMode.auto); + + EnumSet 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); + } + } } diff --git a/test/unit/org/apache/cassandra/db/RecoveryManagerFlushedTest.java b/test/unit/org/apache/cassandra/db/RecoveryManagerFlushedTest.java index c9c8ac158b..4db364f24e 100644 --- a/test/unit/org/apache/cassandra/db/RecoveryManagerFlushedTest.java +++ b/test/unit/org/apache/cassandra/db/RecoveryManagerFlushedTest.java @@ -62,6 +62,7 @@ public class RecoveryManagerFlushedTest { DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setEncryptionContext(encryptionContext); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } @Parameters() diff --git a/test/unit/org/apache/cassandra/db/RecoveryManagerMissingHeaderTest.java b/test/unit/org/apache/cassandra/db/RecoveryManagerMissingHeaderTest.java index 21058671b8..e8444fba4f 100644 --- a/test/unit/org/apache/cassandra/db/RecoveryManagerMissingHeaderTest.java +++ b/test/unit/org/apache/cassandra/db/RecoveryManagerMissingHeaderTest.java @@ -61,6 +61,7 @@ public class RecoveryManagerMissingHeaderTest { DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setEncryptionContext(encryptionContext); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } @Parameters() diff --git a/test/unit/org/apache/cassandra/db/RecoveryManagerTest.java b/test/unit/org/apache/cassandra/db/RecoveryManagerTest.java index c60bb97aa4..2218619c72 100644 --- a/test/unit/org/apache/cassandra/db/RecoveryManagerTest.java +++ b/test/unit/org/apache/cassandra/db/RecoveryManagerTest.java @@ -78,6 +78,7 @@ public class RecoveryManagerTest { DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setEncryptionContext(encryptionContext); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } @Parameters() diff --git a/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java b/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java index a51cd21510..6270b25eb7 100644 --- a/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java +++ b/test/unit/org/apache/cassandra/db/RecoveryManagerTruncateTest.java @@ -59,6 +59,7 @@ public class RecoveryManagerTruncateTest { DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setEncryptionContext(encryptionContext); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } @Parameters() diff --git a/test/unit/org/apache/cassandra/db/commitlog/CommitLogChainedMarkersTest.java b/test/unit/org/apache/cassandra/db/commitlog/CommitLogChainedMarkersTest.java index 319e75b346..b5e417b6e4 100644 --- a/test/unit/org/apache/cassandra/db/commitlog/CommitLogChainedMarkersTest.java +++ b/test/unit/org/apache/cassandra/db/commitlog/CommitLogChainedMarkersTest.java @@ -23,8 +23,8 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Random; -import org.apache.cassandra.io.util.File; import org.junit.Assert; +import org.junit.Assume; import org.junit.Test; 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.marshal.AsciiType; import org.apache.cassandra.db.marshal.BytesType; +import org.apache.cassandra.io.util.File; import org.apache.cassandra.schema.KeyspaceParams; import org.jboss.byteman.contrib.bmunit.BMRule; import org.jboss.byteman.contrib.bmunit.BMUnitRunner; @@ -61,6 +62,10 @@ public class CommitLogChainedMarkersTest { // this method is blend of CommitLogSegmentBackpressureTest & CommitLogReaderTest methods 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.setCommitLogSync(Config.CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSyncPeriod(10000 * 1000); diff --git a/test/unit/org/apache/cassandra/db/commitlog/CommitLogSegmentBackpressureTest.java b/test/unit/org/apache/cassandra/db/commitlog/CommitLogSegmentBackpressureTest.java index e28c25e2bc..6a8521485a 100644 --- a/test/unit/org/apache/cassandra/db/commitlog/CommitLogSegmentBackpressureTest.java +++ b/test/unit/org/apache/cassandra/db/commitlog/CommitLogSegmentBackpressureTest.java @@ -85,6 +85,7 @@ public class CommitLogSegmentBackpressureTest DatabaseDescriptor.setCommitLogSync(CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSyncPeriod(10 * 1000); DatabaseDescriptor.setCommitLogMaxCompressionBuffersPerPool(3); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); SchemaLoader.prepareServer(); SchemaLoader.createKeyspace(KEYSPACE1, KeyspaceParams.simple(1), diff --git a/test/unit/org/apache/cassandra/db/commitlog/CommitLogTest.java b/test/unit/org/apache/cassandra/db/commitlog/CommitLogTest.java index 92aa8d81ab..79a32c0a8b 100644 --- a/test/unit/org/apache/cassandra/db/commitlog/CommitLogTest.java +++ b/test/unit/org/apache/cassandra/db/commitlog/CommitLogTest.java @@ -140,6 +140,7 @@ public abstract class CommitLogTest { DatabaseDescriptor.setCommitLogCompression(commitLogCompression); DatabaseDescriptor.setEncryptionContext(encryptionContext); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); } @Parameters() diff --git a/test/unit/org/apache/cassandra/db/commitlog/CommitlogShutdownTest.java b/test/unit/org/apache/cassandra/db/commitlog/CommitlogShutdownTest.java index e962450a80..9ec6efca53 100644 --- a/test/unit/org/apache/cassandra/db/commitlog/CommitlogShutdownTest.java +++ b/test/unit/org/apache/cassandra/db/commitlog/CommitlogShutdownTest.java @@ -67,6 +67,7 @@ public class CommitlogShutdownTest DatabaseDescriptor.setCommitLogSegmentSize(1); DatabaseDescriptor.setCommitLogSync(Config.CommitLogSync.periodic); DatabaseDescriptor.setCommitLogSyncPeriod(10 * 1000); + DatabaseDescriptor.initializeCommitLogDiskAccessMode(); SchemaLoader.prepareServer(); SchemaLoader.createKeyspace(KEYSPACE1, KeyspaceParams.simple(1), diff --git a/test/unit/org/apache/cassandra/db/commitlog/DirectIOSegmentTest.java b/test/unit/org/apache/cassandra/db/commitlog/DirectIOSegmentTest.java new file mode 100644 index 0000000000..27d3946b2a --- /dev/null +++ b/test/unit/org/apache/cassandra/db/commitlog/DirectIOSegmentTest.java @@ -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 channelFactory = path -> channel; + ArgumentCaptor 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 channelFactory = path -> channel; + ArgumentCaptor 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); + } + } +} \ No newline at end of file diff --git a/test/unit/org/apache/cassandra/utils/Generators.java b/test/unit/org/apache/cassandra/utils/Generators.java index 00d84a24ed..f23aff7498 100644 --- a/test/unit/org/apache/cassandra/utils/Generators.java +++ b/test/unit/org/apache/cassandra/utils/Generators.java @@ -34,6 +34,7 @@ import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.function.Predicate; +import com.google.common.collect.Range; import org.apache.commons.lang3.ArrayUtils; import org.slf4j.Logger; 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"); } } + + public static Gen> forwardRanges(int min, int max) + { + return SourceDSL.integers().between(min, max) + .flatMap(start -> SourceDSL.integers().between(start, max) + .map(end -> Range.closed(start, end))); + } }