diff --git a/CHANGES.txt b/CHANGES.txt index 0c50b0ac0c..7f1930d121 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -22,6 +22,7 @@ * Minimize BTree iterator allocations (CASSANDRA-15389) Merged from 3.11: Merged from 3.0: + * liveDiskSpaceUsed and totalDiskSpaceUsed get corrupted if IndexSummaryRedistribution gets interrupted (CASSANDRA-15674) * Fix Debian init start/stop (CASSANDRA-15770) * Fix infinite loop on index query paging in tables with clustering (CASSANDRA-14242) * Fix chunk index overflow due to large sstable with small chunk length (CASSANDRA-15595) diff --git a/src/java/org/apache/cassandra/db/lifecycle/LifecycleTransaction.java b/src/java/org/apache/cassandra/db/lifecycle/LifecycleTransaction.java index a129e41e9a..574c6a4499 100644 --- a/src/java/org/apache/cassandra/db/lifecycle/LifecycleTransaction.java +++ b/src/java/org/apache/cassandra/db/lifecycle/LifecycleTransaction.java @@ -20,7 +20,6 @@ package org.apache.cassandra.db.lifecycle; import java.io.File; import java.nio.file.Path; import java.util.*; -import java.util.function.BiFunction; import java.util.function.BiPredicate; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Predicate; @@ -36,6 +35,7 @@ import org.apache.cassandra.db.compaction.OperationType; import org.apache.cassandra.io.sstable.SSTable; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.io.sstable.format.SSTableReader.UniqueIdentifier; +import org.apache.cassandra.utils.Throwables; import org.apache.cassandra.utils.concurrent.Transactional; import static com.google.common.base.Functions.compose; @@ -124,6 +124,10 @@ public class LifecycleTransaction extends Transactional.AbstractTransactional im // the tidier and their readers, to be used for marking readers obsoleted during a commit private List obsoletions; + // commit/rollback hooks + private List commitHooks = new ArrayList<>(); + private List abortHooks = new ArrayList<>(); + /** * construct a Transaction for use in an offline operation */ @@ -224,12 +228,14 @@ public class LifecycleTransaction extends Transactional.AbstractTransactional im accumulate = markObsolete(obsoletions, accumulate); accumulate = tracker.updateSizeTracking(logged.obsolete, logged.update, accumulate); + accumulate = runOnCommitHooks(accumulate); accumulate = release(selfRefs(logged.obsolete), accumulate); accumulate = tracker.notifySSTablesChanged(originals, logged.update, log.type(), accumulate); return accumulate; } + /** * undo all of the changes made by this transaction, resetting the state to its original form */ @@ -260,6 +266,7 @@ public class LifecycleTransaction extends Transactional.AbstractTransactional im accumulate = tracker.notifySSTablesChanged(invalid, restored, OperationType.COMPACTION, accumulate); // setReplaced immediately preceding versions that have not been obsoleted accumulate = setReplaced(logged.update, accumulate); + accumulate = runOnAbortooks(accumulate); // we have replaced all of logged.update and never made visible staged.update, // and the files we have logged as obsolete we clone fresh versions of, so they are no longer needed either // any _staged_ obsoletes should either be in staged.update already, and dealt with there, @@ -271,6 +278,32 @@ public class LifecycleTransaction extends Transactional.AbstractTransactional im return accumulate; } + private Throwable runOnCommitHooks(Throwable accumulate) + { + return runHooks(commitHooks, accumulate); + } + + private Throwable runOnAbortooks(Throwable accumulate) + { + return runHooks(abortHooks, accumulate); + } + + private static Throwable runHooks(Iterable hooks, Throwable accumulate) + { + for (Runnable hook : hooks) + { + try + { + hook.run(); + } + catch (Exception e) + { + accumulate = Throwables.merge(accumulate, e); + } + } + return accumulate; + } + @Override protected Throwable doPostCleanup(Throwable accumulate) { @@ -367,6 +400,16 @@ public class LifecycleTransaction extends Transactional.AbstractTransactional im staged.obsolete.add(reader); } + public void runOnCommit(Runnable fn) + { + commitHooks.add(fn); + } + + public void runOnAbort(Runnable fn) + { + abortHooks.add(fn); + } + /** * obsolete every file in the original transaction */ diff --git a/src/java/org/apache/cassandra/io/sstable/IndexSummaryRedistribution.java b/src/java/org/apache/cassandra/io/sstable/IndexSummaryRedistribution.java index 1300c99d4a..90a8621556 100644 --- a/src/java/org/apache/cassandra/io/sstable/IndexSummaryRedistribution.java +++ b/src/java/org/apache/cassandra/io/sstable/IndexSummaryRedistribution.java @@ -28,7 +28,6 @@ import java.util.Map; import java.util.UUID; import com.google.common.annotations.VisibleForTesting; -import com.google.common.collect.Iterables; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -36,11 +35,11 @@ import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.db.compaction.CompactionInfo; import org.apache.cassandra.db.compaction.CompactionInterruptedException; -import org.apache.cassandra.db.compaction.CompactionManager; import org.apache.cassandra.db.compaction.OperationType; import org.apache.cassandra.db.compaction.CompactionInfo.Unit; import org.apache.cassandra.db.lifecycle.LifecycleTransaction; import org.apache.cassandra.io.sstable.format.SSTableReader; +import org.apache.cassandra.metrics.StorageMetrics; import org.apache.cassandra.schema.TableId; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Pair; @@ -257,14 +256,39 @@ public class IndexSummaryRedistribution extends CompactionInfo.Holder sstable, sstable.getIndexSummarySamplingLevel(), Downsampling.BASE_SAMPLING_LEVEL, entry.newSamplingLevel, Downsampling.BASE_SAMPLING_LEVEL); ColumnFamilyStore cfs = Keyspace.open(sstable.metadata().keyspace).getColumnFamilyStore(sstable.metadata().id); + long oldSize = sstable.bytesOnDisk(); SSTableReader replacement = sstable.cloneWithNewSummarySamplingLevel(cfs, entry.newSamplingLevel); + long newSize = replacement.bytesOnDisk(); newSSTables.add(replacement); transactions.get(sstable.metadata().id).update(replacement, true); + addHooks(cfs, transactions, oldSize, newSize); } return newSSTables; } + /** + * Add hooks to correctly update the storage load metrics once the transaction is closed/aborted + */ + @SuppressWarnings("resource") // Transactions are closed in finally outside of this method + private void addHooks(ColumnFamilyStore cfs, Map transactions, long oldSize, long newSize) + { + LifecycleTransaction txn = transactions.get(cfs.metadata.id); + txn.runOnCommit(() -> { + // The new size will be added in Transactional.commit() as an updated SSTable, more details: CASSANDRA-13738 + StorageMetrics.load.dec(oldSize); + cfs.metric.liveDiskSpaceUsed.dec(oldSize); + cfs.metric.totalDiskSpaceUsed.dec(oldSize); + }); + txn.runOnAbort(() -> { + // the local disk was modified but book keeping couldn't be commited, apply the delta + long delta = oldSize - newSize; // if new is larger this will be negative, so dec will become a inc + StorageMetrics.load.dec(delta); + cfs.metric.liveDiskSpaceUsed.dec(delta); + cfs.metric.totalDiskSpaceUsed.dec(delta); + }); + } + @VisibleForTesting static Pair, List> distributeRemainingSpace(List toDownsample, long remainingSpace) { diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java index c7fea7eb74..e85abb7be6 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java @@ -1177,7 +1177,6 @@ public abstract class SSTableReader extends SSTable implements SelfRefCounted sstables = Lists.newArrayList(cfs.getSSTables(SSTableSet.CANONICAL)); + LifecycleTransaction txn = cfs.getTracker().tryModify(sstables, OperationType.UNKNOWN); + Map txns = ImmutableMap.of(cfs.metadata.id, txn); + // fail on the last file (* 3 because we call isStopRequested 3 times for each sstable, and we should fail on the last) + AtomicInteger countdown = new AtomicInteger(3 * sstables.size() - 1); + IndexSummaryRedistribution redistribution = new IndexSummaryRedistribution(txns, 0, 0) { + public boolean isStopRequested() + { + return countdown.decrementAndGet() == 0; + } + }; + try + { + IndexSummaryManager.redistributeSummaries(redistribution); + Assert.fail("Should throw CompactionInterruptedException"); + } + catch (CompactionInterruptedException e) + { + // trying to get this to happen + } + catch (IOException e) + { + throw new RuntimeException(e); + } + finally + { + try + { + FBUtilities.closeAll(txns.values()); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + } + } +}