Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.txt
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
7.0
* Extend background direct reads to SSTable streaming and SASI index builds (CASSANDRA-19988)
* Allow CQLSSTableWriter to specify SSTable id generator to use (CASSANDRA-21012)
* Reject LIKE patterns with a wildcard (%) anywhere other than the start or end (CASSANDRA-21068)
* Support pluggable default role initialization (CASSANDRA-21546)
Expand Down
3 changes: 2 additions & 1 deletion NEWS.txt
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,8 @@ Upgrading

Deprecation
-----------

- `compaction_read_disk_access_mode` has been renamed to `background_read_disk_access_mode`. The old name
is still accepted for compatibility. (see CASSANDRA-19988)

6.0
===
Expand Down
10 changes: 7 additions & 3 deletions conf/cassandra.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -711,10 +711,14 @@ commitlog_disk_access_mode: legacy
# Experimental. New garbage free compaction path. Does not yet support collections, counters, BTI.
cursor_compaction_enabled: false

# Set the disk access mode for reading SSTables during compaction. The allowed values are:
# Set the disk access mode for reading SSTables during background operations
# (compaction, scrub, verify, index build, streaming). The allowed values are:
# - auto: inherit from disk_access_mode (default)
# - direct: use direct I/O for compaction reads, bypassing the OS page cache
# compaction_read_disk_access_mode: auto
# - direct: use direct I/O for background reads, bypassing the OS page cache where supported.
# On macOS this is only a cache retention hint and tmpfs accepts it without effect; falls
# back to normal reads when direct I/O is unsupported.
# The former compaction_read_disk_access_mode name is accepted for compatibility.
# background_read_disk_access_mode: auto

# Set the disk access mode for writing compressed SSTables during background operations
# (compaction, streaming, cleanup, repair, etc.). The allowed values are:
Expand Down
10 changes: 7 additions & 3 deletions conf/cassandra_latest.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -718,10 +718,14 @@ commitlog_disk_access_mode: auto
# Experimental. New garbage free compaction path. Does not yet support collections, counters, BTI.
cursor_compaction_enabled: true

# Set the disk access mode for reading SSTables during compaction. The allowed values are:
# Set the disk access mode for reading SSTables during background operations
# (compaction, scrub, verify, index build, streaming). The allowed values are:
# - auto: inherit from disk_access_mode (default)
# - direct: use direct I/O for compaction reads, bypassing the OS page cache
# compaction_read_disk_access_mode: auto
# - direct: use direct I/O for background reads, bypassing the OS page cache where supported.
# On macOS this is only a cache retention hint and tmpfs accepts it without effect; falls
# back to normal reads when direct I/O is unsupported.
# The former compaction_read_disk_access_mode name is accepted for compatibility.
# background_read_disk_access_mode: auto

# Set the disk access mode for writing compressed SSTables during background operations
# (compaction, streaming, cleanup, repair, etc.). The allowed values are:
Expand Down
3 changes: 2 additions & 1 deletion src/java/org/apache/cassandra/config/Config.java
Original file line number Diff line number Diff line change
Expand Up @@ -485,7 +485,8 @@ public static class SSTableConfig
public FlushCompression flush_compression = FlushCompression.fast;
public int commitlog_max_compression_buffers_in_pool = 3;
public DiskAccessMode commitlog_disk_access_mode = DiskAccessMode.legacy;
public DiskAccessMode compaction_read_disk_access_mode = DiskAccessMode.auto;
@Replaces(oldName = "compaction_read_disk_access_mode", deprecated = true)
public DiskAccessMode background_read_disk_access_mode = DiskAccessMode.auto;
@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();
Expand Down
26 changes: 13 additions & 13 deletions src/java/org/apache/cassandra/config/DatabaseDescriptor.java
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ public class DatabaseDescriptor

private static DiskAccessMode commitLogWriteDiskAccessMode;

private static DiskAccessMode compactionReadDiskAccessMode;
private static DiskAccessMode backgroundReadDiskAccessMode;

private static DiskAccessMode backgroundWriteDiskAccessMode;

Expand Down Expand Up @@ -701,20 +701,20 @@ else if (conf.disk_access_mode == DiskAccessMode.direct)
}
logger.info("DiskAccessMode is {}, indexAccessMode is {}", conf.disk_access_mode, indexAccessMode);

if (DiskAccessMode.auto == conf.compaction_read_disk_access_mode)
if (DiskAccessMode.auto == conf.background_read_disk_access_mode)
{
compactionReadDiskAccessMode = conf.disk_access_mode;
backgroundReadDiskAccessMode = conf.disk_access_mode;
}
else if (DiskAccessMode.direct == conf.compaction_read_disk_access_mode)
else if (DiskAccessMode.direct == conf.background_read_disk_access_mode)
{
compactionReadDiskAccessMode = DiskAccessMode.direct;
backgroundReadDiskAccessMode = DiskAccessMode.direct;
}
else
{
throw new IllegalArgumentException("Unsupported disk access mode for compaction_read_disk_access_mode " +
"(options: direct/auto) " + conf.compaction_read_disk_access_mode);
throw new IllegalArgumentException("Unsupported disk access mode for background_read_disk_access_mode " +
"(options: direct/auto) " + conf.background_read_disk_access_mode);
}
logger.info("compaction_read_disk_access_mode resolved to: {}", compactionReadDiskAccessMode);
logger.info("background_read_disk_access_mode resolved to: {}", backgroundReadDiskAccessMode);

/* phi convict threshold for FailureDetector */
if (conf.phi_convict_threshold < 5 || conf.phi_convict_threshold > 16)
Expand Down Expand Up @@ -3459,16 +3459,16 @@ public static void setCommitLogSegmentSize(int sizeMebibytes)
conf.commitlog_segment_size = new DataStorageSpec.IntMebibytesBound(sizeMebibytes);
}

public static DiskAccessMode getCompactionReadDiskAccessMode()
public static DiskAccessMode getBackgroundReadDiskAccessMode()
{
return compactionReadDiskAccessMode;
return backgroundReadDiskAccessMode;
}

@VisibleForTesting
public static void setCompactionReadDiskAccessMode(DiskAccessMode scanDiskAccessMode)
public static void setBackgroundReadDiskAccessMode(DiskAccessMode scanDiskAccessMode)
{
compactionReadDiskAccessMode = scanDiskAccessMode;
conf.compaction_read_disk_access_mode = scanDiskAccessMode;
backgroundReadDiskAccessMode = scanDiskAccessMode;
conf.background_read_disk_access_mode = scanDiskAccessMode;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -263,7 +263,7 @@ public ScannerList getScanners(Collection<SSTableReader> sstables, Collection<Ra
try
{
for (SSTableReader sstable : sstables)
scanners.add(sstable.getScanner(ranges, DatabaseDescriptor.getCompactionReadDiskAccessMode()));
scanners.add(sstable.getScanner(ranges, DatabaseDescriptor.getBackgroundReadDiskAccessMode()));
}
catch (Throwable t)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1777,7 +1777,7 @@ public ISSTableScanner getScanner(SSTableReader sstable)
{
rangesToScan = Collections2.filter(ranges, range -> !transientRanges.contains(range));
}
return sstable.getScanner(rangesToScan, DatabaseDescriptor.getCompactionReadDiskAccessMode());
return sstable.getScanner(rangesToScan, DatabaseDescriptor.getBackgroundReadDiskAccessMode());
}

@Override
Expand All @@ -1800,7 +1800,7 @@ public Full(ColumnFamilyStore cfs, Collection<Range<Token>> ranges, long nowInSe
@Override
public ISSTableScanner getScanner(SSTableReader sstable)
{
return sstable.getScanner(DatabaseDescriptor.getCompactionReadDiskAccessMode());
return sstable.getScanner(DatabaseDescriptor.getBackgroundReadDiskAccessMode());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -513,7 +513,7 @@ private CursorCompactor(OperationType type,
* {@link CompactionIterator#CompactionIterator(OperationType, List, AbstractCompactionController, long, TimeUUID, ActiveCompactionsTracker)}
*/

this.sstableCursors = convertScannersToCursors(scanners, sstables, DatabaseDescriptor.getCompactionReadDiskAccessMode());
this.sstableCursors = convertScannersToCursors(scanners, sstables, DatabaseDescriptor.getBackgroundReadDiskAccessMode());
this.sstableCursorsEqualsNext = new boolean[sstables.size()];
this.enforceStrictLiveness = controller.cfs.metadata.get().enforceStrictLiveness();
this.probeCursorOrder = enforceStrictLiveness ? new StatefulCursor[sstableCursors.length] : null;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -344,7 +344,7 @@ public ScannerList getScanners(Collection<SSTableReader> sstables, Collection<Ra
{
// L0 makes no guarantees about overlapping-ness. Just create a direct scanner for each
for (SSTableReader sstable : byLevel.get(level))
scanners.add(sstable.getScanner(ranges, DatabaseDescriptor.getCompactionReadDiskAccessMode()));
scanners.add(sstable.getScanner(ranges, DatabaseDescriptor.getBackgroundReadDiskAccessMode()));
}
else
{
Expand Down Expand Up @@ -446,7 +446,7 @@ public LeveledScanner(TableMetadata metadata, Collection<SSTableReader> sstables
sstableIterator = this.sstables.iterator();
assert sstableIterator.hasNext(); // caller should check intersecting first
SSTableReader currentSSTable = sstableIterator.next();
currentScanner = currentSSTable.getScanner(ranges, DatabaseDescriptor.getCompactionReadDiskAccessMode());
currentScanner = currentSSTable.getScanner(ranges, DatabaseDescriptor.getBackgroundReadDiskAccessMode());
}

@Override
Expand Down Expand Up @@ -498,7 +498,7 @@ protected UnfilteredRowIterator computeNext()
return endOfData();
}
SSTableReader currentSSTable = sstableIterator.next();
currentScanner = currentSSTable.getScanner(ranges, DatabaseDescriptor.getCompactionReadDiskAccessMode());
currentScanner = currentSSTable.getScanner(ranges, DatabaseDescriptor.getBackgroundReadDiskAccessMode());
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@
import org.apache.cassandra.io.compress.CompressionMetadata;
import org.apache.cassandra.io.sstable.format.SSTableFormat.Components;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.util.ChannelProxy;
import org.apache.cassandra.streaming.ProgressInfo;
import org.apache.cassandra.streaming.StreamSession;
import org.apache.cassandra.streaming.StreamingDataOutputPlus;
Expand Down Expand Up @@ -61,7 +60,7 @@ public void write(StreamingDataOutputPlus out) throws IOException
long totalSize = totalSize();
logger.debug("[Stream #{}] Start streaming file {} to {}, repairedAt = {}, totalSize = {}", session.planId(),
sstable.getFilename(), session.peer, sstable.getSSTableMetadata().repairedAt, totalSize);
try (ChannelProxy fc = sstable.getDataChannel().newChannel())
try (StreamingFileReader fc = StreamingFileReader.open(sstable.descriptor.fileFor(Components.DATA)))
{
long progress = 0L;

Expand All @@ -88,8 +87,8 @@ public void write(StreamingDataOutputPlus out) throws IOException

out.writeToChannel(bufferSupplier -> {
ByteBuffer outBuffer = bufferSupplier.get(toTransfer);
long read = fc.read(outBuffer, position);
assert read == toTransfer : String.format("could not read required number of bytes from file to be streamed: read %d bytes, wanted %d bytes", read, toTransfer);
outBuffer.limit(toTransfer);
fc.readFully(outBuffer, position);
outBuffer.flip();
}, limiter);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,14 @@
package org.apache.cassandra.db.streaming;

import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.apache.cassandra.config.Config.DiskAccessMode;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.io.sstable.Component;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.streaming.ProgressInfo;
Expand Down Expand Up @@ -88,8 +91,49 @@ public void write(StreamingDataOutputPlus out) throws IOException
component,
prettyPrintMemory(length));

FileChannel channel = context.channel(sstable.descriptor, component, length);
long bytesWritten = out.writeFileToChannel(channel, limiter);
long bytesWritten;
// sendfile from an O_DIRECT descriptor either falls back to the page cache or is bounce-buffered
// by the kernel, so direct reads are staged in user space; sendfile is kept when not direct.
StreamingFileReader opened = DatabaseDescriptor.getBackgroundReadDiskAccessMode() == DiskAccessMode.direct
? StreamingFileReader.open(context.file(sstable.descriptor, component))
: null;
if (opened != null && !opened.isDirect())
{
opened.close();
opened = null;
}
final StreamingFileReader reader = opened;
try
{
if (reader != null)
{
if (reader.size() != length)
throw new IOException("Component size changed while streaming " + component);
bytesWritten = 0;
while (bytesWritten < length)
{
long position = bytesWritten;
int count = (int) Math.min(StreamingFileReader.BUFFER_SIZE, length - position);
out.writeToChannel(supplier -> {
ByteBuffer buffer = supplier.get(count);
buffer.limit(count);
reader.readFully(buffer, position);
buffer.flip();
}, limiter);
bytesWritten += count;
}
}
else
{
FileChannel channel = context.channel(sstable.descriptor, component, length);
bytesWritten = out.writeFileToChannel(channel, limiter);
}
}
finally
{
if (reader != null)
reader.close();
}
progress += bytesWritten;

session.progress(sstable.descriptor.fileFor(component).toString(), ProgressInfo.Direction.OUT, bytesWritten, bytesWritten, length);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@
import org.apache.cassandra.io.compress.BufferType;
import org.apache.cassandra.io.sstable.format.SSTableFormat.Components;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.util.ChannelProxy;
import org.apache.cassandra.io.util.DataIntegrityMetadata.ChecksumValidator;
import org.apache.cassandra.streaming.ProgressInfo;
import org.apache.cassandra.streaming.StreamManager;
Expand Down Expand Up @@ -80,7 +79,7 @@ public void write(StreamingDataOutputPlus out) throws IOException
logger.debug("[Stream #{}] Start streaming file {} to {}, repairedAt = {}, totalSize = {}", session.planId(),
sstable.getFilename(), session.peer, sstable.getSSTableMetadata().repairedAt, totalSize);

try(ChannelProxy proxy = sstable.getDataChannel().newChannel();
try(StreamingFileReader proxy = StreamingFileReader.open(sstable.descriptor.fileFor(Components.DATA));
ChecksumValidator validator = sstable.maybeGetChecksumValidator())
{
int bufferSize = validator == null ? DEFAULT_CHUNK_SIZE: validator.chunkSize;
Expand Down Expand Up @@ -140,7 +139,7 @@ protected long totalSize()
*
* @throws java.io.IOException on any I/O error
*/
protected long write(ChannelProxy proxy, ChecksumValidator validator, StreamingDataOutputPlus output, long start, int transferOffset, int toTransfer, int bufferSize) throws IOException
protected long write(StreamingFileReader proxy, ChecksumValidator validator, StreamingDataOutputPlus output, long start, int transferOffset, int toTransfer, int bufferSize) throws IOException
{
// the count of bytes to read off disk
int minReadable = (int) Math.min(bufferSize, proxy.size() - start);
Expand All @@ -150,8 +149,8 @@ protected long write(ChannelProxy proxy, ChecksumValidator validator, StreamingD
ByteBuffer buffer = BufferPools.forNetworking().get(minReadable, BufferType.OFF_HEAP);
try
{
int readCount = proxy.read(buffer, start);
assert readCount == minReadable : String.format("could not read required number of bytes from file to be streamed: read %d bytes, wanted %d bytes", readCount, minReadable);
buffer.limit(minReadable);
proxy.readFully(buffer, start);
buffer.flip();

if (validator != null)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ public ComponentManifest manifest()
*/
public FileChannel channel(Descriptor descriptor, Component component, long size) throws IOException
{
File toTransfer = hardLinks.containsKey(component) ? hardLinks.get(component) : descriptor.fileFor(component);
File toTransfer = file(descriptor, component);
@SuppressWarnings("resource") // file channel will be closed by Caller
FileChannel channel = toTransfer.newReadChannel();

Expand All @@ -83,6 +83,11 @@ public FileChannel channel(Descriptor descriptor, Component component, long size
return channel;
}

File file(Descriptor descriptor, Component component)
{
return hardLinks.getOrDefault(component, descriptor.fileFor(component));
}

@Override
public void close()
{
Expand Down
Loading