Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,24 @@

package org.apache.cassandra.io.sstable.format;

import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.apache.cassandra.io.sstable.Component;
import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.util.File;

public abstract class AbstractSSTableFormat<R extends SSTableReader, W extends SSTableWriter> implements SSTableFormat<R, W>
{
private final Logger logger = LoggerFactory.getLogger(getClass());

public final String name;
protected final Map<String, String> options;

Expand All @@ -38,6 +51,41 @@ public final String name()
return name;
}

@Override
public final void deleteOrphanedComponents(Descriptor descriptor, Set<Component> components)
{
File dataFile = descriptor.fileFor(Components.DATA);
if (components.contains(Components.DATA) && dataFile.length() > 0)
// everything appears to be in order... moving on.
return;

// missing the DATA file! all components are orphaned
logger.warn("[{}] Removing orphans for {}: {}", getClass().getSimpleName(), descriptor, components);
for (Component component : components)
{
File file = descriptor.fileFor(component);
if (file.exists())
descriptor.fileFor(component).delete();
}
}

protected final void deleteComponentsOldestFirst(Descriptor desc, List<Component> components)
{
logger.info("[{}] Deleting sstable: {}", getClass().getSimpleName(), desc);

// delete older files first so the overall SSTable timestamp stays the same on partial deletes
Map<Component, Long> lastModified = new HashMap<>();
for (Component c : components)
lastModified.put(c, desc.fileFor(c).lastModified());
components.sort(Comparator.comparingLong(lastModified::get));

for (Component component : components)
{
logger.trace("[{}] Deleting component {} of {}", getClass().getSimpleName(), component, desc);
desc.fileFor(component).deleteIfExists();
}
}

@Override
public final boolean equals(Object o)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,6 @@
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.Set;
import java.util.SortedSet;
import java.util.TreeSet;
import java.util.concurrent.locks.Lock;
Expand Down Expand Up @@ -61,12 +60,9 @@
import org.apache.cassandra.db.rows.UnfilteredRowIterators;
import org.apache.cassandra.db.rows.WrappingUnfilteredRowIterator;
import org.apache.cassandra.exceptions.ConfigurationException;
import org.apache.cassandra.io.sstable.Component;
import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.IScrubber;
import org.apache.cassandra.io.sstable.SSTableIdentityIterator;
import org.apache.cassandra.io.sstable.SSTableRewriter;
import org.apache.cassandra.io.sstable.format.SSTableFormat.Components;
import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
import org.apache.cassandra.io.util.File;
import org.apache.cassandra.io.util.RandomAccessReader;
Expand Down Expand Up @@ -164,23 +160,6 @@ protected SortedTableScrubber(ColumnFamilyStore cfs,
outputHandler.output("Starting scrub with reinsert overflowed TTL option");
}

public static void deleteOrphanedComponents(Descriptor descriptor, Set<Component> components)
{
File dataFile = descriptor.fileFor(Components.DATA);
if (components.contains(Components.DATA) && dataFile.length() > 0)
// everything appears to be in order... moving on.
return;

// missing the DATA file! all components are orphaned
logger.warn("Removing orphans for {}: {}", descriptor, components);
for (Component component : components)
{
File file = descriptor.fileFor(component);
if (file.exists())
descriptor.fileFor(component).delete();
}
}

@Override
public void scrub()
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@

import java.io.IOException;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
Expand All @@ -30,9 +29,6 @@
import com.google.common.collect.Lists;
import com.google.common.collect.Sets;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.apache.cassandra.cache.KeyCacheKey;
import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
Expand All @@ -51,7 +47,6 @@
import org.apache.cassandra.io.sstable.format.SSTableFormat;
import org.apache.cassandra.io.sstable.format.SSTableReaderLoadingBuilder;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.format.SortedTableScrubber;
import org.apache.cassandra.io.sstable.format.Version;
import org.apache.cassandra.io.sstable.indexsummary.IndexSummaryMetrics;
import org.apache.cassandra.io.sstable.keycache.KeyCacheMetrics;
Expand All @@ -64,8 +59,6 @@
import org.apache.cassandra.utils.OutputHandler;
import org.apache.cassandra.utils.Pair;

import static org.apache.cassandra.io.sstable.format.SSTableFormat.Components.DATA;

/**
* Legacy bigtable format. Components and approximate lifecycle:
* <br>
Expand Down Expand Up @@ -163,8 +156,6 @@
*/
public class BigFormat extends AbstractSSTableFormat<BigTableReader, BigTableWriter>
{
private final static Logger logger = LoggerFactory.getLogger(BigFormat.class);

public static final String NAME = "big";

private final Version latestVersion = new BigVersion(this, BigVersion.current_version);
Expand Down Expand Up @@ -313,28 +304,6 @@ public MetricsProviders getFormatSpecificMetricsProviders()
return BigTableSpecificMetricsProviders.instance;
}

@Override
public void deleteOrphanedComponents(Descriptor descriptor, Set<Component> components)
{
SortedTableScrubber.deleteOrphanedComponents(descriptor, components);
}

private void delete(Descriptor desc, List<Component> components)
{
logger.info("Deleting sstable: {}", desc);

if (components.remove(DATA))
components.add(0, DATA); // DATA component should be first
if (components.remove(Components.SUMMARY))
components.add(Components.SUMMARY); // SUMMARY component should be last (IDK why)

for (Component component : components)
{
logger.trace("Deleting component {} of {}", component, desc);
desc.fileFor(component).deleteIfExists();
}
}

@Override
public void delete(Descriptor desc)
{
Expand All @@ -352,7 +321,7 @@ public void delete(Descriptor desc)
}
}

delete(desc, Lists.newArrayList(Sets.intersection(allComponents(), desc.discoverComponents())));
deleteComponentsOldestFirst(desc, Lists.newArrayList(Sets.intersection(allComponents(), desc.discoverComponents())));
}
catch (Throwable t)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
package org.apache.cassandra.io.sstable.format.bti;

import java.io.IOException;
import java.util.List;
import java.util.Map;
import java.util.Set;

Expand All @@ -27,9 +26,6 @@
import com.google.common.collect.Lists;
import com.google.common.collect.Sets;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.ColumnFamilyStore;
import org.apache.cassandra.db.DecoratedKey;
Expand All @@ -47,7 +43,6 @@
import org.apache.cassandra.io.sstable.format.SSTableReader;
import org.apache.cassandra.io.sstable.format.SSTableReaderLoadingBuilder;
import org.apache.cassandra.io.sstable.format.SSTableWriter;
import org.apache.cassandra.io.sstable.format.SortedTableScrubber;
import org.apache.cassandra.io.sstable.format.Version;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.schema.TableMetadataRef;
Expand All @@ -60,8 +55,6 @@
*/
public class BtiFormat extends AbstractSSTableFormat<BtiTableReader, BtiTableWriter>
{
private final static Logger logger = LoggerFactory.getLogger(BtiFormat.class);

public static final String NAME = "bti";

private final Version latestVersion = new BtiVersion(this, BtiVersion.current_version);
Expand Down Expand Up @@ -206,32 +199,12 @@ public MetricsProviders getFormatSpecificMetricsProviders()
return BtiTableSpecificMetricsProviders.instance;
}

@Override
public void deleteOrphanedComponents(Descriptor descriptor, Set<Component> components)
{
SortedTableScrubber.deleteOrphanedComponents(descriptor, components);
}

private void delete(Descriptor desc, List<Component> components)
{
logger.info("Deleting sstable: {}", desc);

if (components.remove(SSTableFormat.Components.DATA))
components.add(0, SSTableFormat.Components.DATA); // DATA component should be first

for (Component component : components)
{
logger.trace("Deleting component {} of {}", component, desc);
desc.fileFor(component).deleteIfExists();
}
}

@Override
public void delete(Descriptor desc)
{
try
{
delete(desc, Lists.newArrayList(Sets.intersection(allComponents(), desc.discoverComponents())));
deleteComponentsOldestFirst(desc, Lists.newArrayList(Sets.intersection(allComponents(), desc.discoverComponents())));
}
catch (Throwable t)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,11 @@
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.NoSuchFileException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
Expand Down Expand Up @@ -1634,4 +1636,88 @@ public void useAfterCompletedTest()
txnFile.abort(); // this should complete the txn
txnFile.trackNew(dummySSTable()); // expect an IllegalStateException here
}
}}
}

@Test
public void testBigFormatSSTableLastModifiedTimestampNeverChangesOnPartialDeletes() throws IOException
{
assertLastModifiedTimestampNeverChangesOnPartialDeletes(BigFormat.getInstance());
}

@Test
public void testBtiFormatSSTableLastModifiedTimestampNeverChangesOnPartialDeletes() throws IOException
{
assertLastModifiedTimestampNeverChangesOnPartialDeletes(DatabaseDescriptor.getSSTableFormats().get(BtiFormat.NAME));
}

// Verifies that a format's deletion of an sstable's component files never lets the highest
// update time among the sstable's files change, even mid-deletion.
private static void assertLastModifiedTimestampNeverChangesOnPartialDeletes(SSTableFormat<?, ?> format) throws IOException
{
// 1) create an sstable's real component files
File dataDir = new File(Files.createTempDirectory("LastModifiedOnPartialDeletesTest").toFile());
Descriptor realDesc = new Descriptor(dataDir, "ks", "cf", new SequenceBasedSSTableId(1), format);

Map<Component, File> realFiles = new HashMap<>();
for (Component component : format.allComponents())
{
File file = realDesc.fileFor(component);
file.parent().createDirectoriesIfNotExists();
assertTrue(file.createFileIfNotExists());
realFiles.put(component, file);
}

// 2) change the update times so that DATA ends up with the highest one
long base = System.currentTimeMillis();
long offset = 0;
for (Component component : realFiles.keySet())
{
if (!component.equals(SSTableFormat.Components.DATA))
assertTrue(realFiles.get(component).trySetLastModified(base + (offset++) * 1000));
}
assertTrue(realFiles.get(SSTableFormat.Components.DATA).trySetLastModified(base + realFiles.size() * 1000));

// 3) get the highest timestamp and save it
long expectedMaxUpdateTime = realFiles.values().stream().mapToLong(File::lastModified).max().getAsLong();

Set<Component> deleted = new HashSet<>();
List<String> violations = new ArrayList<>();
Descriptor stubDesc = stubDeletionOrder(realDesc, realFiles, component ->
{
deleted.add(component);
long remainingMax = realFiles.entrySet().stream()
.filter(e -> !deleted.contains(e.getKey()))
.mapToLong(e -> e.getValue().lastModified())
.max().orElse(-1);
if (remainingMax >= 0 && remainingMax != expectedMaxUpdateTime)
violations.add("after deleting " + component + ": last-modified calculation changed from "
+ expectedMaxUpdateTime + " to " + remainingMax);
});

// 4) run the real delete path and confirm each partial delete preserved the invariant
format.delete(stubDesc);

assertTrue("expected no violations, got: " + violations, violations.isEmpty());
assertEquals(realFiles.size(), deleted.size());
}

private static Descriptor stubDeletionOrder(Descriptor realDesc, Map<Component, File> realFiles, Consumer<Component> onDelete)
{
return new Descriptor(realDesc.version, realDesc.directory, realDesc.ksname,
realDesc.cfname, realDesc.id)
{
@Override
public File fileFor(Component component)
{
return new File(realFiles.get(component).toPath())
{
@Override
public void deleteIfExists()
{
onDelete.accept(component);
}
};
}
};
}
}