package org.apache.cassandra.db.compaction.writers;
import java.util.Arrays;
import java.util.Set;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.Directories;
import org.apache.cassandra.db.RowIndexEntry;
import org.apache.cassandra.db.SerializationHeader;
import org.apache.cassandra.db.rows.UnfilteredRowIterator;
import org.apache.cassandra.db.lifecycle.LifecycleTransaction;
import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
public class SplittingSizeTieredCompactionWriter extends CompactionAwareWriter
{
private static final Logger logger = LoggerFactory.getLogger(SplittingSizeTieredCompactionWriter.class);
public static final long DEFAULT_SMALLEST_SSTABLE_BYTES = 50_000_000;
private final double[] ratios;
private final long totalSize;
private final Set<SSTableReader> allSSTables;
private long currentBytesToWrite;
private int currentRatioIndex = 0;
private Directories.DataDirectory location;
public SplittingSizeTieredCompactionWriter(ColumnFamilyStore cfs, Directories directories, LifecycleTransaction txn, Set<SSTableReader> nonExpiredSSTables)
{
this(cfs, directories, txn, nonExpiredSSTables, DEFAULT_SMALLEST_SSTABLE_BYTES);
}
public SplittingSizeTieredCompactionWriter(ColumnFamilyStore cfs, Directories directories, LifecycleTransaction txn, Set<SSTableReader> nonExpiredSSTables, long smallestSSTable)
{
super(cfs, directories, txn, nonExpiredSSTables, false, false);
this.allSSTables = txn.originals();
totalSize = cfs.getExpectedCompactedFileSize(nonExpiredSSTables, txn.opType());
double[] potentialRatios = new double[20];
double currentRatio = 1;
for (int i = 0; i < potentialRatios.length; i++)
{
currentRatio /= 2;
potentialRatios[i] = currentRatio;
}
int noPointIndex = 0;
for (double ratio : potentialRatios)
{
noPointIndex++;
if (ratio * totalSize < smallestSSTable)
{
break;
}
}
ratios = Arrays.copyOfRange(potentialRatios, 0, noPointIndex);
currentBytesToWrite = Math.round(totalSize * ratios[currentRatioIndex]);
}
@Override
public boolean realAppend(UnfilteredRowIterator partition)
{
RowIndexEntry rie = sstableWriter.append(partition);
if (sstableWriter.currentWriter().getEstimatedOnDiskBytesWritten() > currentBytesToWrite && currentRatioIndex < ratios.length - 1)
{
currentRatioIndex++;
currentBytesToWrite = Math.round(totalSize * ratios[currentRatioIndex]);
switchCompactionLocation(location);
logger.debug("Switching writer, currentBytesToWrite = {}", currentBytesToWrite);
}
return rie != null;
}
@Override
public void switchCompactionLocation(Directories.DataDirectory location)
{
this.location = location;
long currentPartitionsToWrite = Math.round(ratios[currentRatioIndex] * estimatedTotalKeys);
@SuppressWarnings("resource")
SSTableWriter writer = SSTableWriter.create(Descriptor.fromFilename(cfs.getSSTablePath(getDirectories().getLocationForDisk(location))),
currentPartitionsToWrite,
minRepairedAt,
cfs.metadata,
new MetadataCollector(allSSTables, cfs.metadata.comparator, 0),
SerializationHeader.make(cfs.metadata, nonExpiredSSTables),
cfs.indexManager.listIndexes(),
txn);
logger.trace("Switching writer, currentPartitionsToWrite = {}", currentPartitionsToWrite);
sstableWriter.switchWriter(writer);
}
}