This is an automated email from the ASF dual-hosted git repository. asf-gitbox-commits pushed a commit to branch cassandra-4.1 in repository https://gitbox.apache.org/repos/asf/cassandra.git
commit be37969e89a5741cc41687a24eeee77d461917b4 Merge: 4827aafd02 7b4170491d Author: Caleb Rackliffe <[email protected]> AuthorDate: Mon Aug 17 14:15:19 2026 -0500 Merge branch 'cassandra-4.0' into cassandra-4.1 * cassandra-4.0: Make runWithCompactionsDisabled return non-null on success CHANGES.txt | 1 + .../org/apache/cassandra/db/ColumnFamilyStore.java | 89 ++++++++++----- .../cassandra/db/compaction/CompactionManager.java | 5 +- .../org/apache/cassandra/db/view/TableViews.java | 3 + .../cassandra/exceptions/RequestFailureReason.java | 8 +- src/java/org/apache/cassandra/net/InboundSink.java | 5 +- .../apache/cassandra/db/TruncateBlockingTest.java | 99 +++++++++++++++++ .../db/compaction/CancelCompactionsTest.java | 121 +++++++++++++++++++++ 8 files changed, 300 insertions(+), 31 deletions(-) diff --cc CHANGES.txt index b76d2a3754,50ad389065..75813be2a5 --- a/CHANGES.txt +++ b/CHANGES.txt @@@ -1,5 -1,5 +1,6 @@@ -4.0.22 +4.1.13 +Merged from 4.0: + * Make runWithCompactionsDisabled return non-null on success (CASSANDRA-21527) * Validate authz before performing role check in LIST ROLES/PERMISSIONS (CASSANDRA-21560) * Ensure transferred_ranges reset on decommision re-attempt when pending ranges cannot be proven continous (CASSANDRA-16290) * Add validation to uncompressed length during decompression (CASSANDRA-21567) diff --cc src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 2d23522280,2ac4783d89..6192781ed7 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@@ -73,46 -41,25 +73,47 @@@ import com.google.common.util.concurren import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.cache.*; -import org.apache.cassandra.concurrent.*; -import org.apache.cassandra.config.*; +import org.apache.cassandra.cache.CounterCacheKey; +import org.apache.cassandra.cache.IRowCacheEntry; +import org.apache.cassandra.cache.RowCacheKey; +import org.apache.cassandra.cache.RowCacheSentinel; +import org.apache.cassandra.concurrent.ExecutorPlus; +import org.apache.cassandra.concurrent.FutureTask; +import org.apache.cassandra.config.CassandraRelevantProperties; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.config.DurationSpec; import org.apache.cassandra.db.commitlog.CommitLog; import org.apache.cassandra.db.commitlog.CommitLogPosition; -import org.apache.cassandra.db.compaction.*; +import org.apache.cassandra.db.compaction.AbstractCompactionStrategy; +import org.apache.cassandra.db.compaction.CompactionManager; +import org.apache.cassandra.db.compaction.CompactionStrategyManager; +import org.apache.cassandra.db.compaction.OperationType; +import org.apache.cassandra.db.compaction.Verifier; import org.apache.cassandra.db.filter.ClusteringIndexFilter; import org.apache.cassandra.db.filter.DataLimits; -import org.apache.cassandra.db.streaming.CassandraStreamManager; -import org.apache.cassandra.db.repair.CassandraTableRepairManager; -import org.apache.cassandra.db.view.TableViews; -import org.apache.cassandra.db.lifecycle.*; +import org.apache.cassandra.db.memtable.Flushing; +import org.apache.cassandra.db.memtable.Memtable; +import org.apache.cassandra.db.memtable.ShardBoundaries; +import org.apache.cassandra.db.lifecycle.LifecycleNewTracker; +import org.apache.cassandra.db.lifecycle.LifecycleTransaction; +import org.apache.cassandra.db.lifecycle.SSTableSet; +import org.apache.cassandra.db.lifecycle.Tracker; +import org.apache.cassandra.db.lifecycle.View; import org.apache.cassandra.db.partitions.CachedPartition; import org.apache.cassandra.db.partitions.PartitionUpdate; +import org.apache.cassandra.db.repair.CassandraTableRepairManager; import org.apache.cassandra.db.rows.CellPath; -import org.apache.cassandra.dht.*; +import org.apache.cassandra.db.streaming.CassandraStreamManager; +import org.apache.cassandra.db.view.TableViews; +import org.apache.cassandra.dht.AbstractBounds; +import org.apache.cassandra.dht.Bounds; +import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Splitter; +import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; import org.apache.cassandra.exceptions.StartupException; ++import org.apache.cassandra.exceptions.TruncateException; import org.apache.cassandra.index.SecondaryIndexManager; import org.apache.cassandra.index.internal.CassandraIndex; import org.apache.cassandra.index.transactions.UpdateTransaction; @@@ -1801,11 -1670,18 +1810,18 @@@ public class ColumnFamilyStore implemen if (force) { Predicate<SSTableReader> predicate = sst -> { - UUID session = sst.getPendingRepair(); + TimeUUID session = sst.getPendingRepair(); return session != null && sessions.contains(session); }; - return runWithCompactionsDisabled(() -> compactionStrategyManager.releaseRepairData(sessions), - predicate, false, true, true); + CleanupSummary summary = runWithCompactionsDisabled(() -> compactionStrategyManager.releaseRepairData(sessions), + predicate, false, true, true); + if (summary == null) + { + logger.warn("Unable to cancel in-progress compactions for {}.{}, could not force release repair data for sessions {}", + keyspace.getName(), name, sessions); + return new CleanupSummary(this, Collections.emptySet(), new HashSet<>(sessions)); + } + return summary; } else { @@@ -2665,39 -2344,58 +2681,58 @@@ for (SSTableReader sstable : cfs.getLiveSSTables()) now = Math.max(now, sstable.maxDataAge); truncatedAt = now; - - Runnable truncateRunnable = new Runnable() + Throwable failure = null; + try { - public void run() - { - logger.info("Truncating {}.{} with truncatedAt={}", keyspace.getName(), getTableName(), truncatedAt); - // since truncation can happen at different times on different nodes, we need to make sure - // that any repairs are aborted, otherwise we might clear the data on one node and then - // stream in data that is actually supposed to have been deleted - ActiveRepairService.instance.abort((prs) -> prs.getTableIds().contains(metadata.id), - "Stopping parent sessions {} due to truncation of tableId="+metadata.id); - data.notifyTruncated(truncatedAt); + Boolean succeeded = runWithCompactionsDisabled(() -> runTruncate(truncatedAt, noSnapshot, replayAfter), + true, true); + // null means compactions couldn't be disabled, not failure of the truncate work itself + if (succeeded == null) + failure = new TruncateException("Unable to stop in-progress compactions. Please retry when current compaction tasks have completed or overall system load has decreased."); + else if (!succeeded) + failure = new TruncateException("Truncate failed."); + } + catch (Throwable t) + { + failure = t; + } - if (!noSnapshot && DatabaseDescriptor.isAutoSnapshot()) - snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, SNAPSHOT_TRUNCATE_PREFIX), DatabaseDescriptor.getAutoSnapshotTtl()); + try + { + viewManager.build(); + } + catch (Throwable t) + { + failure = merge(failure, t); + } - discardSSTables(truncatedAt); + maybeFail(failure); - indexManager.truncateAllIndexesBlocking(truncatedAt); - viewManager.truncateBlocking(replayAfter, truncatedAt); + logger.info("Truncate of {}.{} is complete", keyspace.getName(), name); + } - SystemKeyspace.saveTruncationRecord(ColumnFamilyStore.this, truncatedAt, replayAfter); - logger.trace("cleaning out row cache"); - invalidateCaches(); + private boolean runTruncate(long truncatedAt, boolean noSnapshot, CommitLogPosition replayAfter) + { + logger.info("Truncating {}.{} with truncatedAt={}", keyspace.getName(), getTableName(), truncatedAt); + // since truncation can happen at different times on different nodes, we need to make sure + // that any repairs are aborted, otherwise we might clear the data on one node and then + // stream in data that is actually supposed to have been deleted + ActiveRepairService.instance.abort((prs) -> prs.getTableIds().contains(metadata.id), + "Stopping parent sessions {} due to truncation of tableId="+metadata.id); + data.notifyTruncated(truncatedAt); - } - }; + if (!noSnapshot && DatabaseDescriptor.isAutoSnapshot()) - snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, SNAPSHOT_TRUNCATE_PREFIX)); ++ snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, SNAPSHOT_TRUNCATE_PREFIX), DatabaseDescriptor.getAutoSnapshotTtl()); - runWithCompactionsDisabled(FutureTask.callable(truncateRunnable), true, true); + discardSSTables(truncatedAt); - viewManager.build(); + indexManager.truncateAllIndexesBlocking(truncatedAt); + viewManager.truncateBlocking(replayAfter, truncatedAt); - logger.info("Truncate of {}.{} is complete", keyspace.getName(), name); + SystemKeyspace.saveTruncationRecord(ColumnFamilyStore.this, truncatedAt, replayAfter); + logger.trace("cleaning out row cache"); + invalidateCaches(); + return true; } /** diff --cc src/java/org/apache/cassandra/exceptions/RequestFailureReason.java index 3d3476a139,6664070e97..3b911b68eb --- a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java +++ b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java @@@ -36,8 -36,7 +36,9 @@@ public enum RequestFailureReaso READ_TOO_MANY_TOMBSTONES (1), TIMEOUT (2), INCOMPATIBLE_SCHEMA (3), + READ_SIZE (4), - NODE_DOWN (5); ++ NODE_DOWN (5), + TRUNCATE_FAILED (12); public static final Serializer serializer = new Serializer(); diff --cc test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java index 67421ba6d2,7ea27a1367..12af15c05e --- a/test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java @@@ -32,8 -33,14 +32,14 @@@ import java.util.stream.Collectors import com.google.common.collect.ImmutableSet; import com.google.common.util.concurrent.Uninterruptibles; + import net.bytebuddy.ByteBuddy; + import net.bytebuddy.agent.ByteBuddyAgent; + import net.bytebuddy.dynamic.loading.ClassReloadingStrategy; + import net.bytebuddy.implementation.StubMethod; ++ + import org.assertj.core.api.Assertions; import org.junit.Test; -import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.cql3.CQLTester; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.lifecycle.LifecycleTransaction; @@@ -51,14 -58,12 +57,17 @@@ import org.apache.cassandra.io.sstable. import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.RangesAtEndpoint; import org.apache.cassandra.locator.Replica; ++import org.apache.cassandra.repair.consistent.admin.CleanupSummary; ++import org.apache.cassandra.schema.CompactionParams.TombstoneOption; import org.apache.cassandra.schema.MockSchema; import org.apache.cassandra.schema.TableMetadata; import org.apache.cassandra.service.ActiveRepairService; import org.apache.cassandra.streaming.PreviewKind; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.TimeUUID; ++import static net.bytebuddy.matcher.ElementMatchers.named; +import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; @@@ -363,6 -372,118 +372,118 @@@ public class CancelCompactionsTest exte return new Murmur3Partitioner.LongToken(t); } + private LifecycleTransaction blockCompactions(ColumnFamilyStore cfs, List<SSTableReader> sstables) + { + LifecycleTransaction txn = cfs.getTracker().tryModify(sstables, OperationType.COMPACTION); + assertNotNull(txn); + return txn; + } + + private ClassReloadingStrategy stubWaitForCessation() + { + ByteBuddyAgent.install(); + ClassReloadingStrategy strategy = ClassReloadingStrategy.fromInstalledAgent(); + new ByteBuddy().redefine(CompactionManager.class) + .method(named("waitForCessation")) + .intercept(StubMethod.INSTANCE) + .make() + .load(CompactionManager.class.getClassLoader(), strategy); + return strategy; + } + + @Test + public void testForceCompactionThrowsWhenCompactionsCannotBeDisabled() throws Exception + { + ColumnFamilyStore cfs = MockSchema.newCFS(); + List<SSTableReader> sstables = createSSTables(cfs, 3, 0); + + LifecycleTransaction txn = blockCompactions(cfs, sstables); + ClassReloadingStrategy strategy = stubWaitForCessation(); + try + { + Range<Token> allData = new Range<>(cfs.getPartitioner().getMinimumToken(), cfs.getPartitioner().getMaximumToken()); + + // forceCompaction should fail at the null runWithCompactionsDisabled result + Assertions.assertThatThrownBy(() -> cfs.forceCompactionForTokenRange(Collections.singleton(allData))) + .as("Unable to cancel in-progress compactions. Usually retrying will work") + .isInstanceOf(RuntimeException.class); + } + finally + { + strategy.reset(CompactionManager.class); + txn.abort(); + } + } + + @Test + public void testGarbageCollectReturnsUnableToCancelWhenCompactionsCannotBeDisabled() throws Throwable + { + ColumnFamilyStore cfs = MockSchema.newCFS(); + List<SSTableReader> sstables = createSSTables(cfs, 3, 0); + + LifecycleTransaction txn = blockCompactions(cfs, sstables); + ClassReloadingStrategy strategy = stubWaitForCessation(); + try + { + // garbageCollect goes through parallelAllSSTableOperation, which must handle a null + // LifecycleTransaction from markAllCompacting without NPEing. + assertEquals(CompactionManager.AllSSTableOpStatus.UNABLE_TO_CANCEL, cfs.garbageCollect(TombstoneOption.ROW, 0)); + } + finally + { + strategy.reset(CompactionManager.class); + txn.abort(); + } + } + + @Test + public void testReleaseRepairDataReturnsUnsuccessfulWhenCompactionsCannotBeDisabled() throws Exception + { + ColumnFamilyStore cfs = MockSchema.newCFS(); + List<SSTableReader> sstables = createSSTables(cfs, 3, 0); + + // releaseRepairData only tries to cancel compactions on sstables pending one of these sessions, + // so tag one of our blocked sstables with a pending session to make it match - UUID pendingSession = UUID.randomUUID(); - Set<UUID> sessions = ImmutableSet.of(pendingSession, UUID.randomUUID()); ++ TimeUUID pendingSession = nextTimeUUID(); ++ Set<TimeUUID> sessions = ImmutableSet.of(pendingSession, nextTimeUUID()); + AbstractPendingRepairTest.mutateRepaired(sstables.get(0), pendingSession, false); + + LifecycleTransaction txn = blockCompactions(cfs, sstables); + ClassReloadingStrategy strategy = stubWaitForCessation(); + try + { + // force release should report the sessions as unsuccessful rather than throwing + CleanupSummary summary = cfs.releaseRepairData(sessions, true); + assertTrue(summary.successful.isEmpty()); + assertEquals(sessions, summary.unsuccessful); + } + finally + { + strategy.reset(CompactionManager.class); + txn.abort(); + } + } + + @Test + public void testSubmitMaximalNoOpsWhenCompactionsCannotBeDisabled() throws Exception + { + ColumnFamilyStore cfs = MockSchema.newCFS(); + List<SSTableReader> sstables = createSSTables(cfs, 3, 0); + + LifecycleTransaction txn = blockCompactions(cfs, sstables); + ClassReloadingStrategy strategy = stubWaitForCessation(); + try + { + // submitMaximal should no-op on null getMaximalTasks result + assertTrue(CompactionManager.instance.submitMaximal(cfs, -1, false).isEmpty()); + } + finally + { + strategy.reset(CompactionManager.class); + txn.abort(); + } + } + private List<SSTableReader> createSSTables(ColumnFamilyStore cfs, int count, int startGeneration) { List<SSTableReader> sstables = new ArrayList<>(); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
