Switch to CRC32 for sstable metadata checksums

patch by Aleksey Yeschenko; reviewed by Marcus Eriksson for
CASSANDRA-13953
This commit is contained in:
Aleksey Yeschenko 2017-10-13 14:26:02 +01:00
parent 7f989edd26
commit 912fdb3ea4
43 changed files with 170 additions and 118 deletions

View File

@ -1,7 +1,7 @@
4.0 4.0
* Checksum sstable metadata (CASSANDRA-13321, CASSANDRA-13593)
* Add result set metadata to prepared statement MD5 hash calculation (CASSANDRA-10786) * Add result set metadata to prepared statement MD5 hash calculation (CASSANDRA-10786)
* Refactor GcCompactionTest to avoid boxing (CASSANDRA-13941) * Refactor GcCompactionTest to avoid boxing (CASSANDRA-13941)
* Checksum sstable metadata (CASSANDRA-13321)
* Expose recent histograms in JmxHistograms (CASSANDRA-13642) * Expose recent histograms in JmxHistograms (CASSANDRA-13642)
* Fix buffer length comparison when decompressing in netty-based streaming (CASSANDRA-13899) * Fix buffer length comparison when decompressing in netty-based streaming (CASSANDRA-13899)
* Properly close StreamCompressionInputStream to release any ByteBuf (CASSANDRA-13906) * Properly close StreamCompressionInputStream to release any ByteBuf (CASSANDRA-13906)

View File

@ -19,11 +19,9 @@ package org.apache.cassandra.io.sstable.metadata;
import java.io.*; import java.io.*;
import java.util.*; import java.util.*;
import java.util.zip.CRC32;
import com.google.common.collect.Lists; import com.google.common.collect.Lists;
import com.google.common.hash.HashFunction;
import com.google.common.hash.Hasher;
import com.google.common.hash.Hashing;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
@ -41,11 +39,13 @@ import org.apache.cassandra.io.util.BufferedDataOutputStreamPlus;
import org.apache.cassandra.io.util.RandomAccessReader; import org.apache.cassandra.io.util.RandomAccessReader;
import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.FBUtilities;
import static org.apache.cassandra.utils.FBUtilities.updateChecksumInt;
/** /**
* Metadata serializer for SSTables {@code version >= 'k'}. * Metadata serializer for SSTables {@code version >= 'na'}.
* *
* <pre> * <pre>
* File format := | number of components (4 bytes) | toc | component1 | c1 hash | component2 | c2 hash | ... | * File format := | number of components (4 bytes) | crc | toc | crc | component1 | c1 crc | component2 | c2 crc | ... |
* toc := | component type (4 bytes) | position of component | * toc := | component type (4 bytes) | position of component |
* </pre> * </pre>
* *
@ -54,31 +54,40 @@ import org.apache.cassandra.utils.FBUtilities;
public class MetadataSerializer implements IMetadataSerializer public class MetadataSerializer implements IMetadataSerializer
{ {
private static final Logger logger = LoggerFactory.getLogger(MetadataSerializer.class); private static final Logger logger = LoggerFactory.getLogger(MetadataSerializer.class);
private static final HashFunction hashFunction = Hashing.md5();
private static final int CHECKSUM_LENGTH = 4; // CRC32
public void serialize(Map<MetadataType, MetadataComponent> components, DataOutputPlus out, Version version) throws IOException public void serialize(Map<MetadataType, MetadataComponent> components, DataOutputPlus out, Version version) throws IOException
{ {
boolean checksum = version.hasMetadataChecksum(); boolean checksum = version.hasMetadataChecksum();
CRC32 crc = new CRC32();
// sort components by type // sort components by type
List<MetadataComponent> sortedComponents = Lists.newArrayList(components.values()); List<MetadataComponent> sortedComponents = Lists.newArrayList(components.values());
Collections.sort(sortedComponents); Collections.sort(sortedComponents);
// write number of component // write number of component
out.writeInt(components.size()); out.writeInt(components.size());
updateChecksumInt(crc, components.size());
maybeWriteChecksum(crc, out, version);
// build and write toc // build and write toc
int lastPosition = 4 + (8 * sortedComponents.size()); int lastPosition = 4 + (8 * sortedComponents.size()) + (checksum ? 2 * CHECKSUM_LENGTH : 0);
Map<MetadataType, Integer> sizes = new EnumMap<>(MetadataType.class); Map<MetadataType, Integer> sizes = new EnumMap<>(MetadataType.class);
for (MetadataComponent component : sortedComponents) for (MetadataComponent component : sortedComponents)
{ {
MetadataType type = component.getType(); MetadataType type = component.getType();
// serialize type // serialize type
out.writeInt(type.ordinal()); out.writeInt(type.ordinal());
updateChecksumInt(crc, type.ordinal());
// serialize position // serialize position
out.writeInt(lastPosition); out.writeInt(lastPosition);
updateChecksumInt(crc, lastPosition);
int size = type.serializer.serializedSize(version, component); int size = type.serializer.serializedSize(version, component);
lastPosition += size + (checksum ? 8 : 0); // checksum is long lastPosition += size + (checksum ? CHECKSUM_LENGTH : 0);
sizes.put(type, size); sizes.put(type, size);
} }
maybeWriteChecksum(crc, out, version);
// serialize components // serialize components
for (MetadataComponent component : sortedComponents) for (MetadataComponent component : sortedComponents)
{ {
@ -89,11 +98,18 @@ public class MetadataSerializer implements IMetadataSerializer
bytes = dob.getData(); bytes = dob.getData();
} }
out.write(bytes); out.write(bytes);
if (checksum)
out.writeLong(hashFunction.hashBytes(bytes).asLong()); crc.reset(); crc.update(bytes);
maybeWriteChecksum(crc, out, version);
} }
} }
private static void maybeWriteChecksum(CRC32 crc, DataOutputPlus out, Version version) throws IOException
{
if (version.hasMetadataChecksum())
out.writeInt((int) crc.getValue());
}
public Map<MetadataType, MetadataComponent> deserialize( Descriptor descriptor, EnumSet<MetadataType> types) throws IOException public Map<MetadataType, MetadataComponent> deserialize( Descriptor descriptor, EnumSet<MetadataType> types) throws IOException
{ {
Map<MetadataType, MetadataComponent> components; Map<MetadataType, MetadataComponent> components;
@ -120,66 +136,87 @@ public class MetadataSerializer implements IMetadataSerializer
return deserialize(descriptor, EnumSet.of(type)).get(type); return deserialize(descriptor, EnumSet.of(type)).get(type);
} }
public Map<MetadataType, MetadataComponent> deserialize(Descriptor descriptor, FileDataInput in, EnumSet<MetadataType> types) throws IOException public Map<MetadataType, MetadataComponent> deserialize(Descriptor descriptor,
FileDataInput in,
EnumSet<MetadataType> selectedTypes)
throws IOException
{ {
int totalSize = (int) in.bytesRemaining(); boolean isChecksummed = descriptor.version.hasMetadataChecksum();
CRC32 crc = new CRC32();
/*
* Read TOC
*/
int length = (int) in.bytesRemaining();
int count = in.readInt();
updateChecksumInt(crc, count);
maybeValidateChecksum(crc, in, descriptor);
int[] ordinals = new int[count];
int[] offsets = new int[count];
int[] lengths = new int[count];
for (int i = 0; i < count; i++)
{
ordinals[i] = in.readInt();
updateChecksumInt(crc, ordinals[i]);
offsets[i] = in.readInt();
updateChecksumInt(crc, offsets[i]);
}
maybeValidateChecksum(crc, in, descriptor);
lengths[count - 1] = length - offsets[count - 1];
for (int i = 0; i < count - 1; i++)
lengths[i] = offsets[i + 1] - offsets[i];
/*
* Read components
*/
MetadataType[] allMetadataTypes = MetadataType.values();
Map<MetadataType, MetadataComponent> components = new EnumMap<>(MetadataType.class); Map<MetadataType, MetadataComponent> components = new EnumMap<>(MetadataType.class);
// read number of components
int numComponents = in.readInt();
// read toc
Map<MetadataType, Integer> toc = new EnumMap<>(MetadataType.class);
MetadataType[] values = MetadataType.values();
Map<MetadataType, Integer> lengths = new EnumMap<>(MetadataType.class);
int start = 0;
MetadataType lastType = null;
for (int i = 0; i < numComponents; i++)
{
int metadataTypeId = in.readInt();
int position = in.readInt();
toc.put(values[metadataTypeId], position); for (int i = 0; i < count; i++)
if (lastType != null)
lengths.put(lastType, position - start);
start = position;
lastType = values[metadataTypeId];
}
lengths.put(lastType, totalSize - start);
for (MetadataType type : types)
{ {
Integer offset = toc.get(type); MetadataType type = allMetadataTypes[ordinals[i]];
if (offset != null)
if (!selectedTypes.contains(type))
{ {
in.seek(offset); in.skipBytes(lengths[i]);
continue;
if (descriptor.version.hasMetadataChecksum())
{
int size = lengths.get(type) - 8; // 8 bytes checksum
byte[] bytes = new byte[size];
in.readFully(bytes);
MetadataComponent component;
try (DataInputBuffer dib = new DataInputBuffer(bytes))
{
component = type.serializer.deserialize(descriptor.version, dib);
}
long writtenChecksum = in.readLong();
if (writtenChecksum != hashFunction.hashBytes(bytes).asLong())
{
String filename = descriptor.filenameFor(Component.STATS);
throw new CorruptSSTableException(new IOException("Checksums do not match for " + filename), filename);
}
components.put(type, component);
}
else
{
MetadataComponent component = type.serializer.deserialize(descriptor.version, in);
components.put(type, component);
}
} }
byte[] buffer = new byte[isChecksummed ? lengths[i] - CHECKSUM_LENGTH : lengths[i]];
in.readFully(buffer);
crc.reset(); crc.update(buffer);
maybeValidateChecksum(crc, in, descriptor);
components.put(type, type.serializer.deserialize(descriptor.version, new DataInputBuffer(buffer)));
} }
return components; return components;
} }
private static void maybeValidateChecksum(CRC32 crc, FileDataInput in, Descriptor descriptor) throws IOException
{
if (!descriptor.version.hasMetadataChecksum())
return;
int actualChecksum = (int) crc.getValue();
int expectedChecksum = in.readInt();
if (actualChecksum != expectedChecksum)
{
String filename = descriptor.filenameFor(Component.STATS);
throw new CorruptSSTableException(new IOException("Checksums do not match for " + filename), filename);
}
}
public void mutateLevel(Descriptor descriptor, int newLevel) throws IOException public void mutateLevel(Descriptor descriptor, int newLevel) throws IOException
{ {
logger.trace("Mutating {} to level {}", descriptor.filenameFor(Component.STATS), newLevel); logger.trace("Mutating {} to level {}", descriptor.filenameFor(Component.STATS), newLevel);

Binary file not shown.

Before

Width:  |  Height:  |  Size: 5.1 KiB

After

Width:  |  Height:  |  Size: 5.1 KiB

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

Binary file not shown.

Before

Width:  |  Height:  |  Size: 5.2 KiB

After

Width:  |  Height:  |  Size: 5.1 KiB

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

Binary file not shown.

Before

Width:  |  Height:  |  Size: 5.7 KiB

After

Width:  |  Height:  |  Size: 5.6 KiB

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

Binary file not shown.

Before

Width:  |  Height:  |  Size: 5.7 KiB

After

Width:  |  Height:  |  Size: 5.4 KiB

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

View File

@ -1,8 +1,8 @@
Index.db Index.db
Data.db
CompressionInfo.db
Statistics.db
Summary.db
TOC.txt TOC.txt
Digest.crc32
Filter.db Filter.db
CompressionInfo.db
Summary.db
Data.db
Statistics.db
Digest.crc32

View File

@ -98,8 +98,13 @@ public class MetadataSerializerTest
String partitioner = RandomPartitioner.class.getCanonicalName(); String partitioner = RandomPartitioner.class.getCanonicalName();
double bfFpChance = 0.1; double bfFpChance = 0.1;
Map<MetadataType, MetadataComponent> originalMetadata = collector.finalizeMetadata(partitioner, bfFpChance, 0, null, SerializationHeader.make(cfm, Collections.emptyList())); return collector.finalizeMetadata(partitioner, bfFpChance, 0, null, SerializationHeader.make(cfm, Collections.emptyList()));
return originalMetadata; }
@Test
public void testMaReadMa() throws IOException
{
testOldReadsNew("ma", "ma");
} }
@Test @Test
@ -114,6 +119,12 @@ public class MetadataSerializerTest
testOldReadsNew("ma", "mc"); testOldReadsNew("ma", "mc");
} }
@Test
public void testMbReadMb() throws IOException
{
testOldReadsNew("mb", "mb");
}
@Test @Test
public void testMbReadMc() throws IOException public void testMbReadMc() throws IOException
{ {
@ -121,9 +132,15 @@ public class MetadataSerializerTest
} }
@Test @Test
public void testNaReadMc() throws IOException public void testMcReadMc() throws IOException
{ {
testOldReadsNew("mc", "na"); testOldReadsNew("mc", "mc");
}
@Test
public void testNaReadNa() throws IOException
{
testOldReadsNew("na", "na");
} }
public void testOldReadsNew(String oldV, String newV) throws IOException public void testOldReadsNew(String oldV, String newV) throws IOException
@ -146,11 +163,9 @@ public class MetadataSerializerTest
for (MetadataType type : MetadataType.values()) for (MetadataType type : MetadataType.values())
{ {
assertEquals(deserializedLa.get(type), deserializedLb.get(type)); assertEquals(deserializedLa.get(type), deserializedLb.get(type));
if (!originalMetadata.get(type).equals(deserializedLb.get(type)))
{ if (MetadataType.STATS != type)
// Currently only STATS can be different. Change if no longer the case assertEquals(originalMetadata.get(type), deserializedLb.get(type));
assertEquals(MetadataType.STATS, type);
}
} }
} }
} }