diff --git a/CHANGES.txt b/CHANGES.txt index 772455c386..cff477b227 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -38,6 +38,7 @@ Merged from 2.1: * (cqlsh) Allow the SSL protocol version to be specified through the config file or environment variables (CASSANDRA-9544) Merged from 2.0: + * Add tool to find why expired sstables are not getting dropped (CASSANDRA-10015) * Remove erroneous pending HH tasks from tpstats/jmx (CASSANDRA-9129) * Don't cast expected bf size to an int (CASSANDRA-9959) * checkForEndpointCollision fails for legitimate collisions (CASSANDRA-9765) diff --git a/src/java/org/apache/cassandra/tools/SSTableExpiredBlockers.java b/src/java/org/apache/cassandra/tools/SSTableExpiredBlockers.java new file mode 100644 index 0000000000..2feee76585 --- /dev/null +++ b/src/java/org/apache/cassandra/tools/SSTableExpiredBlockers.java @@ -0,0 +1,136 @@ +/* + * 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.tools; + +import java.io.IOException; +import java.io.PrintStream; +import java.util.Collections; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; + +import com.google.common.base.Throwables; +import com.google.common.collect.ArrayListMultimap; +import com.google.common.collect.Multimap; + +import org.apache.cassandra.config.CFMetaData; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.config.Schema; +import org.apache.cassandra.db.ColumnFamilyStore; +import org.apache.cassandra.db.Directories; +import org.apache.cassandra.db.Keyspace; +import org.apache.cassandra.io.sstable.Component; +import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.io.sstable.format.SSTableReader; + +/** + * During compaction we can drop entire sstables if they only contain expired tombstones and if it is guaranteed + * to not cover anything in other sstables. An expired sstable can be blocked from getting dropped if its newest + * timestamp is newer than the oldest data in another sstable. + * + * This class outputs all sstables that are blocking other sstables from getting dropped so that a user can + * figure out why certain sstables are still on disk. + */ +public class SSTableExpiredBlockers +{ + public static void main(String[] args) throws IOException + { + PrintStream out = System.out; + if (args.length < 2) + { + out.println("Usage: sstableexpiredblockers "); + System.exit(1); + } + String keyspace = args[args.length - 2]; + String columnfamily = args[args.length - 1]; + Schema.instance.loadFromDisk(false); + + CFMetaData metadata = Schema.instance.getCFMetaData(keyspace, columnfamily); + if (metadata == null) + throw new IllegalArgumentException(String.format("Unknown keyspace/table %s.%s", + keyspace, + columnfamily)); + + Keyspace ks = Keyspace.openWithoutSSTables(keyspace); + ColumnFamilyStore cfs = ks.getColumnFamilyStore(columnfamily); + Directories.SSTableLister lister = cfs.directories.sstableLister().skipTemporary(true); + Set sstables = new HashSet<>(); + for (Map.Entry> sstable : lister.list().entrySet()) + { + if (sstable.getKey() != null) + { + try + { + SSTableReader reader = SSTableReader.open(sstable.getKey()); + sstables.add(reader); + } + catch (Throwable t) + { + out.println("Couldn't open sstable: " + sstable.getKey().filenameFor(Component.DATA)); + Throwables.propagate(t); + } + } + } + if (sstables.isEmpty()) + { + out.println("No sstables for " + keyspace + "." + columnfamily); + System.exit(1); + } + + int gcBefore = (int)(System.currentTimeMillis()/1000) - metadata.getGcGraceSeconds(); + Multimap blockers = checkForExpiredSSTableBlockers(sstables, gcBefore); + for (SSTableReader blocker : blockers.keySet()) + { + out.println(String.format("%s blocks %d expired sstables from getting dropped: %s%n", + formatForExpiryTracing(Collections.singleton(blocker)), + blockers.get(blocker).size(), + formatForExpiryTracing(blockers.get(blocker)))); + } + + System.exit(0); + } + + public static Multimap checkForExpiredSSTableBlockers(Iterable sstables, int gcBefore) + { + Multimap blockers = ArrayListMultimap.create(); + for (SSTableReader sstable : sstables) + { + if (sstable.getSSTableMetadata().maxLocalDeletionTime < gcBefore) + { + for (SSTableReader potentialBlocker : sstables) + { + if (!potentialBlocker.equals(sstable) && + potentialBlocker.getMinTimestamp() <= sstable.getMaxTimestamp() && + potentialBlocker.getSSTableMetadata().maxLocalDeletionTime > gcBefore) + blockers.put(potentialBlocker, sstable); + } + } + } + return blockers; + } + + private static String formatForExpiryTracing(Iterable sstables) + { + StringBuilder sb = new StringBuilder(); + + for (SSTableReader sstable : sstables) + sb.append(String.format("[%s (minTS = %d, maxTS = %d, maxLDT = %d)]", sstable, sstable.getMinTimestamp(), sstable.getMaxTimestamp(), sstable.getSSTableMetadata().maxLocalDeletionTime)).append(", "); + + return sb.toString(); + } +} diff --git a/test/unit/org/apache/cassandra/db/compaction/TTLExpiryTest.java b/test/unit/org/apache/cassandra/db/compaction/TTLExpiryTest.java index 579794dfd2..bd1e559799 100644 --- a/test/unit/org/apache/cassandra/db/compaction/TTLExpiryTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/TTLExpiryTest.java @@ -20,8 +20,8 @@ package org.apache.cassandra.db.compaction; * */ -import org.apache.cassandra.io.sstable.format.SSTableReader; import org.junit.BeforeClass; +import com.google.common.collect.Multimap; import com.google.common.collect.Sets; import org.junit.Test; import org.junit.runner.RunWith; @@ -33,8 +33,10 @@ import org.apache.cassandra.config.KSMetaData; import org.apache.cassandra.db.*; import org.apache.cassandra.db.columniterator.OnDiskAtomIterator; import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.io.sstable.ISSTableScanner; import org.apache.cassandra.locator.SimpleStrategy; +import org.apache.cassandra.tools.SSTableExpiredBlockers; import org.apache.cassandra.utils.ByteBufferUtil; import java.io.IOException; @@ -209,4 +211,30 @@ public class TTLExpiryTest scanner.close(); } + + @Test + public void testCheckForExpiredSSTableBlockers() throws InterruptedException + { + ColumnFamilyStore cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore("Standard1"); + cfs.truncateBlocking(); + cfs.disableAutoCompaction(); + cfs.metadata.gcGraceSeconds(0); + + Mutation rm = new Mutation(KEYSPACE1, Util.dk("test").getKey()); + rm.add("Standard1", Util.cellname("col1"), ByteBufferUtil.EMPTY_BYTE_BUFFER, System.currentTimeMillis()); + rm.applyUnsafe(); + cfs.forceBlockingFlush(); + SSTableReader blockingSSTable = cfs.getSSTables().iterator().next(); + for (int i = 0; i < 10; i++) + { + rm = new Mutation(KEYSPACE1, Util.dk("test").getKey()); + rm.delete("Standard1", System.currentTimeMillis()); + rm.applyUnsafe(); + cfs.forceBlockingFlush(); + } + Multimap blockers = SSTableExpiredBlockers.checkForExpiredSSTableBlockers(cfs.getSSTables(), (int) (System.currentTimeMillis() / 1000) + 100); + assertEquals(1, blockers.keySet().size()); + assertTrue(blockers.keySet().contains(blockingSSTable)); + assertEquals(10, blockers.get(blockingSSTable).size()); + } } diff --git a/tools/bin/sstableexpiredblockers b/tools/bin/sstableexpiredblockers new file mode 100755 index 0000000000..00272081da --- /dev/null +++ b/tools/bin/sstableexpiredblockers @@ -0,0 +1,54 @@ +#!/bin/sh + +# 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. + +if [ "x$CASSANDRA_INCLUDE" = "x" ]; then + for include in /usr/share/cassandra/cassandra.in.sh \ + /usr/local/share/cassandra/cassandra.in.sh \ + /opt/cassandra/cassandra.in.sh \ + ~/.cassandra.in.sh \ + "`dirname "$0"`/cassandra.in.sh"; do + if [ -r "$include" ]; then + . "$include" + break + fi + done +elif [ -r "$CASSANDRA_INCLUDE" ]; then + . "$CASSANDRA_INCLUDE" +fi + +# Use JAVA_HOME if set, otherwise look for java in PATH +if [ -x "$JAVA_HOME/bin/java" ]; then + JAVA="$JAVA_HOME/bin/java" +else + JAVA="`which java`" +fi + +if [ -z "$CLASSPATH" ]; then + echo "You must set the CLASSPATH var" >&2 + exit 1 +fi + +if [ "x$MAX_HEAP_SIZE" = "x" ]; then + MAX_HEAP_SIZE="256M" +fi + +"$JAVA" $JAVA_AGENT -ea -cp "$CLASSPATH" -Xmx$MAX_HEAP_SIZE \ + -Dcassandra.storagedir="$cassandra_storagedir" \ + -Dlogback.configurationFile=logback-tools.xml \ + org.apache.cassandra.tools.SSTableExpiredBlockers "$@" + diff --git a/tools/bin/sstableexpiredblockers.bat b/tools/bin/sstableexpiredblockers.bat new file mode 100644 index 0000000000..7af110575d --- /dev/null +++ b/tools/bin/sstableexpiredblockers.bat @@ -0,0 +1,23 @@ +@REM Licensed to the Apache Software Foundation (ASF) under one or more +@REM contributor license agreements. See the NOTICE file distributed with +@REM this work for additional information regarding copyright ownership. +@REM The ASF licenses this file to You under the Apache License, Version 2.0 +@REM (the "License"); you may not use this file except in compliance with +@REM the License. You may obtain a copy of the License at +@REM +@REM http://www.apache.org/licenses/LICENSE-2.0 +@REM +@REM Unless required by applicable law or agreed to in writing, software +@REM distributed under the License is distributed on an "AS IS" BASIS, +@REM WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +@REM See the License for the specific language governing permissions and +@REM limitations under the License. + +@echo off + +if "%OS%" == "Windows_NT" setlocal + +pushd "%~dp0" +call cassandra.in.bat + +"%JAVA_HOME%\bin\java" -cp %CLASSPATH% org.apache.cassandra.tools.SSTableExpiredBlockers %*