diff --git a/src/java/org/apache/cassandra/db/filter/SSTableNamesIterator.java b/src/java/org/apache/cassandra/db/filter/SSTableNamesIterator.java index 7e4bdc77ae..820977d737 100644 --- a/src/java/org/apache/cassandra/db/filter/SSTableNamesIterator.java +++ b/src/java/org/apache/cassandra/db/filter/SSTableNamesIterator.java @@ -70,7 +70,9 @@ public class SSTableNamesIterator extends SimpleAbstractColumnIterator implement file = ssTable.getFileDataInput(decoratedKey, DatabaseDescriptor.getIndexedReadBufferSizeInKB() * 1024); if (file == null) return; - DecoratedKey keyInDisk = ssTable.getPartitioner().convertFromDiskFormat(FBUtilities.readShortByteArray(file)); + DecoratedKey keyInDisk = SSTableReader.decodeKey(ssTable.getPartitioner(), + ssTable.getDescriptor(), + FBUtilities.readShortByteArray(file)); assert keyInDisk.equals(decoratedKey) : String.format("%s != %s in %s", keyInDisk, decoratedKey, file.getPath()); SSTableReader.readRowSize(file, ssTable.getDescriptor()); diff --git a/src/java/org/apache/cassandra/db/filter/SSTableSliceIterator.java b/src/java/org/apache/cassandra/db/filter/SSTableSliceIterator.java index 56088a1af3..b591d16d67 100644 --- a/src/java/org/apache/cassandra/db/filter/SSTableSliceIterator.java +++ b/src/java/org/apache/cassandra/db/filter/SSTableSliceIterator.java @@ -88,7 +88,9 @@ class SSTableSliceIterator extends AbstractIterator implements IColumnI return; try { - DecoratedKey keyInDisk = ssTable.getPartitioner().convertFromDiskFormat(FBUtilities.readShortByteArray(file)); + DecoratedKey keyInDisk = SSTableReader.decodeKey(ssTable.getPartitioner(), + ssTable.getDescriptor(), + FBUtilities.readShortByteArray(file)); assert keyInDisk.equals(decoratedKey) : String.format("%s != %s in %s", keyInDisk, decoratedKey, file.getPath()); SSTableReader.readRowSize(file, ssTable.getDescriptor()); diff --git a/src/java/org/apache/cassandra/dht/AbstractByteOrderedPartitioner.java b/src/java/org/apache/cassandra/dht/AbstractByteOrderedPartitioner.java index 6e5880af47..5d8a35aa4a 100644 --- a/src/java/org/apache/cassandra/dht/AbstractByteOrderedPartitioner.java +++ b/src/java/org/apache/cassandra/dht/AbstractByteOrderedPartitioner.java @@ -47,11 +47,6 @@ public abstract class AbstractByteOrderedPartitioner implements IPartitioner(getToken(key), key); } - public byte[] convertToDiskFormat(DecoratedKey key) - { - return key.key; - } - public BytesToken midpoint(BytesToken ltoken, BytesToken rtoken) { int sigbytes = Math.max(ltoken.token.length, rtoken.token.length); diff --git a/src/java/org/apache/cassandra/dht/IPartitioner.java b/src/java/org/apache/cassandra/dht/IPartitioner.java index 04e35def0f..c9304c79db 100644 --- a/src/java/org/apache/cassandra/dht/IPartitioner.java +++ b/src/java/org/apache/cassandra/dht/IPartitioner.java @@ -25,20 +25,14 @@ import org.apache.cassandra.db.DecoratedKey; public interface IPartitioner { /** + * @Deprecated: Used by SSTables before version 'e'. + * * Convert the on disk representation to a DecoratedKey object * @param key On disk representation * @return DecoratedKey object */ public DecoratedKey convertFromDiskFormat(byte[] key); - /** - * Convert the DecoratedKey to the on disk format used for - * this partitioner. - * @param key The DecoratedKey in question - * @return - */ - public byte[] convertToDiskFormat(DecoratedKey key); - /** * Transform key to object representation of the on-disk format. * diff --git a/src/java/org/apache/cassandra/dht/LocalPartitioner.java b/src/java/org/apache/cassandra/dht/LocalPartitioner.java index da44d7672c..1dcd22fc3c 100644 --- a/src/java/org/apache/cassandra/dht/LocalPartitioner.java +++ b/src/java/org/apache/cassandra/dht/LocalPartitioner.java @@ -19,11 +19,6 @@ public class LocalPartitioner implements IPartitioner return decorateKey(key); } - public byte[] convertToDiskFormat(DecoratedKey key) - { - return key.token.token; - } - public DecoratedKey decorateKey(byte[] key) { return new DecoratedKey(getToken(key), key); diff --git a/src/java/org/apache/cassandra/dht/OrderPreservingPartitioner.java b/src/java/org/apache/cassandra/dht/OrderPreservingPartitioner.java index 27aa614a9c..b475c501d4 100644 --- a/src/java/org/apache/cassandra/dht/OrderPreservingPartitioner.java +++ b/src/java/org/apache/cassandra/dht/OrderPreservingPartitioner.java @@ -46,11 +46,6 @@ public class OrderPreservingPartitioner implements IPartitioner return new DecoratedKey(getToken(key), key); } - public byte[] convertToDiskFormat(DecoratedKey key) - { - return key.key; - } - public StringToken midpoint(StringToken ltoken, StringToken rtoken) { int sigchars = Math.max(ltoken.token.length(), rtoken.token.length()); diff --git a/src/java/org/apache/cassandra/dht/RandomPartitioner.java b/src/java/org/apache/cassandra/dht/RandomPartitioner.java index 1f9aaaa800..a6a3773ee6 100644 --- a/src/java/org/apache/cassandra/dht/RandomPartitioner.java +++ b/src/java/org/apache/cassandra/dht/RandomPartitioner.java @@ -64,21 +64,6 @@ public class RandomPartitioner implements IPartitioner return new DecoratedKey(new BigIntegerToken(token), key); } - public byte[] convertToDiskFormat(DecoratedKey key) - { - // encode token prefix and calculate final length (with delimiter) - byte[] prefix = key.token.toString().getBytes(UTF_8); - int length = prefix.length + 1 + key.key.length; - assert length <= FBUtilities.MAX_UNSIGNED_SHORT; - - // copy into output bytes - byte[] todisk = new byte[length]; - System.arraycopy(prefix, 0, todisk, 0, prefix.length); - todisk[prefix.length] = DELIMITER_BYTE; - System.arraycopy(key.key, 0, todisk, prefix.length + 1, key.key.length); - return todisk; - } - public BigIntegerToken midpoint(BigIntegerToken ltoken, BigIntegerToken rtoken) { Pair midpair = FBUtilities.midpoint(ltoken.token, rtoken.token, 127); diff --git a/src/java/org/apache/cassandra/io/sstable/Descriptor.java b/src/java/org/apache/cassandra/io/sstable/Descriptor.java index 8893c441d1..a3736d9dfe 100644 --- a/src/java/org/apache/cassandra/io/sstable/Descriptor.java +++ b/src/java/org/apache/cassandra/io/sstable/Descriptor.java @@ -15,7 +15,7 @@ import com.google.common.base.Objects; public class Descriptor { public static final String LEGACY_VERSION = "a"; - public static final String CURRENT_VERSION = "d"; + public static final String CURRENT_VERSION = "e"; public final File directory; public final String version; @@ -160,6 +160,11 @@ public class Descriptor return version.compareTo("d") < 0; } + public boolean hasEncodedKeys() + { + return version.compareTo("e") < 0; + } + public boolean isLatestVersion() { return version.compareTo(CURRENT_VERSION) == 0; diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java index 96a00aaccd..5a4421ee8d 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java @@ -266,7 +266,7 @@ public class SSTableReader extends SSTable implements Comparable if (indexPosition == indexSize) break; - DecoratedKey decoratedKey = partitioner.convertFromDiskFormat(FBUtilities.readShortByteArray(input)); + DecoratedKey decoratedKey = decodeKey(partitioner, desc, FBUtilities.readShortByteArray(input)); if (recreatebloom) bf.add(decoratedKey.key); long dataPosition = input.readLong(); @@ -414,7 +414,7 @@ public class SSTableReader extends SSTable implements Comparable while (!input.isEOF()) { // read key & data position from index entry - DecoratedKey indexDecoratedKey = partitioner.convertFromDiskFormat(FBUtilities.readShortByteArray(input)); + DecoratedKey indexDecoratedKey = decodeKey(partitioner, desc, FBUtilities.readShortByteArray(input)); long dataPosition = input.readLong(); int comparison = indexDecoratedKey.compareTo(decoratedKey); @@ -556,6 +556,16 @@ public class SSTableReader extends SSTable implements Comparable return in.readLong(); } + /** + * Conditionally use the deprecated 'IPartitioner.convertFromDiskFormat' method. + */ + public static DecoratedKey decodeKey(IPartitioner p, Descriptor d, byte[] bytes) + { + if (d.hasEncodedKeys()) + return p.convertFromDiskFormat(bytes); + return p.decorateKey(bytes); + } + /** * TODO: Move someplace reusable */ diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java b/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java index b9a49e61c2..3b284757ec 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java @@ -168,7 +168,9 @@ public class SSTableScanner implements Iterator, Closeable file.seek(finishedAt); assert !file.isEOF(); - DecoratedKey key = StorageService.getPartitioner().convertFromDiskFormat(FBUtilities.readShortByteArray(file)); + DecoratedKey key = SSTableReader.decodeKey(sstable.getPartitioner(), + sstable.getDescriptor(), + FBUtilities.readShortByteArray(file)); long dataSize = SSTableReader.readRowSize(file, sstable.getDescriptor()); dataStart = file.getFilePointer(); finishedAt = dataStart + dataSize; diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java index 2ba7bec6d4..9f8d608989 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java @@ -99,7 +99,7 @@ public class SSTableWriter extends SSTable public void append(AbstractCompactedRow row) throws IOException { long currentPosition = beforeAppend(row.key); - FBUtilities.writeShortByteArray(partitioner.convertToDiskFormat(row.key), dataFile); + FBUtilities.writeShortByteArray(row.key.key, dataFile); row.write(dataFile); afterAppend(row.key, currentPosition); } @@ -107,7 +107,7 @@ public class SSTableWriter extends SSTable public void append(DecoratedKey decoratedKey, ColumnFamily cf) throws IOException { long startPosition = beforeAppend(decoratedKey); - FBUtilities.writeShortByteArray(partitioner.convertToDiskFormat(decoratedKey), dataFile); + FBUtilities.writeShortByteArray(decoratedKey.key, dataFile); // write placeholder for the row size, since we don't know it yet long sizePosition = dataFile.getFilePointer(); dataFile.writeLong(-1); @@ -125,7 +125,7 @@ public class SSTableWriter extends SSTable public void append(DecoratedKey decoratedKey, byte[] value) throws IOException { long currentPosition = beforeAppend(decoratedKey); - FBUtilities.writeShortByteArray(partitioner.convertToDiskFormat(decoratedKey), dataFile); + FBUtilities.writeShortByteArray(decoratedKey.key, dataFile); assert value.length > 0; dataFile.writeLong(value.length); dataFile.write(value); @@ -237,7 +237,7 @@ public class SSTableWriter extends SSTable long dataPosition = 0; while (dataPosition < dfile.length()) { - key = StorageService.getPartitioner().convertFromDiskFormat(FBUtilities.readShortByteArray(dfile)); + key = SSTableReader.decodeKey(StorageService.getPartitioner(), desc, FBUtilities.readShortByteArray(dfile)); long dataSize = SSTableReader.readRowSize(dfile, desc); iwriter.afterAppend(key, dataPosition); dataPosition = dfile.getFilePointer() + dataSize; @@ -301,7 +301,7 @@ public class SSTableWriter extends SSTable { bf.add(key.key); long indexPosition = indexFile.getFilePointer(); - FBUtilities.writeShortByteArray(partitioner.convertToDiskFormat(key), indexFile); + FBUtilities.writeShortByteArray(key.key, indexFile); indexFile.writeLong(dataPosition); if (logger.isTraceEnabled()) logger.trace("wrote index of " + key + " at " + indexPosition); diff --git a/src/java/org/apache/cassandra/tools/SSTableExport.java b/src/java/org/apache/cassandra/tools/SSTableExport.java index 23b3b0d0d2..8935af7a85 100644 --- a/src/java/org/apache/cassandra/tools/SSTableExport.java +++ b/src/java/org/apache/cassandra/tools/SSTableExport.java @@ -31,6 +31,7 @@ import org.apache.cassandra.db.IColumn; import org.apache.cassandra.db.TimestampClock; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.io.sstable.Descriptor; import org.apache.cassandra.io.sstable.SSTable; import org.apache.cassandra.io.sstable.SSTableIdentityIterator; import org.apache.cassandra.io.sstable.SSTableReader; @@ -157,10 +158,13 @@ public class SSTableExport throws IOException { IPartitioner partitioner = StorageService.getPartitioner(); + Descriptor desc = Descriptor.fromFilename(ssTableFile); BufferedRandomAccessFile input = new BufferedRandomAccessFile(SSTable.indexFilename(ssTableFile), "r"); while (!input.isEOF()) { - DecoratedKey decoratedKey = partitioner.convertFromDiskFormat(FBUtilities.readShortByteArray(input)); + DecoratedKey decoratedKey = SSTableReader.decodeKey(partitioner, + desc, + FBUtilities.readShortByteArray(input)); long dataPosition = input.readLong(); outs.println(bytesToHex(decoratedKey.key)); } diff --git a/test/conf/cassandra.yaml b/test/conf/cassandra.yaml index c59af3d2a1..0356a8f138 100644 --- a/test/conf/cassandra.yaml +++ b/test/conf/cassandra.yaml @@ -1,3 +1,7 @@ +# +# Warning! +# Consider the effects on 'o.a.c.i.s.LegacySSTableTest' before changing schemas in this file. +# cluster_name: Test Cluster in_memory_compaction_limit_in_mb: 1 commitlog_sync: batch diff --git a/test/unit/org/apache/cassandra/dht/PartitionerTestCase.java b/test/unit/org/apache/cassandra/dht/PartitionerTestCase.java index 568c493334..5280896287 100644 --- a/test/unit/org/apache/cassandra/dht/PartitionerTestCase.java +++ b/test/unit/org/apache/cassandra/dht/PartitionerTestCase.java @@ -100,15 +100,6 @@ public abstract class PartitionerTestCase assertMidpoint(tok("bbb"), tok("a"), 16); } - @Test - public void testDiskFormat() - { - byte[] key = "key".getBytes(); - DecoratedKey decKey = partitioner.decorateKey(key); - DecoratedKey result = partitioner.convertFromDiskFormat(partitioner.convertToDiskFormat(decKey)); - assertEquals(decKey, result); - } - @Test public void testTokenFactoryBytes() { diff --git a/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java b/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java index 46de761c56..99f4a838c3 100644 --- a/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java @@ -48,7 +48,7 @@ public class LegacySSTableTest extends CleanupHelper { String scp = System.getProperty(LEGACY_SSTABLE_PROP); assert scp != null; - LEGACY_SSTABLE_ROOT = new File(scp); + LEGACY_SSTABLE_ROOT = new File(scp).getAbsoluteFile(); assert LEGACY_SSTABLE_ROOT.isDirectory(); TEST_DATA = new HashMap(); @@ -75,7 +75,7 @@ public class LegacySSTableTest extends CleanupHelper public void buildTestSSTable() throws IOException { // write the output in a version specific directory - SSTable.Descriptor dest = getDescriptor(SSTable.Descriptor.CURRENT_VERSION); + Descriptor dest = getDescriptor(Descriptor.CURRENT_VERSION); assert dest.directory.mkdirs() : "Could not create " + dest.directory + ". Might it already exist?"; SSTableReader ssTable = SSTableUtils.writeRawSSTable(new File(dest.filenameFor(SSTable.COMPONENT_DATA)), @@ -88,22 +88,33 @@ public class LegacySSTableTest extends CleanupHelper } */ - /** - * Between version b and c, on disk bloom filters became incompatible, and needed to be regenerated. - */ @Test - public void testVerB() throws IOException + public void testVersions() throws IOException { - SSTableReader reader = SSTableReader.open(getDescriptor("b")); + for (File version : LEGACY_SSTABLE_ROOT.listFiles()) + testVersion(version.getName()); + } - List keys = new ArrayList(TEST_DATA.keySet()); - Collections.shuffle(keys); - BufferedRandomAccessFile file = new BufferedRandomAccessFile(reader.getFilename(), "r"); - for (byte[] key : keys) + public void testVersion(String version) + { + try { - // confirm that the bloom filter does not reject any keys - file.seek(reader.getPosition(reader.partitioner.decorateKey(key), SSTableReader.Operator.EQ)); - assert Arrays.equals(key, FBUtilities.readShortByteArray(file)); + SSTableReader reader = SSTableReader.open(getDescriptor(version)); + + List keys = new ArrayList(TEST_DATA.keySet()); + Collections.shuffle(keys); + BufferedRandomAccessFile file = new BufferedRandomAccessFile(reader.getFilename(), "r"); + for (byte[] key : keys) + { + // confirm that the bloom filter does not reject any keys + file.seek(reader.getPosition(reader.partitioner.decorateKey(key), SSTableReader.Operator.EQ)); + assert Arrays.equals(key, FBUtilities.readShortByteArray(file)); + } + } + catch (Throwable e) + { + System.err.println("Failed to read " + version); + e.printStackTrace(System.err); } } } diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java index 9691d5608c..3e179e80cc 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java @@ -97,7 +97,9 @@ public class SSTableReaderTest extends CleanupHelper { DecoratedKey dk = Util.dk(String.valueOf(j)); FileDataInput file = sstable.getFileDataInput(dk, DatabaseDescriptor.getIndexedReadBufferSizeInKB() * 1024); - DecoratedKey keyInDisk = sstable.getPartitioner().convertFromDiskFormat(FBUtilities.readShortByteArray(file)); + DecoratedKey keyInDisk = SSTableReader.decodeKey(sstable.getPartitioner(), + sstable.getDescriptor(), + FBUtilities.readShortByteArray(file)); assert keyInDisk.equals(dk) : String.format("%s != %s in %s", keyInDisk, dk, file.getPath()); }