From 62f8399fad0cd0e1bca82d309ea197ef7feb9684 Mon Sep 17 00:00:00 2001 From: Yuqi Yan Date: Fri, 24 Jul 2026 13:29:30 -0700 Subject: [PATCH] Configurable compression algorithm on each level (CASSANDRA-21533) --- CHANGES.txt | 1 + .../compaction/LeveledCompactionStrategy.java | 17 + .../io/sstable/format/big/BigTableWriter.java | 15 +- .../sstable/metadata/MetadataCollector.java | 5 + .../cassandra/schema/CompactionParams.java | 81 ++++ .../io/compress/PerLevelCompressionTest.java | 346 ++++++++++++++++++ 6 files changed, 464 insertions(+), 1 deletion(-) create mode 100644 test/unit/org/apache/cassandra/io/compress/PerLevelCompressionTest.java diff --git a/CHANGES.txt b/CHANGES.txt index fdf3e22695..073244f4ea 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.1.12 + * Configurable compression algorithm on each level (CASSANDRA-21533) * Add Paxos v2 option and informatin in cassandra.yaml (CASSANDRA-21316) * Harden data resurrection startup check with atomic heartbeat file write with fallback (CASSANDRA-21290) Merged from 4.0: diff --git a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java index b2cd3f3152..cf362fced9 100644 --- a/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java +++ b/src/java/org/apache/cassandra/db/compaction/LeveledCompactionStrategy.java @@ -50,6 +50,19 @@ public class LeveledCompactionStrategy extends AbstractCompactionStrategy { private static final Logger logger = LoggerFactory.getLogger(LeveledCompactionStrategy.class); private static final String SSTABLE_SIZE_OPTION = "sstable_size_in_mb"; + + /** + * LeveledCompactionStrategy option that lets each level use a different sstable compression than the + * table default. The value is a JSON object mapping level to a compression definition, e.g. + * {@code {"0":{"class":"LZ4Compressor"},"1":{"class":"ZstdCompressor","compression_level":"9"}}}. Each + * level object is a full, independent compression spec, specified exactly like the table + * {@code compression} option (it requires a {@code "class"}; extra keys are that compressor's options; it + * may also set {@code chunk_length_in_kb} / {@code min_compress_ratio}). Nothing is inherited from the + * table compression: unspecified settings use the same defaults as the {@code compression} option (e.g. + * the default chunk length). Levels not listed use the table compression. + */ + public static final String PER_LEVEL_COMPRESSION_OPTION = "per_level_compression"; + private static final boolean tolerateSstableSize = Boolean.getBoolean(Config.PROPERTY_PREFIX + "tolerate_sstable_size"); private static final String LEVEL_FANOUT_SIZE_OPTION = "fanout_size"; private static final String SINGLE_SSTABLE_UPLEVEL_OPTION = "single_sstable_uplevel"; @@ -621,6 +634,10 @@ public class LeveledCompactionStrategy extends AbstractCompactionStrategy uncheckedOptions.remove(LEVEL_FANOUT_SIZE_OPTION); uncheckedOptions.remove(SINGLE_SSTABLE_UPLEVEL_OPTION); + // Validate the per-level compression override (levels not listed use the table compression default) + CompactionParams.parsePerLevelCompression(options.get(PER_LEVEL_COMPRESSION_OPTION)); + uncheckedOptions.remove(PER_LEVEL_COMPRESSION_OPTION); + uncheckedOptions.remove(CompactionParams.Option.MIN_THRESHOLD.toString()); uncheckedOptions.remove(CompactionParams.Option.MAX_THRESHOLD.toString()); diff --git a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java index e8dff32fbc..7a0e956a4c 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java @@ -80,7 +80,7 @@ public class BigTableWriter extends SSTableWriter TimeUUID pendingRepair, boolean isTransient, TableMetadataRef metadata, - MetadataCollector metadataCollector, + MetadataCollector metadataCollector, SerializationHeader header, Collection observers, LifecycleNewTracker lifecycleNewTracker) @@ -122,6 +122,19 @@ public class BigTableWriter extends SSTableWriter private CompressionParams compressionFor(final OperationType opType) { CompressionParams compressionParams = metadata.getLocal().params.compression; + + // Per-level compression (LCS-only per-table option): use the compression configured for this + // sstable's level. Each level is a full, independent compression spec (same defaults as the table + // compression option; nothing inherited). Keyed on level. Flush is excluded (it has its own + // flush_compression policy below). No-op for non-LCS tables, which never carry the option. + if (compressionParams.isEnabled() && opType != OperationType.FLUSH) + { + CompressionParams override = metadata.getLocal().params.compaction + .perLevelCompression(metadataCollector.getSSTableLevel()); + if (override != null) + compressionParams = override; + } + final ICompressor compressor = compressionParams.getSstableCompressor(); if (null != compressor && opType == OperationType.FLUSH) diff --git a/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java b/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java index 1375331ce5..ee867b7e47 100644 --- a/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java +++ b/src/java/org/apache/cassandra/io/sstable/metadata/MetadataCollector.java @@ -250,6 +250,11 @@ public class MetadataCollector implements PartitionStatisticsCollector return this; } + public int getSSTableLevel() + { + return sstableLevel; + } + public MetadataCollector updateClusteringValues(ClusteringPrefix clustering) { minClustering = minClustering == null || comparator.compare(clustering, minClustering) < 0 ? clustering.minimize() : minClustering; diff --git a/src/java/org/apache/cassandra/schema/CompactionParams.java b/src/java/org/apache/cassandra/schema/CompactionParams.java index daf1f70c66..0522bef5ce 100644 --- a/src/java/org/apache/cassandra/schema/CompactionParams.java +++ b/src/java/org/apache/cassandra/schema/CompactionParams.java @@ -17,6 +17,7 @@ */ package org.apache.cassandra.schema; +import java.io.IOException; import java.lang.reflect.InvocationTargetException; import java.util.Arrays; import java.util.HashMap; @@ -30,6 +31,8 @@ import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.cassandra.db.compaction.AbstractCompactionStrategy; import org.apache.cassandra.db.compaction.LeveledCompactionStrategy; import org.apache.cassandra.db.compaction.SizeTieredCompactionStrategy; @@ -43,6 +46,8 @@ public final class CompactionParams { private static final Logger logger = LoggerFactory.getLogger(CompactionParams.class); + private static final ObjectMapper JSON_MAPPER = new ObjectMapper(); + public enum Option { CLASS, @@ -90,6 +95,7 @@ public final class CompactionParams private final ImmutableMap options; private final boolean isEnabled; private final TombstoneOption tombstoneOption; + private final ImmutableMap perLevelCompression; private CompactionParams(Class klass, Map options, boolean isEnabled, TombstoneOption tombstoneOption) { @@ -97,6 +103,19 @@ public final class CompactionParams this.options = ImmutableMap.copyOf(options); this.isEnabled = isEnabled; this.tombstoneOption = tombstoneOption; + + String perLevelSpec = options.get(LeveledCompactionStrategy.PER_LEVEL_COMPRESSION_OPTION); + ImmutableMap perLevel; + try + { + perLevel = ImmutableMap.copyOf(parsePerLevelCompression(perLevelSpec)); + } + catch (ConfigurationException e) + { + logger.warn("Ignoring malformed per_level_compression spec '{}': {}", perLevelSpec, e.getMessage()); + perLevel = ImmutableMap.of(); + } + this.perLevelCompression = perLevel; } public static CompactionParams create(Class klass, Map options) @@ -252,6 +271,68 @@ public final class CompactionParams return options; } + /** + * Returns the pre-built compression override for {@code level}, or {@code null} if the level has no + * override and should use the table compression. Each level is built independently from its own + * compression spec, exactly like the table {@code compression} option (same defaults, e.g. chunk length); + * nothing is inherited from the table. See {@link LeveledCompactionStrategy#PER_LEVEL_COMPRESSION_OPTION}. + */ + public CompressionParams perLevelCompression(int level) + { + return perLevelCompression.get(level); + } + + /** + * Parses the {@code per_level_compression} spec (see + * {@link LeveledCompactionStrategy#PER_LEVEL_COMPRESSION_OPTION} for the grammar) into a level to + * {@link CompressionParams} map. Each level value is a compression option map that must specify a + * {@code class} and is built and validated exactly like the table {@code compression} option via + * {@link CompressionParams#fromMap}. Returns an empty map for a null or blank spec. + */ + public static Map parsePerLevelCompression(String spec) throws ConfigurationException + { + Map result = new HashMap<>(); + if (spec == null || spec.trim().isEmpty()) + return result; + + Map> raw; + try + { + raw = JSON_MAPPER.readValue(spec, new TypeReference<>() {}); + } + catch (IOException e) + { + throw new ConfigurationException( + format("Invalid per_level_compression JSON '%s': %s", spec, e.getMessage()), e); + } + + for (Map.Entry> entry : raw.entrySet()) + { + int level; + try + { + level = Integer.parseInt(entry.getKey().trim()); + } + catch (NumberFormatException ex) + { + throw new ConfigurationException( + format("Invalid level '%s' in per_level_compression", entry.getKey()), ex); + } + + if (level < 0) + throw new ConfigurationException(format("per_level_compression level must be >= 0, but was %s", level)); + + Map levelOptions = entry.getValue(); + if (levelOptions == null || !levelOptions.containsKey(CompressionParams.CLASS)) + throw new ConfigurationException(format("per_level_compression entry for level %s must specify a '%s'", + level, CompressionParams.CLASS)); + + // A level is built exactly like the table 'compression' option (independently, same defaults). + result.put(level, CompressionParams.fromMap(levelOptions)); + } + return result; + } + public boolean isEnabled() { return isEnabled; diff --git a/test/unit/org/apache/cassandra/io/compress/PerLevelCompressionTest.java b/test/unit/org/apache/cassandra/io/compress/PerLevelCompressionTest.java new file mode 100644 index 0000000000..4ad7bd0b65 --- /dev/null +++ b/test/unit/org/apache/cassandra/io/compress/PerLevelCompressionTest.java @@ -0,0 +1,346 @@ +/* + * 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.io.compress; + +import java.util.Map; +import java.util.Set; + +import com.google.common.collect.ImmutableMap; +import org.junit.After; +import org.junit.Before; +import org.junit.BeforeClass; +import org.junit.Test; + +import org.apache.cassandra.config.Config; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.cql3.CQLTester; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.compaction.LeveledCompactionStrategy; +import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.schema.CompactionParams; +import org.apache.cassandra.schema.CompressionParams; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +/** + * Tests for the per-table {@code per_level_compression} compaction option on + * {@code LeveledCompactionStrategy}. See {@link LeveledCompactionStrategy#PER_LEVEL_COMPRESSION_OPTION}. + */ +public class PerLevelCompressionTest extends CQLTester +{ + private static Config.FlushCompression defaultFlush; + + @BeforeClass + public static void setUpClass() + { + CQLTester.setUpClass(); + defaultFlush = DatabaseDescriptor.getFlushCompression(); + } + + @Before + public void setUpCompressionState() + { + // Flush as the table codec so flushed sstables are distinguishable from compaction output. + DatabaseDescriptor.setFlushCompression(Config.FlushCompression.table); + } + + @After + public void tearDownCompressionState() + { + DatabaseDescriptor.setFlushCompression(defaultFlush); + } + + // ---- Pure parsing / building unit tests ---- + + @Test + public void parseAndResolveTest() + { + assertTrue(CompactionParams.parsePerLevelCompression(null).isEmpty()); + assertTrue(CompactionParams.parsePerLevelCompression(" ").isEmpty()); + + Map parsed = CompactionParams.parsePerLevelCompression( + "{\"0\":{\"class\":\"LZ4Compressor\"},\"2\":{\"class\":\"ZstdCompressor\",\"compression_level\":\"9\"}}"); + assertEquals(2, parsed.size()); + assertTrue(parsed.get(0).getSstableCompressor() instanceof LZ4Compressor); + assertTrue(parsed.get(2).getSstableCompressor() instanceof ZstdCompressor); + assertEquals("9", parsed.get(2).getOtherOptions().get("compression_level")); + + // Each level is built independently, like the 'compression' option: unspecified chunk length uses the + // compression default; an explicit chunk_length_in_kb is honored (nothing is inherited from the table). + assertEquals(CompressionParams.DEFAULT_CHUNK_LENGTH, parsed.get(0).chunkLength()); + assertEquals(1024 * 64, CompactionParams.parsePerLevelCompression( + "{\"1\":{\"class\":\"LZ4Compressor\",\"chunk_length_in_kb\":\"64\"}}").get(1).chunkLength()); + + // Resolution goes through the cached map on the immutable CompactionParams instance (no re-parsing). + assertTrue(lcsParams("{\"0\":{\"class\":\"LZ4Compressor\"}}").perLevelCompression(0) + .getSstableCompressor() instanceof LZ4Compressor); + assertNull(lcsParams("{\"0\":{\"class\":\"LZ4Compressor\"}}").perLevelCompression(1)); + + // A malformed spec must degrade gracefully to the table default (null) rather than throw. + assertNull(lcsParams("bogus").perLevelCompression(0)); + assertNull(lcsParams("{\"0\":{}}").perLevelCompression(0)); // missing class -> ignored + } + + private static CompactionParams lcsParams(String perLevelSpec) + { + return CompactionParams.create(LeveledCompactionStrategy.class, + ImmutableMap.of(LeveledCompactionStrategy.PER_LEVEL_COMPRESSION_OPTION, + perLevelSpec)); + } + + @Test + public void invalidSpecsRejectedTest() + { + assertBadSpec("not json"); + assertBadSpec("[]"); // not a JSON object + assertBadSpec("{\"0\":\"LZ4Compressor\"}"); // level value is not an object + assertBadSpec("{\"x\":{\"class\":\"LZ4Compressor\"}}"); // non-integer level + assertBadSpec("{\"-1\":{\"class\":\"LZ4Compressor\"}}"); // negative level + assertBadSpec("{\"0\":{}}"); // missing class + assertBadSpec("{\"0\":{\"class\":\"NotARealCompressor\"}}"); // class does not resolve + assertBadSpec("{\"0\":{\"class\":\"LZ4Compressor\",\"bogus_option\":\"1\"}}"); // unsupported option + } + + private static void assertBadSpec(String spec) + { + try + { + CompactionParams.parsePerLevelCompression(spec); + fail("Expected ConfigurationException for spec: " + spec); + } + catch (ConfigurationException expected) + { + // expected + } + } + + // ---- End-to-end tests through a real compaction ---- + + // NOTE: a maximal compaction (compact()) on LCS uses MajorLeveledCompactionWriter, which writes output + // starting at level 1. For a tiny dataset that yields a single L1 sstable, so the tests below target + // level 1 as the deterministic compaction-output level. + + @Test + public void overriddenLevelUsesFastCodecTest() throws Throwable + { + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'," + + "'per_level_compression':'{\"1\":{\"class\":\"LZ4Compressor\"}}'} " + + "AND compression = {'class':'ZstdCompressor'}"); + + ColumnFamilyStore store = flushTwiceAsZstd(); + compact(); + + assertHasLevel(store, 1); + assertLevelCompressionInvariant(store, 1); + } + + @Test + public void unlistedLevelFallsBackToTableCodecTest() throws Throwable + { + // Only level 0 is overridden; the compaction output (L1) must fall back to the table's Zstd. + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'," + + "'per_level_compression':'{\"0\":{\"class\":\"LZ4Compressor\"}}'} " + + "AND compression = {'class':'ZstdCompressor'}"); + + ColumnFamilyStore store = flushTwiceAsZstd(); + compact(); + + assertHasLevel(store, 1); + assertLevelCompressionInvariant(store, 0); + } + + @Test + public void perLevelCompressorOptionsAppliedTest() throws Throwable + { + // L1 override uses Zstd with a non-default compression_level; the table default Zstd carries no level. + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'," + + "'per_level_compression':'{\"1\":{\"class\":\"ZstdCompressor\",\"compression_level\":\"9\"}}'} " + + "AND compression = {'class':'ZstdCompressor'}"); + + ColumnFamilyStore store = flushTwiceAsZstd(); + compact(); + + assertHasLevel(store, 1); + store.getLiveSSTables().stream() + .filter(s -> s.getSSTableLevel() == 1) + .forEach(s -> { + CompressionParams params = s.getCompressionMetadata().parameters; + assertTrue("Level 1 should be Zstd", params.getSstableCompressor() instanceof ZstdCompressor); + assertEquals("per-level compression_level should be applied and persisted", + "9", params.getOtherOptions().get("compression_level")); + }); + } + + @Test + public void perLevelChunkLengthOverrideAppliedTest() throws Throwable + { + // Table uses the default 16 KB chunk length; L1 overrides chunk_length_in_kb to 64. + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'," + + "'per_level_compression':'{\"1\":{\"class\":\"LZ4Compressor\",\"chunk_length_in_kb\":\"64\"}}'} " + + "AND compression = {'class':'ZstdCompressor'}"); + + ColumnFamilyStore store = flushTwiceAsZstd(); + compact(); + + assertHasLevel(store, 1); + store.getLiveSSTables().stream() + .filter(s -> s.getSSTableLevel() == 1) + .forEach(s -> assertEquals("L1 chunk length should be the per-level override (64 KB)", + 1024 * 64, s.getCompressionMetadata().chunkLength())); + } + + @Test + public void levelChunkLengthIsIndependentOfTableTest() throws Throwable + { + // Table uses a non-default 64 KB chunk length; L1 overrides only the class, so it must fall back to + // the compression default chunk length (16 KB), NOT inherit the table's 64 KB. + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'," + + "'per_level_compression':'{\"1\":{\"class\":\"LZ4Compressor\"}}'} " + + "AND compression = {'class':'ZstdCompressor','chunk_length_in_kb':'64'}"); + + ColumnFamilyStore store = flushTwiceAsZstd(); + compact(); + + assertHasLevel(store, 1); + store.getLiveSSTables().stream() + .filter(s -> s.getSSTableLevel() == 1) + .forEach(s -> assertEquals("L1 chunk length should be the compression default, not the table's", + CompressionParams.DEFAULT_CHUNK_LENGTH, + s.getCompressionMetadata().chunkLength())); + } + + @Test + public void badOptionRejectedAtDdlTest() throws Throwable + { + try + { + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy','per_level_compression':'bogus'} " + + "AND compression = {'class':'ZstdCompressor'}"); + fail("Expected ConfigurationException for a malformed per_level_compression option"); + } + catch (RuntimeException e) + { + assertTrue("Unexpected exception: " + e, hasConfigurationException(e)); + } + } + + @Test + public void badOptionRejectedAtAlterTest() throws Throwable + { + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'} " + + "AND compression = {'class':'ZstdCompressor'}"); + + try + { + alterTableMayThrow("ALTER TABLE %s WITH compaction = " + + "{'class':'LeveledCompactionStrategy','per_level_compression':'bogus'}"); + fail("Expected ConfigurationException for a malformed per_level_compression option on ALTER"); + } + catch (Throwable t) + { + assertTrue("Unexpected exception: " + t, hasConfigurationException(t)); + } + } + + /** + * Returns true if {@code t} or any exception in its cause chain is a {@link ConfigurationException}. + * The validation failure is wrapped at different depths on CREATE vs ALTER, so walk the whole chain. + */ + private static boolean hasConfigurationException(Throwable t) + { + for (Throwable cause = t; cause != null; cause = cause.getCause()) + { + if (cause instanceof ConfigurationException) + return true; + } + return false; + } + + @Test + public void validOptionAcceptedAtAlterTest() throws Throwable + { + createTable("CREATE TABLE %s (k text PRIMARY KEY, v text) WITH " + + "compaction = {'class':'LeveledCompactionStrategy'} " + + "AND compression = {'class':'ZstdCompressor'}"); + + String spec = "{\"1\":{\"class\":\"LZ4Compressor\"}}"; + alterTable("ALTER TABLE %s WITH compaction = " + + "{'class':'LeveledCompactionStrategy','per_level_compression':'" + spec + "'}"); + + Map compaction = getCurrentColumnFamilyStore().metadata().params.compaction.options(); + assertEquals(spec, compaction.get(LeveledCompactionStrategy.PER_LEVEL_COMPRESSION_OPTION)); + } + + private ColumnFamilyStore flushTwiceAsZstd() throws Throwable + { + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + // Keep the two flushed L0 sstables around for the explicit compact() below; otherwise background + // LCS compaction can merge them first and make the sstable count non-deterministic. + cfs.disableAutoCompaction(); + + execute("INSERT INTO %s (k, v) values (?, ?)", "k1", "v1"); + flush(); + execute("INSERT INTO %s (k, v) values (?, ?)", "k2", "v2"); + flush(); + assertEquals(2, cfs.getLiveSSTables().size()); + + // Flushed sstables use the table codec (flush_compression=table) so they are Zstd. + cfs.getLiveSSTables().forEach(sstable -> + assertTrue(sstable.getCompressionMetadata().parameters.getSstableCompressor() instanceof ZstdCompressor)); + return cfs; + } + + private static void assertHasLevel(ColumnFamilyStore store, int level) + { + boolean found = store.getLiveSSTables().stream().anyMatch(s -> s.getSSTableLevel() == level); + assertTrue("Expected at least one sstable at level " + level, found); + } + + /** + * Asserts every live sstable's compressor matches the override: sstables at {@code fastLevel} use LZ4, + * every other level uses the table's Zstd. + */ + private static void assertLevelCompressionInvariant(ColumnFamilyStore store, int fastLevel) + { + Set sstables = store.getLiveSSTables(); + assertFalse("Expected at least one sstable after compaction", sstables.isEmpty()); + sstables.forEach(sstable -> { + Object compressor = sstable.getCompressionMetadata().parameters.getSstableCompressor(); + if (sstable.getSSTableLevel() == fastLevel) + assertTrue("Level " + fastLevel + " should be LZ4 but was " + compressor.getClass().getSimpleName(), + compressor instanceof LZ4Compressor); + else + assertTrue("Level " + sstable.getSSTableLevel() + " should be Zstd but was " + + compressor.getClass().getSimpleName(), + compressor instanceof ZstdCompressor); + }); + } +}