diff --git a/CHANGES.txt b/CHANGES.txt index 78cacab692..32b6f2d435 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -6,6 +6,7 @@ Merged from 6.0: * Differentiate between legitimate cases where the first entry is the same as the last entry and empty bounds in SSTableCursorWriter#addIndexBlock() (CASSANDRA-21255) * Introduce minimum_threshold for data resurrection startup check (CASSANDRA-21293) Merged from 5.0: + * Use estimated compressed size for tables to check if there is enough free space for a compaction (CASSANDRA-21245) * Fix failing select on system_views.settings for non-string keys (CASSANDRA-21348) diff --git a/src/java/org/apache/cassandra/cache/AutoSavingCache.java b/src/java/org/apache/cassandra/cache/AutoSavingCache.java index 3c3d9a8cdd..fc15aacdb9 100644 --- a/src/java/org/apache/cassandra/cache/AutoSavingCache.java +++ b/src/java/org/apache/cassandra/cache/AutoSavingCache.java @@ -330,6 +330,7 @@ public class AutoSavingCache extends InstrumentingCache extends InstrumentingCache estimatedRemainingWriteBytes() + public Map estimatedRemainingWriteToDiskBytes() { synchronized (compactions) { @@ -66,7 +67,7 @@ public class ActiveCompactions implements ActiveCompactionsTracker List directories = compactionInfo.getTargetDirectories(); if (directories == null || directories.isEmpty()) continue; - long remainingWriteBytesPerDataDir = compactionInfo.estimatedRemainingWriteBytes() / directories.size(); + long remainingWriteBytesPerDataDir = compactionInfo.estimatedRemainingWriteToDiskBytes() / directories.size(); for (File directory : directories) writeBytesPerSSTableDir.merge(directory, remainingWriteBytesPerDataDir, Long::sum); } diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java b/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java index 0bfc925a7d..73dc06911a 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionInfo.java @@ -42,6 +42,7 @@ public final class CompactionInfo public static final String COLUMNFAMILY = "columnfamily"; public static final String COMPLETED = "completed"; public static final String TOTAL = "total"; + public static final String TOTAL_COMPRESSED = "totalCompressed"; public static final String TASK_TYPE = "taskType"; public static final String UNIT = "unit"; public static final String COMPACTION_ID = "compactionId"; @@ -52,16 +53,18 @@ public final class CompactionInfo private final OperationType tasktype; private final long completed; private final long total; + private final long totalCompressed; private final Unit unit; private final TimeUUID compactionId; private final ImmutableSet sstables; private final String targetDirectory; - public CompactionInfo(TableMetadata metadata, OperationType tasktype, long completed, long total, Unit unit, TimeUUID compactionId, Collection sstables, String targetDirectory) + public CompactionInfo(TableMetadata metadata, OperationType tasktype, long completed, long total, long totalCompressed, Unit unit, TimeUUID compactionId, Collection sstables, String targetDirectory) { this.tasktype = tasktype; this.completed = completed; this.total = total; + this.totalCompressed = totalCompressed; this.metadata = metadata; this.unit = unit; this.compactionId = compactionId; @@ -69,38 +72,38 @@ public final class CompactionInfo this.targetDirectory = targetDirectory; } - public CompactionInfo(TableMetadata metadata, OperationType tasktype, long completed, long total, TimeUUID compactionId, Collection sstables, String targetDirectory) + public CompactionInfo(TableMetadata metadata, OperationType tasktype, long completed, long total, long totalCompressed, TimeUUID compactionId, Collection sstables, String targetDirectory) { - this(metadata, tasktype, completed, total, Unit.BYTES, compactionId, sstables, targetDirectory); + this(metadata, tasktype, completed, total, totalCompressed, Unit.BYTES, compactionId, sstables, targetDirectory); } - public CompactionInfo(TableMetadata metadata, OperationType tasktype, long completed, long total, TimeUUID compactionId, Collection sstables) + public CompactionInfo(TableMetadata metadata, OperationType tasktype, long completed, long total, long totalCompressed, TimeUUID compactionId, Collection sstables) { - this(metadata, tasktype, completed, total, Unit.BYTES, compactionId, sstables, null); + this(metadata, tasktype, completed, total, totalCompressed, Unit.BYTES, compactionId, sstables, null); } /** * Special compaction info where we always need to cancel the compaction - for example ViewBuilderTask where we don't know * the sstables at construction */ - public static CompactionInfo withoutSSTables(TableMetadata metadata, OperationType tasktype, long completed, long total, Unit unit, TimeUUID compactionId) + public static CompactionInfo withoutSSTables(TableMetadata metadata, OperationType tasktype, long completed, long total, long totalCompressed, Unit unit, TimeUUID compactionId) { - return withoutSSTables(metadata, tasktype, completed, total, unit, compactionId, null); + return withoutSSTables(metadata, tasktype, completed, total, totalCompressed, unit, compactionId, null); } /** * Special compaction info where we always need to cancel the compaction - for example AutoSavingCache where we don't know * the sstables at construction */ - public static CompactionInfo withoutSSTables(TableMetadata metadata, OperationType tasktype, long completed, long total, Unit unit, TimeUUID compactionId, String targetDirectory) + public static CompactionInfo withoutSSTables(TableMetadata metadata, OperationType tasktype, long completed, long total, long totalCompressed, Unit unit, TimeUUID compactionId, String targetDirectory) { - return new CompactionInfo(metadata, tasktype, completed, total, unit, compactionId, ImmutableSet.of(), targetDirectory); + return new CompactionInfo(metadata, tasktype, completed, total, totalCompressed, unit, compactionId, ImmutableSet.of(), targetDirectory); } /** @return A copy of this CompactionInfo with updated progress. */ - public CompactionInfo forProgress(long complete, long total) + public CompactionInfo forProgress(long complete, long total, long totalCompressed) { - return new CompactionInfo(metadata, tasktype, complete, total, unit, compactionId, sstables, targetDirectory); + return new CompactionInfo(metadata, tasktype, complete, total, totalCompressed, unit, compactionId, sstables, targetDirectory); } public Optional getKeyspace() @@ -128,6 +131,11 @@ public final class CompactionInfo return total; } + public long getTotalCompressed() + { + return totalCompressed; + } + public OperationType getTaskType() { return tasktype; @@ -183,12 +191,16 @@ public final class CompactionInfo /** * Note that this estimate is based on the amount of data we have left to read - it assumes input * size == output size for a compaction, which is not really true, but should most often provide a worst case - * remaining write size. + * remaining write size. We also scale by the effective compression ratio since total/completed are for the uncompressed size. */ - public long estimatedRemainingWriteBytes() + public long estimatedRemainingWriteToDiskBytes() { if (unit == Unit.BYTES && tasktype.writesData) - return getTotal() - getCompleted(); + { + final long total = getTotal(); + double compressionRatio = total == 0 ? 1 : ((double) totalCompressed / (double)total); + return (long)(compressionRatio * (total - getCompleted())); + } return 0; } @@ -216,6 +228,7 @@ public final class CompactionInfo ret.put(COLUMNFAMILY, getTable().orElse(null)); ret.put(COMPLETED, Long.toString(completed)); ret.put(TOTAL, Long.toString(total)); + ret.put(TOTAL_COMPRESSED, Long.toString(totalCompressed)); ret.put(TASK_TYPE, tasktype.toString()); ret.put(UNIT, unit.toString()); ret.put(COMPACTION_ID, compactionId == null ? "" : compactionId.toString()); diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java index 9c1ec42def..35739ec1bc 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionIterator.java @@ -158,6 +158,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte private final long nowInSec; private final TimeUUID compactionId; private final long totalBytes; + private final long totalCompressedBytes; private long bytesRead; private long totalSourceCQLRows; @@ -231,9 +232,14 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte this.bytesRead = 0; long bytes = 0; + long compressedBytes = 0; for (ISSTableScanner scanner : scanners) + { bytes += scanner.getLengthInBytes(); + compressedBytes += scanner.getCompressedLengthInBytes(); + } this.totalBytes = bytes; + this.totalCompressedBytes = compressedBytes; this.mergeCounters = new long[scanners.size()]; // note that we leak `this` from the constructor when calling beginCompaction below, this means we have to get the sstables before // calling that to avoid a NPE. @@ -281,6 +287,7 @@ public class CompactionIterator extends CompactionInfo.Holder implements Unfilte type, bytesRead, totalBytes, + totalCompressedBytes, compactionId, sstables, targetDirectory); diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java index 7336c4543a..bf62db23f5 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java @@ -489,7 +489,7 @@ public class CompactionTask extends AbstractCompactionTask for (File directory : newCompactionDatadirs) expectedNewWriteSize.put(directory, writeSizePerOutputDatadir); - Map expectedWriteSize = CompactionManager.instance.active.estimatedRemainingWriteBytes(); + Map expectedWriteSize = CompactionManager.instance.active.estimatedRemainingWriteToDiskBytes(); // todo: abort streams if they block compactions if (cfs.getDirectories().hasDiskSpaceForCompactionsAndStreams(expectedNewWriteSize, expectedWriteSize)) diff --git a/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java b/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java index 438c5a06dc..7ffe51c896 100644 --- a/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java +++ b/src/java/org/apache/cassandra/db/compaction/CursorCompactor.java @@ -210,6 +210,7 @@ public class CursorCompactor extends CompactionInfo.Holder private final long nowInSec; private final TimeUUID compactionId; private final long totalInputBytes; + private final long totalCompressedInputBytes; private final StatefulCursor[] sstableCursors; private final boolean[] sstableCursorsEqualsNext; private final boolean hasStaticColumns; @@ -268,9 +269,14 @@ public class CursorCompactor extends CompactionInfo.Holder this.compactionId = compactionId; long inputBytes = 0; + long compressedInputBytes = 0; for (ISSTableScanner scanner : scanners) + { inputBytes += scanner.getLengthInBytes(); + compressedInputBytes += scanner.getCompressedLengthInBytes(); + } this.totalInputBytes = inputBytes; + this.totalCompressedInputBytes = compressedInputBytes; this.partitionMergeCounters = new long[scanners.size()]; this.staticRowMergeCounters = new long[partitionMergeCounters.length]; this.rowMergeCounters = new long[partitionMergeCounters.length]; @@ -1461,6 +1467,7 @@ public class CursorCompactor extends CompactionInfo.Holder type, getBytesRead(), totalInputBytes, + totalCompressedInputBytes, compactionId, sstables, targetDirectory); @@ -1666,4 +1673,4 @@ public class CursorCompactor extends CompactionInfo.Holder } preSortedArray[insertInto] = newElement; } -} \ No newline at end of file +} diff --git a/src/java/org/apache/cassandra/db/view/ViewBuilderTask.java b/src/java/org/apache/cassandra/db/view/ViewBuilderTask.java index daf08794d5..548751667f 100644 --- a/src/java/org/apache/cassandra/db/view/ViewBuilderTask.java +++ b/src/java/org/apache/cassandra/db/view/ViewBuilderTask.java @@ -205,13 +205,13 @@ public class ViewBuilderTask extends CompactionInfo.Holder implements Callable imp OperationType.SCRUB, dataFile.getFilePointer(), dataFile.length(), + sstable.onDiskLength(), scrubCompactionId, ImmutableSet.of(sstable), File.getPath(sstable.getFilename()).getParent().toString()); 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 305c337a47..b586c1b9f6 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SortedTableVerifier.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SortedTableVerifier.java @@ -499,6 +499,7 @@ public abstract class SortedTableVerifier imp OperationType.VERIFY, dataFile.getFilePointer(), dataFile.length(), + sstable.onDiskLength(), verificationCompactionId, ImmutableSet.of(sstable)); } diff --git a/src/java/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryRedistribution.java b/src/java/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryRedistribution.java index 335a659354..e37014c550 100644 --- a/src/java/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryRedistribution.java +++ b/src/java/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryRedistribution.java @@ -363,7 +363,7 @@ public class IndexSummaryRedistribution extends CompactionInfo.Holder public CompactionInfo getCompactionInfo() { - return CompactionInfo.withoutSSTables(null, OperationType.INDEX_SUMMARY, (memoryPoolBytes - remainingSpace), memoryPoolBytes, Unit.BYTES, compactionId); + return CompactionInfo.withoutSSTables(null, OperationType.INDEX_SUMMARY, (memoryPoolBytes - remainingSpace), memoryPoolBytes, memoryPoolBytes, Unit.BYTES, compactionId); } public boolean isGlobal() diff --git a/src/java/org/apache/cassandra/metrics/CompactionMetrics.java b/src/java/org/apache/cassandra/metrics/CompactionMetrics.java index 03d51e2ecf..f908b82491 100644 --- a/src/java/org/apache/cassandra/metrics/CompactionMetrics.java +++ b/src/java/org/apache/cassandra/metrics/CompactionMetrics.java @@ -54,6 +54,8 @@ public class CompactionMetrics public final Meter totalCompactionsCompleted; /** Total number of bytes compacted since server [re]start */ public final Counter bytesCompacted; + /** Estimated compressed bytes compacted since server [re]start, computed by scaling uncompressed bytes by the compression ratio */ + public final Counter compressedBytesCompacted; /** Recent/current throughput of compactions take */ public final Meter bytesCompactedThroughput; /** Time spent redistributing index summaries */ @@ -150,6 +152,7 @@ public class CompactionMetrics }); totalCompactionsCompleted = Metrics.meter(factory.createMetricName("TotalCompactionsCompleted")); bytesCompacted = Metrics.counter(factory.createMetricName("BytesCompacted")); + compressedBytesCompacted = Metrics.counter(factory.createMetricName("CompressedBytesCompacted")); bytesCompactedThroughput = Metrics.meter(factory.createMetricName("BytesCompactedThroughput")); // compaction failure metrics diff --git a/src/java/org/apache/cassandra/streaming/StreamSession.java b/src/java/org/apache/cassandra/streaming/StreamSession.java index 12997e94f5..9caf0b706b 100644 --- a/src/java/org/apache/cassandra/streaming/StreamSession.java +++ b/src/java/org/apache/cassandra/streaming/StreamSession.java @@ -969,7 +969,7 @@ public class StreamSession for (FileStore fs : allWriteableFileStores) newStreamBytesToWritePerFileStore.merge(fs, totalBytesInPerFileStore, Long::sum); } - Map totalCompactionWriteRemaining = Directories.perFileStore(CompactionManager.instance.active.estimatedRemainingWriteBytes(), + Map totalCompactionWriteRemaining = Directories.perFileStore(CompactionManager.instance.active.estimatedRemainingWriteToDiskBytes(), fileStoreMapper); long totalStreamRemaining = StreamManager.instance.getTotalRemainingOngoingBytes(); long totalBytesStreamRemainingPerFileStore = totalStreamRemaining / Math.max(1, allFileStores.size()); diff --git a/src/java/org/apache/cassandra/tools/NodeProbe.java b/src/java/org/apache/cassandra/tools/NodeProbe.java index 5e9fee058d..353324e0ea 100644 --- a/src/java/org/apache/cassandra/tools/NodeProbe.java +++ b/src/java/org/apache/cassandra/tools/NodeProbe.java @@ -2225,6 +2225,7 @@ public class NodeProbe implements AutoCloseable switch(metricName) { case "BytesCompacted": + case "CompressedBytesCompacted": case "CompactionsAborted": case "CompactionsReduced": case "SSTablesDroppedFromCompaction": diff --git a/src/java/org/apache/cassandra/tools/nodetool/CompactionStats.java b/src/java/org/apache/cassandra/tools/nodetool/CompactionStats.java index 4dab0158a1..7ca3ffe04a 100644 --- a/src/java/org/apache/cassandra/tools/nodetool/CompactionStats.java +++ b/src/java/org/apache/cassandra/tools/nodetool/CompactionStats.java @@ -89,11 +89,17 @@ public class CompactionStats extends AbstractCommand tableBuilder.add("compactions completed", String.valueOf(totalCompactionsCompletedMetrics.getCount())); CassandraMetricsRegistry.JmxCounterMBean bytesCompacted = (CassandraMetricsRegistry.JmxCounterMBean) probe.getCompactionMetric("BytesCompacted"); + CassandraMetricsRegistry.JmxCounterMBean compressedBytesCompacted = (CassandraMetricsRegistry.JmxCounterMBean) probe.getCompactionMetric("CompressedBytesCompacted"); if (humanReadable) + { tableBuilder.add("data compacted", FileUtils.stringifyFileSize(Double.parseDouble(Long.toString(bytesCompacted.getCount())))); + tableBuilder.add("compressed data compacted", FileUtils.stringifyFileSize(Double.parseDouble(Long.toString(compressedBytesCompacted.getCount())))); + } else + { tableBuilder.add("data compacted", Long.toString(bytesCompacted.getCount())); - + tableBuilder.add("compressed data compacted", Long.toString(compressedBytesCompacted.getCount())); + } CassandraMetricsRegistry.JmxCounterMBean compactionsAborted = (CassandraMetricsRegistry.JmxCounterMBean) probe.getCompactionMetric("CompactionsAborted"); tableBuilder.add("compactions aborted", Long.toString(compactionsAborted.getCount())); @@ -132,13 +138,14 @@ public class CompactionStats extends AbstractCommand long remainingBytes = 0; if (vtableOutput) - table.add("keyspace", "table", "task id", "completion ratio", "kind", "progress", "sstables", "total", "unit", "target directory"); + table.add("keyspace", "table", "task id", "completion ratio", "kind", "progress", "sstables", "total", "total compressed", "unit", "target directory"); else table.add("id", "compaction type", "keyspace", "table", "completed", "total", "unit", "progress"); for (Map c : compactions) { long total = Long.parseLong(c.get(CompactionInfo.TOTAL)); + String totalCompressedValue = c.get(CompactionInfo.TOTAL_COMPRESSED); long completed = Long.parseLong(c.get(CompactionInfo.COMPLETED)); String taskType = c.get(CompactionInfo.TASK_TYPE); String keyspace = c.get(CompactionInfo.KEYSPACE); @@ -148,12 +155,22 @@ public class CompactionStats extends AbstractCommand String[] tables = c.get(CompactionInfo.SSTABLES).split(","); String progressStr = toFileSize ? FileUtils.stringifyFileSize(completed) : Long.toString(completed); String totalStr = toFileSize ? FileUtils.stringifyFileSize(total) : Long.toString(total); + String totalCompressedStr; + if (totalCompressedValue != null) + { + long totalCompressed = Long.parseLong(totalCompressedValue); + totalCompressedStr = toFileSize ? FileUtils.stringifyFileSize(totalCompressed) : Long.toString(totalCompressed); + } + else + { + totalCompressedStr = "n/a"; + } String percentComplete = total == 0 ? "n/a" : new DecimalFormat("0.00").format((double) completed / total * 100) + '%'; String id = c.get(CompactionInfo.COMPACTION_ID); if (vtableOutput) { String targetDirectory = c.get(CompactionInfo.TARGET_DIRECTORY); - table.add(keyspace, columnFamily, id, percentComplete, taskType, progressStr, String.valueOf(tables.length), totalStr, unit, targetDirectory); + table.add(keyspace, columnFamily, id, percentComplete, taskType, progressStr, String.valueOf(tables.length), totalStr, totalCompressedStr, unit, targetDirectory); } else table.add(id, taskType, keyspace, columnFamily, progressStr, totalStr, unit, percentComplete); @@ -172,4 +189,4 @@ public class CompactionStats extends AbstractCommand table.printTo(out); } -} \ No newline at end of file +} diff --git a/test/distributed/org/apache/cassandra/distributed/test/CompactionDiskSpaceTest.java b/test/distributed/org/apache/cassandra/distributed/test/CompactionDiskSpaceTest.java index 36c5163e53..7ee7c9cc4a 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/CompactionDiskSpaceTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/CompactionDiskSpaceTest.java @@ -115,7 +115,7 @@ public class CompactionDiskSpaceTest extends TestBaseImpl public static void install(ClassLoader cl, Integer node) { new ByteBuddy().rebase(ActiveCompactions.class) - .method(named("estimatedRemainingWriteBytes")) + .method(named("estimatedRemainingWriteToDiskBytes")) .intercept(MethodDelegation.to(BB.class)) .make() .load(cl, ClassLoadingStrategy.Default.INJECTION); @@ -127,7 +127,7 @@ public class CompactionDiskSpaceTest extends TestBaseImpl .load(cl, ClassLoadingStrategy.Default.INJECTION); } - public static Map estimatedRemainingWriteBytes() + public static Map estimatedRemainingWriteToDiskBytes() { if (sstableDir != null) return ImmutableMap.of(sstableDir, estimatedRemaining.get()); diff --git a/test/distributed/org/apache/cassandra/distributed/test/SecondaryIndexCompactionTest.java b/test/distributed/org/apache/cassandra/distributed/test/SecondaryIndexCompactionTest.java index 9d168145c5..35d47bc4bd 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/SecondaryIndexCompactionTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/SecondaryIndexCompactionTest.java @@ -60,7 +60,7 @@ public class SecondaryIndexCompactionTest extends TestBaseImpl // emulate ongoing index compaction: CompactionInfo.Holder h = new MockHolder(i.getIndexCfs().metadata(), idxSSTables); CompactionManager.instance.active.beginCompaction(h); - CompactionManager.instance.active.estimatedRemainingWriteBytes(); + CompactionManager.instance.active.estimatedRemainingWriteToDiskBytes(); CompactionManager.instance.active.finishCompaction(h); }); } @@ -79,7 +79,7 @@ public class SecondaryIndexCompactionTest extends TestBaseImpl @Override public CompactionInfo getCompactionInfo() { - return new CompactionInfo(metadata, OperationType.COMPACTION, 0, 1000, nextTimeUUID(), sstables); + return new CompactionInfo(metadata, OperationType.COMPACTION, 0, 1000, 300, nextTimeUUID(), sstables); } @Override @@ -88,4 +88,4 @@ public class SecondaryIndexCompactionTest extends TestBaseImpl return false; } } -} \ No newline at end of file +} diff --git a/test/distributed/org/apache/cassandra/distributed/test/StreamsDiskSpaceTest.java b/test/distributed/org/apache/cassandra/distributed/test/StreamsDiskSpaceTest.java index f2a20fdf8e..9f7336444f 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/StreamsDiskSpaceTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/StreamsDiskSpaceTest.java @@ -73,7 +73,7 @@ public class StreamsDiskSpaceTest extends TestBaseImpl .withConfig(config -> config.set("hinted_handoff_enabled", false) .with(GOSSIP) .with(NETWORK)) - .withInstanceInitializer((cl, id) -> BB.doInstall(cl, id, ActiveCompactions.class, "estimatedRemainingWriteBytes")) + .withInstanceInitializer((cl, id) -> BB.doInstall(cl, id, ActiveCompactions.class, "estimatedRemainingWriteToDiskBytes")) .start())) { cluster.schemaChange("create table " + KEYSPACE + ".tbl (id int primary key, t int) with compaction={'class': 'SizeTieredCompactionStrategy'}"); @@ -139,7 +139,7 @@ public class StreamsDiskSpaceTest extends TestBaseImpl return ongoing.get(); } - public static Map estimatedRemainingWriteBytes() + public static Map estimatedRemainingWriteToDiskBytes() { Map ret = new HashMap<>(); if (datadir != null) diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionInfoTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionInfoTest.java index c4ae274809..2fd20d25aa 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionInfoTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionInfoTest.java @@ -39,7 +39,7 @@ public class CompactionInfoTest extends AbstractPendingAntiCompactionTest { ColumnFamilyStore cfs = MockSchema.newCFS(); TimeUUID expectedTaskId = nextTimeUUID(); - CompactionInfo compactionInfo = new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, 0, 1000, expectedTaskId, new ArrayList<>()); + CompactionInfo compactionInfo = new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, 0, 1000, 1000, expectedTaskId, new ArrayList<>()); Assertions.assertThat(compactionInfo.toString()) .contains(expectedTaskId.toString()); } @@ -50,7 +50,7 @@ public class CompactionInfoTest extends AbstractPendingAntiCompactionTest UUID tableId = UUID.randomUUID(); TimeUUID taskId = nextTimeUUID(); ColumnFamilyStore cfs = MockSchema.newCFS(builder -> builder.id(TableId.fromUUID(tableId))); - CompactionInfo compactionInfo = new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, 0, 1000, taskId, new ArrayList<>()); + CompactionInfo compactionInfo = new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, 0, 1000, 300, taskId, new ArrayList<>()); Assertions.assertThat(compactionInfo.toString()) .isEqualTo("Compaction(%s, 0 / 1000 bytes)@%s(mockks, mockcf1)", taskId, tableId); } diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionsCQLTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionsCQLTest.java index 57113b4feb..3bd133ca8d 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionsCQLTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionsCQLTest.java @@ -85,7 +85,7 @@ public class CompactionsCQLTest extends CQLTester public void before() throws IOException { strategy = DatabaseDescriptor.getCorruptedTombstoneStrategy(); - + CommitLog.instance.resetUnsafe(true); } @@ -877,7 +877,9 @@ public class CompactionsCQLTest extends CQLTester execute("insert into %s (id, i) values (?,?)", i, i); getCurrentColumnFamilyStore().forceBlockingFlush(ColumnFamilyStore.FlushReason.UNIT_TESTS); } - CompactionInfo.Holder holder = holder(OperationType.COMPACTION); + // When we have an existing compaction with sstables of total size more than double the available space, + // we should not be able to then run a major compaction + CompactionInfo.Holder holder = holder(OperationType.COMPACTION, 2); CompactionManager.instance.active.beginCompaction(holder); try { @@ -893,7 +895,19 @@ public class CompactionsCQLTest extends CQLTester CompactionManager.instance.active.finishCompaction(holder); } // don't block compactions if there is a huge validation - holder = holder(OperationType.VALIDATION); + holder = holder(OperationType.VALIDATION, 2); + CompactionManager.instance.active.beginCompaction(holder); + try + { + getCurrentColumnFamilyStore().forceMajorCompaction(); + } + finally + { + CompactionManager.instance.active.finishCompaction(holder); + } + + // Should be able to run when the sstables in question are 90% of the total available space + holder = holder(OperationType.COMPACTION, 0.9); CompactionManager.instance.active.beginCompaction(holder); try { @@ -905,7 +919,7 @@ public class CompactionsCQLTest extends CQLTester } } - private CompactionInfo.Holder holder(OperationType opType) + private CompactionInfo.Holder holder(OperationType opType, double availableSpaceMultiplier) { CompactionInfo.Holder holder = new CompactionInfo.Holder() { @@ -915,12 +929,17 @@ public class CompactionsCQLTest extends CQLTester for (File f : getCurrentColumnFamilyStore().getDirectories().getCFDirectories()) availableSpace += PathUtils.tryGetSpace(f.toPath(), FileStore::getUsableSpace); + Set liveSSTables = getCurrentColumnFamilyStore().getLiveSSTables(); + long totalDiskUsage = (long)(availableSpace * availableSpaceMultiplier); + // Arbitrary compression ratio of 3.4 + long totalUncompressedSize = (long) ((double) totalDiskUsage * 3.4); return new CompactionInfo(getCurrentColumnFamilyStore().metadata(), opType, +0, - +availableSpace * 2, + totalUncompressedSize, + totalDiskUsage, nextTimeUUID(), - getCurrentColumnFamilyStore().getLiveSSTables()); + liveSSTables); } public boolean isGlobal() diff --git a/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java b/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java index f39eeaf1be..387cd4f1b4 100644 --- a/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java +++ b/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java @@ -609,7 +609,7 @@ public class PendingAntiCompactionTest extends AbstractPendingAntiCompactionTest { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.ANTICOMPACTION, 0, 1000, nextTimeUUID(), compacting); + return new CompactionInfo(cfs.metadata(), OperationType.ANTICOMPACTION, 0, 1000, 1000, nextTimeUUID(), compacting); } public boolean isGlobal() @@ -650,7 +650,7 @@ public class PendingAntiCompactionTest extends AbstractPendingAntiCompactionTest { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.ANTICOMPACTION, 0, 0, nextTimeUUID(), cfs.getLiveSSTables()); + return new CompactionInfo(cfs.metadata(), OperationType.ANTICOMPACTION, 0, 0, 0, nextTimeUUID(), cfs.getLiveSSTables()); } public boolean isGlobal() @@ -703,7 +703,7 @@ public class PendingAntiCompactionTest extends AbstractPendingAntiCompactionTest { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.ANTICOMPACTION, 0, 0, nextTimeUUID(), cfs.getLiveSSTables()); + return new CompactionInfo(cfs.metadata(), OperationType.ANTICOMPACTION, 0, 0, 0, nextTimeUUID(), cfs.getLiveSSTables()); } public boolean isGlobal() diff --git a/test/unit/org/apache/cassandra/db/virtual/SSTableTasksTableTest.java b/test/unit/org/apache/cassandra/db/virtual/SSTableTasksTableTest.java index 3a7ec83894..f215bf4d3c 100644 --- a/test/unit/org/apache/cassandra/db/virtual/SSTableTasksTableTest.java +++ b/test/unit/org/apache/cassandra/db/virtual/SSTableTasksTableTest.java @@ -70,6 +70,7 @@ public class SSTableTasksTableTest extends CQLTester long bytesCompacted = 123; long bytesTotal = 123456; + long totalCompressedBytes = 112233; TimeUUID compactionId = nextTimeUUID(); List sstables = IntStream.range(0, 10) .mapToObj(i -> MockSchema.sstable(i, i * 10L, i * 10L + 9, cfs)) @@ -81,7 +82,7 @@ public class SSTableTasksTableTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, compactionId, sstables, directory); + return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables, directory); } public boolean isGlobal() @@ -94,7 +95,7 @@ public class SSTableTasksTableTest extends CQLTester UntypedResultSet result = execute("SELECT * FROM vts.sstable_tasks"); assertRows(result, row(CQLTester.KEYSPACE, currentTable(), compactionId, 1.0 * bytesCompacted / bytesTotal, toLowerCaseLocalized(OperationType.COMPACTION.toString()), bytesCompacted, sstables.size(), - directory, bytesTotal, CompactionInfo.Unit.BYTES.toString())); + directory, bytesTotal, totalCompressedBytes, CompactionInfo.Unit.BYTES.toString())); CompactionManager.instance.active.finishCompaction(compactionHolder); result = execute("SELECT * FROM vts.sstable_tasks"); diff --git a/test/unit/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryManagerTest.java b/test/unit/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryManagerTest.java index b7596ae07b..366f4d20ad 100644 --- a/test/unit/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryManagerTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/indexsummary/IndexSummaryManagerTest.java @@ -655,7 +655,7 @@ public class IndexSummaryManagerTest sstables = IntStream.range(0, 10) .mapToObj(i -> MockSchema.sstable(i, i * 10L, i * 10L + 9, cfs)) @@ -65,7 +66,7 @@ public class CompactionStatsTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, compactionId, sstables); + return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables); } public boolean isGlobal() @@ -87,6 +88,7 @@ public class CompactionStatsTest extends CQLTester assertThat(stdout).containsPattern("pending tasks\\s+[0-9]*"); assertThat(stdout).containsPattern("compactions completed\\s+[0-9]*"); assertThat(stdout).containsPattern("data compacted\\s+[0-9]*"); + assertThat(stdout).containsPattern("compressed data compacted\\s+[0-9]*"); assertThat(stdout).containsPattern("compactions aborted\\s+[0-9]*"); assertThat(stdout).containsPattern("compactions reduced\\s+[0-9]*"); assertThat(stdout).containsPattern("sstables dropped from compaction\\s+[0-9]*"); @@ -109,6 +111,7 @@ public class CompactionStatsTest extends CQLTester long bytesCompacted = 123; long bytesTotal = 123456; + long totalCompressedBytes = 112233; TimeUUID compactionId = nextTimeUUID(); List sstables = IntStream.range(0, 10) .mapToObj(i -> MockSchema.sstable(i, i * 10L, i * 10L + 9, cfs)) @@ -118,7 +121,7 @@ public class CompactionStatsTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, compactionId, sstables, targetDirectory); + return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables, targetDirectory); } public boolean isGlobal() @@ -131,7 +134,7 @@ public class CompactionStatsTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.CLEANUP, bytesCompacted, bytesTotal, compactionId, sstables); + return new CompactionInfo(cfs.metadata(), OperationType.CLEANUP, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables); } public boolean isGlobal() @@ -143,16 +146,16 @@ public class CompactionStatsTest extends CQLTester CompactionManager.instance.active.beginCompaction(compactionHolder); CompactionManager.instance.active.beginCompaction(nonCompactionHolder); String stdout = waitForNumberOfPendingTasks(2, "compactionstats", "-V"); - assertThat(stdout).containsPattern("keyspace\\s+table\\s+task id\\s+completion ratio\\s+kind\\s+progress\\s+sstables\\s+total\\s+unit\\s+target directory"); - String expectedStatsPattern = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", + assertThat(stdout).containsPattern("keyspace\\s+table\\s+task id\\s+completion ratio\\s+kind\\s+progress\\s+sstables\\s+total\\s+total compressed\\s+unit\\s+target directory"); + String expectedStatsPattern = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", CQLTester.KEYSPACE, currentTable(), compactionId, (double) bytesCompacted / bytesTotal * 100, - OperationType.COMPACTION, bytesCompacted, sstables.size(), bytesTotal, CompactionInfo.Unit.BYTES, + OperationType.COMPACTION, bytesCompacted, sstables.size(), bytesTotal, totalCompressedBytes, CompactionInfo.Unit.BYTES, targetDirectory); assertThat(stdout).containsPattern(expectedStatsPattern); - String expectedStatsPatternForNonCompaction = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", + String expectedStatsPatternForNonCompaction = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", CQLTester.KEYSPACE, currentTable(), compactionId, (double) bytesCompacted / bytesTotal * 100, - OperationType.COMPACTION, bytesCompacted, sstables.size(), bytesTotal, CompactionInfo.Unit.BYTES); + OperationType.COMPACTION, bytesCompacted, sstables.size(), bytesTotal, totalCompressedBytes, CompactionInfo.Unit.BYTES); assertThat(stdout).containsPattern(expectedStatsPatternForNonCompaction); CompactionManager.instance.active.finishCompaction(compactionHolder); @@ -168,6 +171,7 @@ public class CompactionStatsTest extends CQLTester long bytesCompacted = 123; long bytesTotal = 123456; + long totalCompressedBytes = 112233; TimeUUID compactionId = nextTimeUUID(); List sstables = IntStream.range(0, 10) .mapToObj(i -> MockSchema.sstable(i, i * 10L, i * 10L + 9, cfs)) @@ -176,7 +180,7 @@ public class CompactionStatsTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, compactionId, sstables); + return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables); } public boolean isGlobal() @@ -205,6 +209,7 @@ public class CompactionStatsTest extends CQLTester long bytesCompacted = 123; long bytesTotal = 123456; + long totalCompressedBytes = 112233; TimeUUID compactionId = nextTimeUUID(); List sstables = IntStream.range(0, 10) .mapToObj(i -> MockSchema.sstable(i, i * 10L, i * 10L + 9, cfs)) @@ -214,7 +219,7 @@ public class CompactionStatsTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, compactionId, sstables, targetDirectory); + return new CompactionInfo(cfs.metadata(), OperationType.COMPACTION, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables, targetDirectory); } public boolean isGlobal() @@ -227,7 +232,7 @@ public class CompactionStatsTest extends CQLTester { public CompactionInfo getCompactionInfo() { - return new CompactionInfo(cfs.metadata(), OperationType.CLEANUP, bytesCompacted, bytesTotal, compactionId, sstables); + return new CompactionInfo(cfs.metadata(), OperationType.CLEANUP, bytesCompacted, bytesTotal, totalCompressedBytes, compactionId, sstables); } public boolean isGlobal() @@ -239,15 +244,15 @@ public class CompactionStatsTest extends CQLTester CompactionManager.instance.active.beginCompaction(compactionHolder); CompactionManager.instance.active.beginCompaction(nonCompactionHolder); String stdout = waitForNumberOfPendingTasks(2, "compactionstats", "--vtable", "--human-readable"); - assertThat(stdout).containsPattern("keyspace\\s+table\\s+task id\\s+completion ratio\\s+kind\\s+progress\\s+sstables\\s+total\\s+unit\\s+target directory"); - String expectedStatsPattern = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", + assertThat(stdout).containsPattern("keyspace\\s+table\\s+task id\\s+completion ratio\\s+kind\\s+progress\\s+sstables\\s+total\\s+total compressed\\s+unit\\s+target directory"); + String expectedStatsPattern = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", CQLTester.KEYSPACE, currentTable(), compactionId, (double) bytesCompacted / bytesTotal * 100, - OperationType.COMPACTION, "123 bytes", sstables.size(), "120.56 KiB", CompactionInfo.Unit.BYTES, + OperationType.COMPACTION, "123 bytes", sstables.size(), "120.56 KiB", "109.6 KiB", CompactionInfo.Unit.BYTES, targetDirectory); assertThat(stdout).containsPattern(expectedStatsPattern); - String expectedStatsPatternForNonCompaction = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", + String expectedStatsPatternForNonCompaction = String.format("%s\\s+%s\\s+%s\\s+%.2f%%\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s\\s+%s", CQLTester.KEYSPACE, currentTable(), compactionId, (double) bytesCompacted / bytesTotal * 100, - OperationType.CLEANUP, "123 bytes", sstables.size(), "120.56 KiB", CompactionInfo.Unit.BYTES); + OperationType.CLEANUP, "123 bytes", sstables.size(), "120.56 KiB", "109.6 KiB", CompactionInfo.Unit.BYTES); assertThat(stdout).containsPattern(expectedStatsPatternForNonCompaction); CompactionManager.instance.active.finishCompaction(compactionHolder); @@ -277,4 +282,4 @@ public class CompactionStatsTest extends CQLTester return stdout.get(); } -} \ No newline at end of file +}