This is an automated email from the ASF dual-hosted git repository. marcuse pushed a commit to branch cassandra-4.0 in repository https://gitbox.apache.org/repos/asf/cassandra.git
commit 702824a99ddf890603113992b09bf5c7780f6a6d Merge: 894d33c e53ad64 Author: Marcus Eriksson <[email protected]> AuthorDate: Thu Feb 17 10:33:48 2022 +0100 Merge branch 'cassandra-3.11' into cassandra-4.0 CHANGES.txt | 2 + .../cassandra/db/compaction/CompactionTask.java | 6 +- .../db/compaction/LeveledCompactionTask.java | 45 +++++- .../compaction/writers/CompactionAwareWriter.java | 7 +- .../writers/MajorLeveledCompactionWriter.java | 6 + .../compaction/writers/MaxSSTableSizeWriter.java | 6 + .../SplittingSizeTieredCompactionWriter.java | 8 +- .../compaction/LeveledCompactionStrategyTest.java | 161 ++++++++++++++++++++- .../{ => writers}/CompactionAwareWriterTest.java | 48 +++++- 9 files changed, 271 insertions(+), 18 deletions(-) diff --cc CHANGES.txt index c64b602,513a8af..4fa9ba3 --- a/CHANGES.txt +++ b/CHANGES.txt @@@ -1,46 -1,21 +1,48 @@@ -3.11.13 +4.0.4 + Merged from 3.0: * Lazy transaction log replica creation allows incorrect replica content divergence during anticompaction (CASSANDRA-17273) + * LeveledCompactionStrategy disk space check improvements (CASSANDRA-17272) +4.0.3 + * Deprecate otc_coalescing_strategy, otc_coalescing_window_us, otc_coalescing_enough_coalesced_messages, + otc_backlog_expiration_interval_ms (CASSANDRA-17377) + * Improve start up processing of Incremental Repair information read from system.repairs (CASSANDRA-17342) -3.11.12 - * Upgrade snakeyaml to 1.26 in 3.11 (CASSANDRA=17028) +4.0.2 + * Full Java 11 support (CASSANDRA-16894) + * Remove unused 'geomet' package from cqlsh path (CASSANDRA-17271) + * Removed unused 'cql' dependency (CASSANDRA-17247) + * Don't block gossip when clearing repair snapshots (CASSANDRA-17168) + * Deduplicate warnings for deprecated parameters (changed names) (CASSANDRA-17160) + * Update ant-junit to version 1.10.12 (CASSANDRA-17218) + * Add droppable tombstone metrics to nodetool tablestats (CASSANDRA-16308) + * Fix disk failure triggered when enabling FQL on an unclean directory (CASSANDRA-17136) + * Fixed broken classpath when multiple jars in build directory (CASSANDRA-17129) + * DebuggableThreadPoolExecutor does not propagate client warnings (CASSANDRA-17072) + * internode_send_buff_size_in_bytes and internode_recv_buff_size_in_bytes have new names. Backward compatibility with the old names added (CASSANDRA-17141) + * Remove unused configuration parameters from cassandra.yaml (CASSANDRA-17132) + * Queries performed with NODE_LOCAL consistency level do not update request metrics (CASSANDRA-17052) + * Fix multiple full sources can be select unexpectedly for bootstrap streaming (CASSANDRA-16945) + * Fix cassandra.yaml formatting of parameters (CASSANDRA-17131) + * Add backward compatibility for CQLSSTableWriter Date fields (CASSANDRA-17117) + * Push initial client connection messages to trace (CASSANDRA-17038) + * Correct the internode message timestamp if sending node has wrapped (CASSANDRA-16997) + * Avoid race causing us to return null in RangesAtEndpoint (CASSANDRA-16965) + * Avoid rewriting all sstables during cleanup when transient replication is enabled (CASSANDRA-16966) + * Prevent CQLSH from failure on Python 3.10 (CASSANDRA-16987) + * Avoid trying to acquire 0 permits from the rate limiter when taking snapshot (CASSANDRA-16872) + * Upgrade Caffeine to 2.5.6 (CASSANDRA-15153) + * Include SASI components to snapshots (CASSANDRA-15134) + * Fix missed wait latencies in the output of `nodetool tpstats -F` (CASSANDRA-16938) + * Remove all the state pollution between tests in SSTableReaderTest (CASSANDRA-16888) + * Delay auth setup until after gossip has settled to avoid unavailables on startup (CASSANDRA-16783) + * Fix clustering order logic in CREATE MATERIALIZED VIEW (CASSANDRA-16898) + * org.apache.cassandra.db.rows.ArrayCell#unsharedHeapSizeExcludingData includes data twice (CASSANDRA-16900) + * Exclude Jackson 1.x transitive dependency of hadoop* provided dependencies (CASSANDRA-16854) +Merged from 3.11: * Add key validation to ssstablescrub (CASSANDRA-16969) * Update Jackson from 2.9.10 to 2.12.5 (CASSANDRA-16851) - * Include SASI components to snapshots (CASSANDRA-15134) * Make assassinate more resilient to missing tokens (CASSANDRA-16847) - * Exclude Jackson 1.x transitive dependency of hadoop* provided dependencies (CASSANDRA-16854) - * Validate SASI tokenizer options before adding index to schema (CASSANDRA-15135) - * Fixup scrub output when no data post-scrub and clear up old use of row, which really means partition (CASSANDRA-16835) - * Fix ant-junit dependency issue (CASSANDRA-16827) - * Reduce thread contention in CommitLogSegment and HintsBuffer (CASSANDRA-16072) - * Avoid sending CDC column if not enabled (CASSANDRA-16770) Merged from 3.0: * Fix conversion from megabits to bytes in streaming rate limiter (CASSANDRA-17243) * Upgrade logback to 1.2.9 (CASSANDRA-17204) diff --cc src/java/org/apache/cassandra/db/compaction/CompactionTask.java index 74bf246,34c249c..3b0e172 --- a/src/java/org/apache/cassandra/db/compaction/CompactionTask.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionTask.java @@@ -83,13 -91,13 +83,13 @@@ public class CompactionTask extends Abs if (partialCompactionsAcceptable() && transaction.originals().size() > 1) { // Try again w/o the largest one. - logger.warn("insufficient space to compact all requested files. {}MB required, {}", + logger.warn("insufficient space to compact all requested files. {}MB required, {} for compaction {}", (float) expectedSize / 1024 / 1024, - StringUtils.join(transaction.originals(), ", ")); - + StringUtils.join(transaction.originals(), ", "), + transaction.opId()); // Note that we have removed files that are still marked as compacting. // This suboptimal but ok since the caller will unmark all the sstables at the end. - SSTableReader removedSSTable = cfs.getMaxSizeFile(transaction.originals()); + SSTableReader removedSSTable = cfs.getMaxSizeFile(nonExpiredSSTables); transaction.cancel(removedSSTable); return true; } diff --cc src/java/org/apache/cassandra/db/compaction/LeveledCompactionTask.java index c633937,c40582c..5b94c54 --- a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionTask.java +++ b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionTask.java @@@ -62,4 -63,46 +63,46 @@@ public class LeveledCompactionTask exte { return level; } + + @Override - public boolean reduceScopeForLimitedSpace(long expectedSize) ++ public boolean reduceScopeForLimitedSpace(Set<SSTableReader> nonExpiredSSTables, long expectedSize) + { + if (transaction.originals().size() > 1 && level <= 1) + { + // Try again w/o the largest one. + logger.warn("insufficient space to do L0 -> L{} compaction. {}MiB required, {} for compaction {}", + level, + (float) expectedSize / 1024 / 1024, + transaction.originals() + .stream() + .map(sstable -> String.format("%s (level=%s, size=%s)", sstable, sstable.getSSTableLevel(), sstable.onDiskLength())) + .collect(Collectors.joining(",")), + transaction.opId()); + // Note that we have removed files that are still marked as compacting. + // This suboptimal but ok since the caller will unmark all the sstables at the end. + int l0SSTableCount = 0; + SSTableReader largestL0SSTable = null; - for (SSTableReader sstable : transaction.originals()) ++ for (SSTableReader sstable : nonExpiredSSTables) + { + if (sstable.getSSTableLevel() == 0) + { + l0SSTableCount++; + if (largestL0SSTable == null || sstable.onDiskLength() > largestL0SSTable.onDiskLength()) + largestL0SSTable = sstable; + } + } + // no point doing a L0 -> L{0,1} compaction if we have cancelled all L0 sstables + if (largestL0SSTable != null && l0SSTableCount > 1) + { + logger.info("Removing {} (level={}, size={}) from compaction {}", + largestL0SSTable, + largestL0SSTable.getSSTableLevel(), + largestL0SSTable.onDiskLength(), + transaction.opId()); + transaction.cancel(largestL0SSTable); + return true; + } + } + return false; + } } diff --cc src/java/org/apache/cassandra/db/compaction/writers/MajorLeveledCompactionWriter.java index 1c53600,f1326e9..b7fb881 --- a/src/java/org/apache/cassandra/db/compaction/writers/MajorLeveledCompactionWriter.java +++ b/src/java/org/apache/cassandra/db/compaction/writers/MajorLeveledCompactionWriter.java @@@ -104,5 -115,12 +104,11 @@@ public class MajorLeveledCompactionWrit txn)); partitionsWritten = 0; sstablesWritten = 0; - } + + @Override + protected long getExpectedWriteSize() + { + return Math.min(maxSSTableSize, super.getExpectedWriteSize()); + } } diff --cc test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java index 41039d5,91b9b3b..00b460e --- a/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/LeveledCompactionStrategyTest.java @@@ -32,21 -33,28 +33,23 @@@ import java.util.Set import java.util.UUID; import java.util.stream.Collectors; - import org.junit.Assert; -import junit.framework.Assert; ++import com.google.common.collect.Iterables; ++import com.google.common.collect.Sets; import org.junit.After; ++import org.junit.Assert; import org.junit.Before; import org.junit.BeforeClass; - -import com.google.common.collect.Iterables; -import com.google.common.collect.Sets; -import org.apache.cassandra.db.lifecycle.LifecycleTransaction; import org.junit.Test; -import org.junit.runner.RunWith; -- import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.MockSchema; -import org.apache.cassandra.OrderedJUnit4ClassRunner; import org.apache.cassandra.SchemaLoader; --import org.apache.cassandra.Util; import org.apache.cassandra.UpdateBuilder; ++import org.apache.cassandra.Util; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.db.Keyspace; ++import org.apache.cassandra.db.lifecycle.LifecycleTransaction; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; @@@ -54,22 -62,23 +57,25 @@@ import org.apache.cassandra.io.sstable. import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.notifications.SSTableAddedNotification; import org.apache.cassandra.notifications.SSTableRepairStatusChanged; - import org.apache.cassandra.repair.ValidationManager; - import org.apache.cassandra.schema.MockSchema; - import org.apache.cassandra.streaming.PreviewKind; import org.apache.cassandra.repair.RepairJobDesc; ++import org.apache.cassandra.repair.ValidationManager; import org.apache.cassandra.repair.Validator; import org.apache.cassandra.schema.CompactionParams; import org.apache.cassandra.schema.KeyspaceParams; ++import org.apache.cassandra.schema.MockSchema; import org.apache.cassandra.service.ActiveRepairService; ++import org.apache.cassandra.streaming.PreviewKind; import org.apache.cassandra.utils.FBUtilities; + import org.apache.cassandra.utils.Pair; import static java.util.Collections.singleton; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; + import static org.junit.Assert.assertNotNull; + import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; -@RunWith(OrderedJUnit4ClassRunner.class) public class LeveledCompactionStrategyTest { private static final Logger logger = LoggerFactory.getLogger(LeveledCompactionStrategyTest.class); @@@ -849,4 -732,147 +855,147 @@@ assertThat(compactionCandidates).containsAll(sstablesOnL7); assertThat(compactionCandidates).doesNotContainAnyElementsOf(sstablesOnL8); } + + @Test + public void testReduceScopeL0L1() throws IOException + { + ColumnFamilyStore cfs = MockSchema.newCFS(); + Map<String, String> localOptions = new HashMap<>(); + localOptions.put("class", "LeveledCompactionStrategy"); + localOptions.put("sstable_size_in_mb", "1"); + cfs.setCompactionParameters(localOptions); + List<SSTableReader> l1sstables = new ArrayList<>(); + for (int i = 0; i < 10; i++) + { + SSTableReader l1sstable = MockSchema.sstable(i, 1 * 1024 * 1024, cfs); + l1sstable.descriptor.getMetadataSerializer().mutateLevel(l1sstable.descriptor, 1); + l1sstable.reloadSSTableMetadata(); + l1sstables.add(l1sstable); + } + List<SSTableReader> l0sstables = new ArrayList<>(); + for (int i = 10; i < 20; i++) + l0sstables.add(MockSchema.sstable(i, (i + 1) * 1024 * 1024, cfs)); + try (LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.COMPACTION, Iterables.concat(l0sstables, l1sstables))) + { ++ Set<SSTableReader> nonExpired = Sets.difference(txn.originals(), Collections.emptySet()); + CompactionTask task = new LeveledCompactionTask(cfs, txn, 1, 0, 1024*1024, false); + SSTableReader lastRemoved = null; + boolean removed = true; + for (int i = 0; i < l0sstables.size(); i++) + { + Set<SSTableReader> before = new HashSet<>(txn.originals()); - removed = task.reduceScopeForLimitedSpace(0); ++ removed = task.reduceScopeForLimitedSpace(nonExpired, 0); + SSTableReader removedSSTable = Iterables.getOnlyElement(Sets.difference(before, txn.originals()), null); + if (removed) + { + assertNotNull(removedSSTable); + assertTrue(lastRemoved == null || removedSSTable.onDiskLength() < lastRemoved.onDiskLength()); + assertEquals(0, removedSSTable.getSSTableLevel()); + Pair<Set<SSTableReader>, Set<SSTableReader>> sstables = groupByLevel(txn.originals()); + Set<SSTableReader> l1after = sstables.right; + + assertEquals(l1after, new HashSet<>(l1sstables)); // we don't touch L1 + assertEquals(before.size() - 1, txn.originals().size()); + lastRemoved = removedSSTable; + } + else + { + assertNull(removedSSTable); + Pair<Set<SSTableReader>, Set<SSTableReader>> sstables = groupByLevel(txn.originals()); + Set<SSTableReader> l0after = sstables.left; + Set<SSTableReader> l1after = sstables.right; + assertEquals(l1after, new HashSet<>(l1sstables)); // we don't touch L1 + assertEquals(1, l0after.size()); // and we stop reducing once there is a single sstable left + } + } + assertFalse(removed); + } + } + + @Test + public void testReduceScopeL0() + { + + List<SSTableReader> l0sstables = new ArrayList<>(); + for (int i = 10; i < 20; i++) + l0sstables.add(MockSchema.sstable(i, (i + 1) * 1024 * 1024, cfs)); + + try (LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.COMPACTION, l0sstables)) + { + CompactionTask task = new LeveledCompactionTask(cfs, txn, 0, 0, 1024*1024, false); + + SSTableReader lastRemoved = null; + boolean removed = true; + for (int i = 0; i < l0sstables.size(); i++) + { + Set<SSTableReader> before = new HashSet<>(txn.originals()); - removed = task.reduceScopeForLimitedSpace(0); ++ removed = task.reduceScopeForLimitedSpace(before, 0); + SSTableReader removedSSTable = Sets.difference(before, txn.originals()).stream().findFirst().orElse(null); + if (removed) + { + assertNotNull(removedSSTable); + assertTrue(lastRemoved == null || removedSSTable.onDiskLength() < lastRemoved.onDiskLength()); + assertEquals(0, removedSSTable.getSSTableLevel()); + assertEquals(before.size() - 1, txn.originals().size()); + lastRemoved = removedSSTable; + } + else + { + assertNull(removedSSTable); + Pair<Set<SSTableReader>, Set<SSTableReader>> sstables = groupByLevel(txn.originals()); + Set<SSTableReader> l0after = sstables.left; + assertEquals(1, l0after.size()); // and we stop reducing once there is a single sstable left + } + } + assertFalse(removed); + } + } + + @Test + public void testNoHighLevelReduction() throws IOException + { + List<SSTableReader> sstables = new ArrayList<>(); + int i = 1; + for (; i < 5; i++) + { + SSTableReader sstable = MockSchema.sstable(i, (i + 1) * 1024 * 1024, cfs); + sstable.descriptor.getMetadataSerializer().mutateLevel(sstable.descriptor, 1); + sstable.reloadSSTableMetadata(); + sstables.add(sstable); + } + for (; i < 10; i++) + { + SSTableReader sstable = MockSchema.sstable(i, (i + 1) * 1024 * 1024, cfs); + sstable.descriptor.getMetadataSerializer().mutateLevel(sstable.descriptor, 2); + sstable.reloadSSTableMetadata(); + sstables.add(sstable); + } + try (LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.COMPACTION, sstables)) + { + CompactionTask task = new LeveledCompactionTask(cfs, txn, 0, 0, 1024 * 1024, false); - assertFalse(task.reduceScopeForLimitedSpace(0)); - assertEquals(new HashSet<>(sstables), txn.originals()); ++ assertFalse(task.reduceScopeForLimitedSpace(Sets.newHashSet(sstables), 0)); ++ assertEquals(Sets.newHashSet(sstables), txn.originals()); + } + } + + private Pair<Set<SSTableReader>, Set<SSTableReader>> groupByLevel(Iterable<SSTableReader> sstables) + { + Set<SSTableReader> l1after = new HashSet<>(); + Set<SSTableReader> l0after = new HashSet<>(); + for (SSTableReader sstable : sstables) + { + switch (sstable.getSSTableLevel()) + { + case 0: + l0after.add(sstable); + break; + case 1: + l1after.add(sstable); + break; + default: + throw new RuntimeException("only l0 & l1 sstables"); + } + } + return Pair.create(l0after, l1after); + } - } diff --cc test/unit/org/apache/cassandra/db/compaction/writers/CompactionAwareWriterTest.java index 68936f5,c25a7af..5e127dd --- a/test/unit/org/apache/cassandra/db/compaction/writers/CompactionAwareWriterTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/writers/CompactionAwareWriterTest.java @@@ -26,13 -27,13 +27,13 @@@ import org.junit.* import org.apache.cassandra.cql3.CQLTester; import org.apache.cassandra.cql3.QueryProcessor; import org.apache.cassandra.db.*; - import org.apache.cassandra.db.compaction.writers.CompactionAwareWriter; - import org.apache.cassandra.db.compaction.writers.DefaultCompactionWriter; - import org.apache.cassandra.db.compaction.writers.MajorLeveledCompactionWriter; - import org.apache.cassandra.db.compaction.writers.MaxSSTableSizeWriter; - import org.apache.cassandra.db.compaction.writers.SplittingSizeTieredCompactionWriter; + import org.apache.cassandra.db.compaction.AbstractCompactionStrategy; + import org.apache.cassandra.db.compaction.CompactionController; + import org.apache.cassandra.db.compaction.CompactionIterator; + import org.apache.cassandra.db.compaction.OperationType; import org.apache.cassandra.db.lifecycle.LifecycleTransaction; import org.apache.cassandra.io.sstable.format.SSTableReader; ++import org.apache.cassandra.schema.MockSchema; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.UUIDGen; @@@ -165,6 -166,41 +166,41 @@@ public class CompactionAwareWriterTest cfs.truncateBlocking(); } + @Test + public void testMultiDatadirCheck() + { + createTable("create table %s (id int primary key)"); + Directories.DataDirectory [] dataDirs = new Directories.DataDirectory[] { + new MockDataDirectory(new File("/tmp/1")), + new MockDataDirectory(new File("/tmp/2")), + new MockDataDirectory(new File("/tmp/3")), + new MockDataDirectory(new File("/tmp/4")), + new MockDataDirectory(new File("/tmp/5")) + }; + Set<SSTableReader> sstables = new HashSet<>(); + for (int i = 0; i < 100; i++) + sstables.add(MockSchema.sstable(i, 1000, getCurrentColumnFamilyStore())); + - Directories dirs = new Directories(getCurrentColumnFamilyStore().metadata, dataDirs); ++ Directories dirs = new Directories(getCurrentColumnFamilyStore().metadata(), dataDirs); + LifecycleTransaction txn = LifecycleTransaction.offline(OperationType.COMPACTION, sstables); + CompactionAwareWriter writer = new MaxSSTableSizeWriter(getCurrentColumnFamilyStore(), dirs, txn, sstables, 2000, 1); + // init case + writer.maybeSwitchWriter(null); + } + + private static class MockDataDirectory extends Directories.DataDirectory + { + public MockDataDirectory(File location) + { + super(location); + } + + public long getAvailableSpace() + { + return 5000; + } + } + private int compact(ColumnFamilyStore cfs, LifecycleTransaction txn, CompactionAwareWriter writer) { assert txn.originals().size() == 1; --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
