From e328bce3cad3e5e35a1f684b46fcc68d102d4e41 Mon Sep 17 00:00:00 2001 From: Mykyta Bozhenko <21245729+cheeeee@users.noreply.github.com> Date: Fri, 11 Sep 2026 18:50:21 -0400 Subject: [PATCH] Honor background read mode across SSTable maintenance and streaming Route maintenance readers through the background read setting and retain the legacy YAML setting replacement. Use bounded aligned reads for direct streaming of physical component bytes, including partial compressed chunks. Keep buffered fallback and the existing zero-copy path for non-direct entire-SSTable streaming. CASSANDRA-19988 Generated-by: Claude (Anthropic) --- CHANGES.txt | 1 + NEWS.txt | 3 +- conf/cassandra.yaml | 10 +- conf/cassandra_latest.yaml | 10 +- .../org/apache/cassandra/config/Config.java | 3 +- .../cassandra/config/DatabaseDescriptor.java | 26 +- .../AbstractCompactionStrategy.java | 2 +- .../db/compaction/CompactionManager.java | 4 +- .../db/compaction/CursorCompactor.java | 2 +- .../compaction/LeveledCompactionStrategy.java | 6 +- .../CassandraCompressedStreamWriter.java | 7 +- .../CassandraEntireSSTableStreamWriter.java | 48 ++- .../db/streaming/CassandraStreamWriter.java | 9 +- .../db/streaming/ComponentContext.java | 7 +- .../db/streaming/StreamingFileReader.java | 177 +++++++++++ .../accord/RouteSecondaryIndexBuilder.java | 3 +- .../sai/StorageAttachedIndexBuilder.java | 3 +- .../index/sasi/SASIIndexBuilder.java | 3 +- .../io/sstable/format/SSTableReader.java | 5 + .../sstable/format/SortedTableScrubber.java | 5 +- .../sstable/format/SortedTableVerifier.java | 5 +- .../apache/cassandra/io/util/FileUtils.java | 10 +- .../cassandra/service/StartupChecks.java | 6 +- .../config/YamlConfigurationLoaderTest.java | 9 + .../db/compaction/CompactionsPurgeTest.java | 6 +- .../db/compaction/CompactionsTest.java | 6 +- ...rgeBoundaryDifferentialCompactionTest.java | 6 +- .../simple/SimpleCompactionTest.java | 6 +- .../db/streaming/BackgroundStreamingTest.java | 289 ++++++++++++++++++ ...assandraEntireSSTableStreamWriterTest.java | 5 +- 30 files changed, 618 insertions(+), 64 deletions(-) create mode 100644 src/java/org/apache/cassandra/db/streaming/StreamingFileReader.java create mode 100644 test/unit/org/apache/cassandra/db/streaming/BackgroundStreamingTest.java diff --git a/CHANGES.txt b/CHANGES.txt index aa356a2da163..b05bcf52a70d 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -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) diff --git a/NEWS.txt b/NEWS.txt index 0032736c903d..889213bc7c48 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -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 === diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 92fafa00de2c..15085d9a302c 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -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: diff --git a/conf/cassandra_latest.yaml b/conf/cassandra_latest.yaml index df4ba596e8a6..4ea204ec1974 100644 --- a/conf/cassandra_latest.yaml +++ b/conf/cassandra_latest.yaml @@ -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: diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index d032b7039c95..feb0e2d13036 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -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(); diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 35caed8fe0d1..bd91317cb758 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -228,7 +228,7 @@ public class DatabaseDescriptor private static DiskAccessMode commitLogWriteDiskAccessMode; - private static DiskAccessMode compactionReadDiskAccessMode; + private static DiskAccessMode backgroundReadDiskAccessMode; private static DiskAccessMode backgroundWriteDiskAccessMode; @@ -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) @@ -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; } /** diff --git a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java index b93b5bf66b2e..02f7bfa5b3dd 100644 --- a/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java +++ b/src/java/org/apache/cassandra/db/compaction/AbstractCompactionStrategy.java @@ -263,7 +263,7 @@ public ScannerList getScanners(Collection sstables, Collection !transientRanges.contains(range)); } - return sstable.getScanner(rangesToScan, DatabaseDescriptor.getCompactionReadDiskAccessMode()); + return sstable.getScanner(rangesToScan, DatabaseDescriptor.getBackgroundReadDiskAccessMode()); } @Override @@ -1800,7 +1800,7 @@ public Full(ColumnFamilyStore cfs, Collection> ranges, long nowInSe @Override public ISSTableScanner getScanner(SSTableReader sstable) { - return sstable.getScanner(DatabaseDescriptor.getCompactionReadDiskAccessMode()); + return sstable.getScanner(DatabaseDescriptor.getBackgroundReadDiskAccessMode()); } @Override diff --git a/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java b/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java index d49931d6b398..b554ae4dd649 100644 --- a/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java +++ b/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java @@ -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; diff --git a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java index 609052c0471a..b7f4fa64e4b6 100644 --- a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java +++ b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java @@ -344,7 +344,7 @@ public ScannerList getScanners(Collection sstables, Collection 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 @@ -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()); } } diff --git a/src/java/org/apache/cassandra/db/streaming/CassandraCompressedStreamWriter.java b/src/java/org/apache/cassandra/db/streaming/CassandraCompressedStreamWriter.java index 806a74a35c30..94095d53aa62 100644 --- a/src/java/org/apache/cassandra/db/streaming/CassandraCompressedStreamWriter.java +++ b/src/java/org/apache/cassandra/db/streaming/CassandraCompressedStreamWriter.java @@ -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; @@ -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; @@ -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); diff --git a/src/java/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriter.java b/src/java/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriter.java index 124e9e757905..401c230226fb 100644 --- a/src/java/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriter.java +++ b/src/java/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriter.java @@ -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; @@ -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); diff --git a/src/java/org/apache/cassandra/db/streaming/CassandraStreamWriter.java b/src/java/org/apache/cassandra/db/streaming/CassandraStreamWriter.java index ea2f1b705304..56a1ba91328b 100644 --- a/src/java/org/apache/cassandra/db/streaming/CassandraStreamWriter.java +++ b/src/java/org/apache/cassandra/db/streaming/CassandraStreamWriter.java @@ -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; @@ -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; @@ -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); @@ -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) diff --git a/src/java/org/apache/cassandra/db/streaming/ComponentContext.java b/src/java/org/apache/cassandra/db/streaming/ComponentContext.java index c03e7b4c3436..dda39780c0b4 100644 --- a/src/java/org/apache/cassandra/db/streaming/ComponentContext.java +++ b/src/java/org/apache/cassandra/db/streaming/ComponentContext.java @@ -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(); @@ -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() { diff --git a/src/java/org/apache/cassandra/db/streaming/StreamingFileReader.java b/src/java/org/apache/cassandra/db/streaming/StreamingFileReader.java new file mode 100644 index 000000000000..5091a96fc890 --- /dev/null +++ b/src/java/org/apache/cassandra/db/streaming/StreamingFileReader.java @@ -0,0 +1,177 @@ +/* + * 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.streaming; + +import java.io.EOFException; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.channels.FileChannel; +import java.nio.file.StandardOpenOption; +import java.util.concurrent.TimeUnit; + +import com.sun.nio.file.ExtendedOpenOption; + +import org.agrona.BitUtil; +import org.agrona.BufferUtil; +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.util.File; +import org.apache.cassandra.io.util.FileUtils; +import org.apache.cassandra.utils.NoSpamLogger; +import org.apache.cassandra.utils.memory.MemoryUtil; + +import sun.nio.ch.DirectBuffer; + +/** Reads stored bytes, without decompression, for both partial and entire SSTable streaming. */ +final class StreamingFileReader implements AutoCloseable +{ + static final int BUFFER_SIZE = 64 * 1024; + // Direct I/O disables read-ahead and each read is synchronous, so stage far more than a network batch. + static final int STAGING_BUFFER_SIZE = 1024 * 1024; + + private static final Logger logger = LoggerFactory.getLogger(StreamingFileReader.class); + + private final FileChannel channel; + private final int blockSize; + private final long size; + private ByteBuffer buffer; + private long bufferOffset = -1; + + static StreamingFileReader open(File file) throws IOException + { + if (DatabaseDescriptor.getBackgroundReadDiskAccessMode() == DiskAccessMode.direct) + { + String reason = "unusable block size"; + try + { + // An oversized or non-power-of-two unit cannot be aligned against, or would dwarf the staging window. + int blockSize = FileUtils.getFileBlockSize(file); + if (blockSize > 0 && blockSize <= STAGING_BUFFER_SIZE && (blockSize & (blockSize - 1)) == 0) + return new StreamingFileReader(FileChannel.open(file.toPath(), StandardOpenOption.READ, ExtendedOpenOption.DIRECT), blockSize); + } + catch (IOException | RuntimeException e) + { + // Probe the actual file, not a temporary file on a potentially different backend. + // Genuine file access failures still propagate from the buffered open below. + reason = e.toString(); + } + NoSpamLogger.log(logger, NoSpamLogger.Level.WARN, 1, TimeUnit.MINUTES, + "Direct I/O is unavailable for streaming {} ({}), falling back to buffered reads", file, reason); + } + return new StreamingFileReader(file.newReadChannel(), 0); + } + + private StreamingFileReader(FileChannel channel, int blockSize) throws IOException + { + this.channel = channel; + this.blockSize = blockSize; + try + { + size = channel.size(); + } + catch (Throwable t) + { + channel.close(); + throw t; + } + } + + boolean isDirect() + { + return blockSize > 0; + } + + long size() + { + return size; + } + + void readFully(ByteBuffer destination, long position) throws IOException + { + if (position < 0 || position > size || destination.remaining() > size - position) + throw new EOFException("Streaming read outside file: " + position + " + " + destination.remaining() + " > " + size); + + if (!isDirect()) + { + while (destination.hasRemaining()) + position += read(channel, destination, position); + return; + } + + if (buffer == null) + { + // Never stage more than the file itself; whole-SSTable streaming opens a reader per component. + int window = size >= STAGING_BUFFER_SIZE ? STAGING_BUFFER_SIZE : BitUtil.align((int) Math.max(size, 1), blockSize); + buffer = BufferUtil.allocateDirectAligned(window, blockSize); + buffer.limit(0); + } + while (destination.hasRemaining()) + { + if (position < bufferOffset || position >= bufferOffset + buffer.limit()) + { + bufferOffset = position - position % buffer.capacity(); + buffer.clear(); + int expected = (int) Math.min(buffer.capacity(), size - bufferOffset); + // Address, offset and length are aligned; only the final physical read may be short. + while (buffer.position() < expected) + { + int count = read(channel, buffer, bufferOffset + buffer.position()); + if (buffer.position() < expected && count % blockSize != 0) + throw new EOFException("Unaligned short direct streaming read at " + bufferOffset); + } + buffer.flip(); + } + int offset = (int) (position - bufferOffset); + int count = Math.min(destination.remaining(), buffer.limit() - offset); + int limit = buffer.limit(); + buffer.position(offset).limit(offset + count); + destination.put(buffer); + buffer.limit(limit); + position += count; + } + } + + private static int read(FileChannel channel, ByteBuffer destination, long position) throws IOException + { + int count = channel.read(destination, position); + if (count <= 0) + throw new EOFException("No streaming read progress at " + position); + return count; + } + + @Override + public void close() throws IOException + { + try + { + channel.close(); + } + finally + { + if (buffer != null) + { + MemoryUtil.clean((ByteBuffer) ((DirectBuffer) buffer).attachment()); + buffer = null; + } + } + } +} diff --git a/src/java/org/apache/cassandra/index/accord/RouteSecondaryIndexBuilder.java b/src/java/org/apache/cassandra/index/accord/RouteSecondaryIndexBuilder.java index c8c7ac4ee8bc..036ffd5ba8a2 100644 --- a/src/java/org/apache/cassandra/index/accord/RouteSecondaryIndexBuilder.java +++ b/src/java/org/apache/cassandra/index/accord/RouteSecondaryIndexBuilder.java @@ -24,6 +24,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.compaction.CompactionInfo; import org.apache.cassandra.db.compaction.CompactionInterruptedException; @@ -124,7 +125,7 @@ private boolean indexSSTable(SSTableReader sstable) return false; } - try (RandomAccessReader dataFile = sstable.openDataReader(); + try (RandomAccessReader dataFile = sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode()); LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.INDEX_BUILD, sstable)) { // remove existing per column index files instead of overwriting diff --git a/src/java/org/apache/cassandra/index/sai/StorageAttachedIndexBuilder.java b/src/java/org/apache/cassandra/index/sai/StorageAttachedIndexBuilder.java index 64fc53525985..21bca3cc809e 100644 --- a/src/java/org/apache/cassandra/index/sai/StorageAttachedIndexBuilder.java +++ b/src/java/org/apache/cassandra/index/sai/StorageAttachedIndexBuilder.java @@ -33,6 +33,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.compaction.CompactionInfo; import org.apache.cassandra.db.compaction.CompactionInterruptedException; @@ -143,7 +144,7 @@ private boolean indexSSTable(SSTableReader sstable, Set in return false; } - try (RandomAccessReader dataFile = sstable.openDataReader(); + try (RandomAccessReader dataFile = sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode()); LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.INDEX_BUILD, sstable)) { perSSTableFileLock = shouldWritePerSSTableFiles(sstable); diff --git a/src/java/org/apache/cassandra/index/sasi/SASIIndexBuilder.java b/src/java/org/apache/cassandra/index/sasi/SASIIndexBuilder.java index 4e71538e41f7..e70206321ca8 100644 --- a/src/java/org/apache/cassandra/index/sasi/SASIIndexBuilder.java +++ b/src/java/org/apache/cassandra/index/sasi/SASIIndexBuilder.java @@ -26,6 +26,7 @@ import java.util.Map; import java.util.SortedMap; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.compaction.CompactionInfo; @@ -80,7 +81,7 @@ public void build() SSTableReader sstable = e.getKey(); Map indexes = e.getValue(); - try (RandomAccessReader dataFile = sstable.openDataReader()) + try (RandomAccessReader dataFile = sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode())) { PerSSTableIndexWriter indexWriter = SASIIndex.newWriter(keyValidator, sstable.descriptor, indexes, OperationType.COMPACTION); targetDirectory = indexWriter.getDescriptor().directory.path(); diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java index ab3171200f17..7dd4ca8b22e6 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java @@ -1435,6 +1435,11 @@ public RandomAccessReader openDataReader(DiskAccessMode diskAccessMode) return openDataReaderInternal(diskAccessMode, null, false); } + public RandomAccessReader openDataReader(DiskAccessMode diskAccessMode, RateLimiter limiter) + { + return openDataReaderInternal(diskAccessMode, limiter, false); + } + public RandomAccessReader openDataReaderForScan() { return openDataReaderInternal(null, null, true); diff --git a/src/java/org/apache/cassandra/io/sstable/format/SortedTableScrubber.java b/src/java/org/apache/cassandra/io/sstable/format/SortedTableScrubber.java index 66effe8e455f..775caf27f492 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SortedTableScrubber.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SortedTableScrubber.java @@ -40,6 +40,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ClusteringComparator; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.DecoratedKey; @@ -155,8 +156,8 @@ protected SortedTableScrubber(ColumnFamilyStore cfs, // partition header (key or data size) is corrupt. (This means our position in the index file will be one // partition "ahead" of the data file.) this.dataFile = transaction.isOffline() - ? sstable.openDataReader() - : sstable.openDataReader(CompactionManager.instance.getRateLimiter()); + ? sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode()) + : sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode(), CompactionManager.instance.getRateLimiter()); this.scrubInfo = new ScrubInfo(dataFile, sstable, fileAccessLock.readLock()); diff --git a/src/java/org/apache/cassandra/io/sstable/format/SortedTableVerifier.java b/src/java/org/apache/cassandra/io/sstable/format/SortedTableVerifier.java index b586c1b9f6ab..3b0d4b62d3a1 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SortedTableVerifier.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SortedTableVerifier.java @@ -37,6 +37,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.compaction.CompactionController; @@ -97,8 +98,8 @@ public SortedTableVerifier(ColumnFamilyStore cfs, R sstable, OutputHandler outpu this.fileAccessLock = new ReentrantReadWriteLock(); this.dataFile = isOffline - ? sstable.openDataReader() - : sstable.openDataReader(CompactionManager.instance.getRateLimiter()); + ? sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode()) + : sstable.openDataReader(DatabaseDescriptor.getBackgroundReadDiskAccessMode(), CompactionManager.instance.getRateLimiter()); this.verifyInfo = new VerifyInfo(dataFile, sstable, fileAccessLock.readLock()); this.options = options; this.isOffline = isOffline; diff --git a/src/java/org/apache/cassandra/io/util/FileUtils.java b/src/java/org/apache/cassandra/io/util/FileUtils.java index 9eb4bd9d10f1..8edcb4e6922e 100644 --- a/src/java/org/apache/cassandra/io/util/FileUtils.java +++ b/src/java/org/apache/cassandra/io/util/FileUtils.java @@ -772,11 +772,19 @@ public static boolean isDirectIOSupported(File file) } } + /** Block size of the filesystem backing an already-existing file; shares {@link #getBlockSize(File)}'s cache. */ public static int getFileBlockSize(File file) { + File directory = file.parent(); + String key = directory != null ? directory.absolutePath() : file.absolutePath(); + Integer cached = blockSizeByDirectory.get(key); + if (cached != null) + return cached; try { - return blockSize(file); + int size = blockSize(file); + blockSizeByDirectory.put(key, size); + return size; } catch (IOException e) { diff --git a/src/java/org/apache/cassandra/service/StartupChecks.java b/src/java/org/apache/cassandra/service/StartupChecks.java index 6148e1f8218a..2441bc197273 100644 --- a/src/java/org/apache/cassandra/service/StartupChecks.java +++ b/src/java/org/apache/cassandra/service/StartupChecks.java @@ -883,7 +883,7 @@ public void execute(StartupChecksConfiguration configuration) throws StartupExce if (configuration.isDisabled(name())) return; - boolean directReads = DatabaseDescriptor.getCompactionReadDiskAccessMode() == Config.DiskAccessMode.direct; + boolean directReads = DatabaseDescriptor.getBackgroundReadDiskAccessMode() == Config.DiskAccessMode.direct; boolean directWrites = DatabaseDescriptor.getBackgroundWriteDiskAccessMode() == Config.DiskAccessMode.direct; if (!directReads && !directWrites) @@ -894,8 +894,8 @@ public void execute(StartupChecksConfiguration configuration) throws StartupExce if (!unsupportedLocations.isEmpty()) { String configuredModes = directReads && directWrites - ? "compaction reads and background writes" - : directReads ? "compaction reads" : "background writes"; + ? "background reads and writes" + : directReads ? "background reads" : "background writes"; throw new StartupException(StartupException.ERR_WRONG_DISK_STATE, String.format("Direct I/O is configured for %s, " + diff --git a/test/unit/org/apache/cassandra/config/YamlConfigurationLoaderTest.java b/test/unit/org/apache/cassandra/config/YamlConfigurationLoaderTest.java index c0ce3d3ea5f5..983d3520eaff 100644 --- a/test/unit/org/apache/cassandra/config/YamlConfigurationLoaderTest.java +++ b/test/unit/org/apache/cassandra/config/YamlConfigurationLoaderTest.java @@ -61,6 +61,15 @@ public class YamlConfigurationLoaderTest { + @Test + public void backgroundReadModeAcceptsLegacyName() + { + Config legacy = YamlConfigurationLoader.fromMap(ImmutableMap.of("compaction_read_disk_access_mode", "direct"), true, Config.class); + Config current = YamlConfigurationLoader.fromMap(ImmutableMap.of("background_read_disk_access_mode", "direct"), true, Config.class); + assertEquals(Config.DiskAccessMode.direct, legacy.background_read_disk_access_mode); + assertEquals(legacy.background_read_disk_access_mode, current.background_read_disk_access_mode); + } + @Test public void repairRetryEmpty() { diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java index 3fa099b6fe9a..0605aef1dd6a 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionsPurgeTest.java @@ -91,16 +91,16 @@ public static Collection params() @Before public void setCompactionParams() { - originalDiskAccessMode = DatabaseDescriptor.getCompactionReadDiskAccessMode(); + originalDiskAccessMode = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); originalCursorCompactionEnabled = DatabaseDescriptor.cursorCompactionEnabled(); - DatabaseDescriptor.setCompactionReadDiskAccessMode(compactionReadDiskAccessMode); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(compactionReadDiskAccessMode); DatabaseDescriptor.setCursorCompactionEnabled(cursorCompactionEnabled); } @After public void restoreCompactionParams() { - DatabaseDescriptor.setCompactionReadDiskAccessMode(originalDiskAccessMode); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(originalDiskAccessMode); DatabaseDescriptor.setCursorCompactionEnabled(originalCursorCompactionEnabled); } diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionsTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionsTest.java index 36ab11a02bc6..541930bdafbd 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionsTest.java @@ -137,10 +137,10 @@ public static Collection params() @Before public void setCompactionParams() { - originalDiskAccessMode = DatabaseDescriptor.getCompactionReadDiskAccessMode(); + originalDiskAccessMode = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); originalCursorCompactionEnabled = DatabaseDescriptor.cursorCompactionEnabled(); originalBackgroundWriteDiskAccessMode = DatabaseDescriptor.getBackgroundWriteDiskAccessMode(); - DatabaseDescriptor.setCompactionReadDiskAccessMode(compactionReadDiskAccessMode); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(compactionReadDiskAccessMode); DatabaseDescriptor.setCursorCompactionEnabled(cursorCompactionEnabled); DatabaseDescriptor.setBackgroundWriteDiskAccessMode(backgroundWriteDiskAccessMode); } @@ -148,7 +148,7 @@ public void setCompactionParams() @After public void restoreCompactionParams() { - DatabaseDescriptor.setCompactionReadDiskAccessMode(originalDiskAccessMode); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(originalDiskAccessMode); DatabaseDescriptor.setCursorCompactionEnabled(originalCursorCompactionEnabled); DatabaseDescriptor.setBackgroundWriteDiskAccessMode(originalBackgroundWriteDiskAccessMode); } diff --git a/test/unit/org/apache/cassandra/db/compaction/differential/PurgeBoundaryDifferentialCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/differential/PurgeBoundaryDifferentialCompactionTest.java index d90cf8b3895b..e3fd4b9cdd3b 100644 --- a/test/unit/org/apache/cassandra/db/compaction/differential/PurgeBoundaryDifferentialCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/differential/PurgeBoundaryDifferentialCompactionTest.java @@ -142,15 +142,15 @@ public void directDiskAccessMode() throws Exception flush(); } - DiskAccessMode original = DatabaseDescriptor.getCompactionReadDiskAccessMode(); - DatabaseDescriptor.setCompactionReadDiskAccessMode(DiskAccessMode.direct); + DiskAccessMode original = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(DiskAccessMode.direct); try { assertCursorMatchesIterator(cfs); } finally { - DatabaseDescriptor.setCompactionReadDiskAccessMode(original); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(original); } } diff --git a/test/unit/org/apache/cassandra/db/compaction/simple/SimpleCompactionTest.java b/test/unit/org/apache/cassandra/db/compaction/simple/SimpleCompactionTest.java index 602b380c3e81..aeb18a0acfd5 100644 --- a/test/unit/org/apache/cassandra/db/compaction/simple/SimpleCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/simple/SimpleCompactionTest.java @@ -75,16 +75,16 @@ public static Collection params() @Before public void setCompactionParams() { - originalDiskAccessMode = DatabaseDescriptor.getCompactionReadDiskAccessMode(); + originalDiskAccessMode = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); originalCursorCompactionEnabled = DatabaseDescriptor.cursorCompactionEnabled(); - DatabaseDescriptor.setCompactionReadDiskAccessMode(compactionReadDiskAccessMode); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(compactionReadDiskAccessMode); DatabaseDescriptor.setCursorCompactionEnabled(cursorCompactionEnabled); } @After public void restoreCompactionParams() { - DatabaseDescriptor.setCompactionReadDiskAccessMode(originalDiskAccessMode); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(originalDiskAccessMode); DatabaseDescriptor.setCursorCompactionEnabled(originalCursorCompactionEnabled); } diff --git a/test/unit/org/apache/cassandra/db/streaming/BackgroundStreamingTest.java b/test/unit/org/apache/cassandra/db/streaming/BackgroundStreamingTest.java new file mode 100644 index 000000000000..d4be4779a16c --- /dev/null +++ b/test/unit/org/apache/cassandra/db/streaming/BackgroundStreamingTest.java @@ -0,0 +1,289 @@ +/* + * 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.streaming; + +import java.io.EOFException; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.channels.FileChannel; +import java.nio.file.Files; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Random; + +import net.jpountz.lz4.LZ4Factory; + +import org.junit.BeforeClass; +import org.junit.Test; + +import org.apache.cassandra.SchemaLoader; +import org.apache.cassandra.Util; +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.RowUpdateBuilder; +import org.apache.cassandra.io.compress.CompressionMetadata; +import org.apache.cassandra.io.sstable.Component; +import org.apache.cassandra.io.sstable.format.SSTableFormat.Components; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.io.util.DataInputBuffer; +import org.apache.cassandra.io.util.DataOutputBuffer; +import org.apache.cassandra.io.util.File; +import org.apache.cassandra.io.util.FileUtils; +import org.apache.cassandra.schema.CompressionParams; +import org.apache.cassandra.schema.KeyspaceParams; +import org.apache.cassandra.streaming.PreviewKind; +import org.apache.cassandra.streaming.SessionInfo; +import org.apache.cassandra.streaming.StreamCoordinator; +import org.apache.cassandra.streaming.StreamOperation; +import org.apache.cassandra.streaming.StreamResultFuture; +import org.apache.cassandra.streaming.StreamSession; +import org.apache.cassandra.streaming.StreamingDataOutputPlus; +import org.apache.cassandra.streaming.async.NettyStreamingConnectionFactory; +import org.apache.cassandra.streaming.async.StreamCompressionSerializer; +import org.apache.cassandra.utils.FBUtilities; + +import io.netty.buffer.ByteBuf; +import io.netty.buffer.UnpooledByteBufAllocator; + +import static java.util.Collections.emptyList; +import static org.apache.cassandra.net.MessagingService.current_version; +import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; + +public class BackgroundStreamingTest +{ + private static final String KEYSPACE = "BackgroundStreamingTest"; + private static SSTableReader compressed; + private static SSTableReader uncompressed; + + @BeforeClass + public static void setup() + { + SchemaLoader.prepareServer(); + SchemaLoader.createKeyspace(KEYSPACE, KeyspaceParams.simple(1), + SchemaLoader.standardCFMD(KEYSPACE, "compressed").compression(CompressionParams.lz4()), + SchemaLoader.standardCFMD(KEYSPACE, "uncompressed").compression(CompressionParams.noCompression())); + compressed = writeSSTable("compressed"); + uncompressed = writeSSTable("uncompressed"); + } + + private static SSTableReader writeSSTable(String table) + { + ColumnFamilyStore cfs = Keyspace.open(KEYSPACE).getColumnFamilyStore(table); + cfs.disableAutoCompaction(); + byte[] value = new byte[3 * StreamingFileReader.BUFFER_SIZE + 37]; + new Random(1).nextBytes(value); + new RowUpdateBuilder(cfs.metadata(), 1, "key").clustering("row").add("val", ByteBuffer.wrap(value)).build().applyUnsafe(); + Util.flush(cfs); + return cfs.getLiveSSTables().iterator().next(); + } + + @Test + public void rawReaderHonorsModeAcrossUnalignedBoundaries() throws Exception + { + DiskAccessMode previous = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); + File file = FileUtils.createTempFile("background-stream", ".db", uncompressed.descriptor.fileFor(Components.DATA).parent()); + byte[] content = new byte[StreamingFileReader.STAGING_BUFFER_SIZE + StreamingFileReader.BUFFER_SIZE + 37]; + new Random(2).nextBytes(content); + Files.write(file.toPath(), content); + try + { + for (DiskAccessMode mode : new DiskAccessMode[]{ DiskAccessMode.standard, DiskAccessMode.direct }) + { + DatabaseDescriptor.setBackgroundReadDiskAccessMode(mode); + try (StreamingFileReader reader = StreamingFileReader.open(file)) + { + assertEquals(mode == DiskAccessMode.direct && FileUtils.isDirectIOSupported(file), reader.isDirect()); + // Includes disjoint/backward reads, a staging-window crossing, and the unaligned EOF tail. + int[][] ranges = { { 17, 3 }, { 65533, 65543 }, { content.length - 39, 39 }, { StreamingFileReader.STAGING_BUFFER_SIZE - 5, 65541 }, { 1, 131077 } }; + for (int[] range : ranges) + { + ByteBuffer actual = ByteBuffer.allocate(range[1]); + reader.readFully(actual, range[0]); + assertArrayEquals(Arrays.copyOfRange(content, range[0], range[0] + range[1]), actual.array()); + } + reader.readFully(ByteBuffer.allocate(0), content.length); + assertThatThrownBy(() -> reader.readFully(ByteBuffer.allocate(2), content.length - 1)) + .isInstanceOf(EOFException.class); + } + } + } + finally + { + DatabaseDescriptor.setBackgroundReadDiskAccessMode(previous); + file.delete(); + } + } + + @Test + public void compressedSectionsPreserveStoredChunksAndChecksums() throws Exception + { + CompressionMetadata metadata = compressed.getCompressionMetadata(); + List sections = Arrays.asList(new SSTableReader.PartitionPositionBounds(1, 2), + new SSTableReader.PartitionPositionBounds(metadata.dataLength - 13, metadata.dataLength)); + CompressionInfo info = CompressionInfo.newLazyInstance(metadata, sections); + byte[] stored = Files.readAllBytes(compressed.descriptor.fileFor(Components.DATA).toPath()); + try (DataOutputBuffer expected = new DataOutputBuffer()) + { + for (CompressionMetadata.Chunk chunk : info.chunks()) + expected.write(stored, (int) chunk.offset, chunk.length + Integer.BYTES); + assertPartialStream(compressed, sections, info, expected.toByteArray()); + } + } + + @Test + public void uncompressedSectionsPreserveChecksumBoundarySlices() throws Exception + { + byte[] stored = Files.readAllBytes(uncompressed.descriptor.fileFor(Components.DATA).toPath()); + List sections = Arrays.asList(new SSTableReader.PartitionPositionBounds(17, 131079), + new SSTableReader.PartitionPositionBounds(stored.length - 29, stored.length)); + try (DataOutputBuffer expected = new DataOutputBuffer()) + { + for (SSTableReader.PartitionPositionBounds section : sections) + expected.write(stored, (int) section.lowerPosition, (int) (section.upperPosition - section.lowerPosition)); + assertPartialStream(uncompressed, sections, null, expected.toByteArray()); + } + } + + private void assertPartialStream(SSTableReader sstable, List sections, + CompressionInfo info, byte[] expected) throws Exception + { + DiskAccessMode previous = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); + CassandraStreamHeader header = CassandraStreamHeader.builder().withSSTableVersion(sstable.descriptor.version) + .withSections(sections).withCompressionInfo(info) + .withSerializationHeader(sstable.header.toComponent()) + .withTableId(sstable.metadata().id).build(); + try + { + for (DiskAccessMode mode : new DiskAccessMode[]{ DiskAccessMode.standard, DiskAccessMode.direct }) + { + DatabaseDescriptor.setBackgroundReadDiskAccessMode(mode); + try (CollectingOutput output = new CollectingOutput()) + { + CassandraStreamWriter writer = info == null ? new CassandraStreamWriter(sstable, header, session()) + : new CassandraCompressedStreamWriter(sstable, header, session()); + writer.write(output); + if (info != null) + assertArrayEquals(expected, output.toByteArray()); + else + { + StreamCompressionSerializer serializer = new StreamCompressionSerializer(UnpooledByteBufAllocator.DEFAULT); + try (DataInputBuffer input = new DataInputBuffer(output.toByteArray()); + DataOutputBuffer decoded = new DataOutputBuffer()) + { + while (input.available() > 0) + { + ByteBuf chunk = serializer.deserialize(LZ4Factory.fastestInstance().safeDecompressor(), input, current_version); + try + { + decoded.write(chunk.nioBuffer()); + } + finally + { + chunk.release(); + } + } + assertArrayEquals(expected, decoded.toByteArray()); + } + } + } + } + } + finally + { + DatabaseDescriptor.setBackgroundReadDiskAccessMode(previous); + } + } + + @Test + public void entireDirectStreamPreservesEveryComponent() throws Exception + { + DiskAccessMode previous = DatabaseDescriptor.getBackgroundReadDiskAccessMode(); + DatabaseDescriptor.setBackgroundReadDiskAccessMode(DiskAccessMode.direct); + try + { + for (SSTableReader sstable : new SSTableReader[]{ compressed, uncompressed }) + { + try (ComponentContext context = ComponentContext.create(sstable); + CollectingOutput output = new CollectingOutput(); + DataOutputBuffer expected = new DataOutputBuffer()) + { + for (Component component : context.manifest().components()) + expected.write(Files.readAllBytes(context.file(sstable.descriptor, component).toPath())); + new CassandraEntireSSTableStreamWriter(sstable, session(), context).write(output); + assertArrayEquals(expected.toByteArray(), output.toByteArray()); + try (StreamingFileReader reader = StreamingFileReader.open(context.file(sstable.descriptor, Components.DATA))) + { + // Only a reader that really obtained direct I/O may cost the zero-copy path. + assertEquals(!reader.isDirect(), output.sendfile); + } + } + } + } + finally + { + DatabaseDescriptor.setBackgroundReadDiskAccessMode(previous); + } + } + + private static StreamSession session() + { + StreamCoordinator coordinator = new StreamCoordinator(StreamOperation.BOOTSTRAP, 1, new NettyStreamingConnectionFactory(), false, false, null, PreviewKind.NONE); + StreamResultFuture future = StreamResultFuture.createInitiator(nextTimeUUID(), StreamOperation.BOOTSTRAP, Collections.emptyList(), coordinator); + coordinator.addSessionInfo(new SessionInfo(FBUtilities.getBroadcastAddressAndPort(), 0, FBUtilities.getBroadcastAddressAndPort(), emptyList(), emptyList(), StreamSession.State.INITIALIZED, null)); + StreamSession session = coordinator.getOrCreateOutboundSession(FBUtilities.getBroadcastAddressAndPort()); + session.init(future); + return session; + } + + private static class CollectingOutput extends DataOutputBuffer implements StreamingDataOutputPlus + { + private boolean sendfile; + + @Override + public int writeToChannel(Write write, RateLimiter limiter) throws IOException + { + ByteBuffer[] supplied = new ByteBuffer[1]; + write.write(size -> supplied[0] = ByteBuffer.allocate(size + 8)); + int count = supplied[0].remaining(); + write(supplied[0]); + return count; + } + + @Override + public long writeFileToChannel(FileChannel file, RateLimiter limiter) throws IOException + { + sendfile = true; + ByteBuffer buffer = ByteBuffer.allocate((int) file.size()); + try (FileChannel closing = file) + { + while (buffer.hasRemaining()) + closing.read(buffer); + } + buffer.flip(); + write(buffer); + return buffer.limit(); + } + } +} diff --git a/test/unit/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriterTest.java b/test/unit/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriterTest.java index 3cb982b89bb4..1bacd4d74ffc 100644 --- a/test/unit/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriterTest.java +++ b/test/unit/org/apache/cassandra/db/streaming/CassandraEntireSSTableStreamWriterTest.java @@ -200,7 +200,10 @@ public void close() throws IOException @Override public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception { - ((SharedDefaultFileRegion) msg).transferTo(wbc, 0); + if (msg instanceof ByteBuf) + serializedFile.writeBytes(((ByteBuf) msg).duplicate()); + else + ((SharedDefaultFileRegion) msg).transferTo(wbc, 0); super.write(ctx, msg, promise); } });