Create compression chunk when sending file only

patch by yukim; reviewed by Paulo Motta for CASSANDRA-10680
This commit is contained in:
Yuki Morishita 2015-11-09 23:01:31 -06:00
parent 9b97766057
commit 8385bb639a
6 changed files with 87 additions and 17 deletions

View File

@ -1,4 +1,5 @@
2.1.12
* Create compression chunk for sending file only (CASSANDRA-10680)
* Make buffered read size configurable (CASSANDRA-10249)
* Forbid compact clustering column type changes in ALTER TABLE (CASSANDRA-8879)
* Reject incremental repair with subrange repair (CASSANDRA-10422)

View File

@ -228,6 +228,36 @@ public class CompressionMetadata
return new Chunk(chunkOffset, (int) (nextChunkOffset - chunkOffset - 4)); // "4" bytes reserved for checksum
}
/**
* @param sections Collection of sections in uncompressed file. Should not contain sections that overlap each other.
* @return Total chunk size in bytes for given sections including checksum.
*/
public long getTotalSizeForSections(Collection<Pair<Long, Long>> sections)
{
long size = 0;
long lastOffset = -1;
for (Pair<Long, Long> section : sections)
{
int startIndex = (int) (section.left / parameters.chunkLength());
int endIndex = (int) (section.right / parameters.chunkLength());
endIndex = section.right % parameters.chunkLength() == 0 ? endIndex - 1 : endIndex;
for (int i = startIndex; i <= endIndex; i++)
{
long offset = i * 8L;
long chunkOffset = chunkOffsets.getLong(offset);
if (chunkOffset > lastOffset)
{
lastOffset = chunkOffset;
long nextChunkOffset = offset + 8 == chunkOffsetsSize
? compressedFileLength
: chunkOffsets.getLong(offset + 8);
size += (nextChunkOffset - chunkOffset);
}
}
}
return size;
}
/**
* @param sections Collection of sections in uncompressed file
* @return Array of chunks which corresponds to given sections of uncompressed file, sorted by chunk offset

View File

@ -37,7 +37,7 @@ import org.apache.cassandra.utils.UUIDSerializer;
*/
public class FileMessageHeader
{
public static IVersionedSerializer<FileMessageHeader> serializer = new FileMessageHeaderSerializer();
public static FileMessageHeaderSerializer serializer = new FileMessageHeaderSerializer();
public final UUID cfId;
public final int sequenceNumber;
@ -45,7 +45,13 @@ public class FileMessageHeader
public final String version;
public final long estimatedKeys;
public final List<Pair<Long, Long>> sections;
/**
* Compression info for SSTable to send. Can be null if SSTable is not compressed.
* On sender, this field is always null to avoid holding large number of Chunks.
* Use compressionMetadata instead.
*/
public final CompressionInfo compressionInfo;
private final CompressionMetadata compressionMetadata;
public final long repairedAt;
public FileMessageHeader(UUID cfId,
@ -62,9 +68,33 @@ public class FileMessageHeader
this.estimatedKeys = estimatedKeys;
this.sections = sections;
this.compressionInfo = compressionInfo;
this.compressionMetadata = null;
this.repairedAt = repairedAt;
}
public FileMessageHeader(UUID cfId,
int sequenceNumber,
String version,
long estimatedKeys,
List<Pair<Long, Long>> sections,
CompressionMetadata compressionMetadata,
long repairedAt)
{
this.cfId = cfId;
this.sequenceNumber = sequenceNumber;
this.version = version;
this.estimatedKeys = estimatedKeys;
this.sections = sections;
this.compressionInfo = null;
this.compressionMetadata = compressionMetadata;
this.repairedAt = repairedAt;
}
public boolean isCompressed()
{
return compressionInfo != null || compressionMetadata != null;
}
/**
* @return total file size to transfer in bytes
*/
@ -77,6 +107,10 @@ public class FileMessageHeader
for (CompressionMetadata.Chunk chunk : compressionInfo.chunks)
size += chunk.length + 4; // 4 bytes for CRC
}
else if (compressionMetadata != null)
{
size = compressionMetadata.getTotalSizeForSections(sections);
}
else
{
for (Pair<Long, Long> section : sections)
@ -94,7 +128,7 @@ public class FileMessageHeader
sb.append(", version: ").append(version);
sb.append(", estimated keys: ").append(estimatedKeys);
sb.append(", transfer size: ").append(size());
sb.append(", compressed?: ").append(compressionInfo != null);
sb.append(", compressed?: ").append(isCompressed());
sb.append(", repairedAt: ").append(repairedAt);
sb.append(')');
return sb.toString();
@ -117,9 +151,9 @@ public class FileMessageHeader
return result;
}
static class FileMessageHeaderSerializer implements IVersionedSerializer<FileMessageHeader>
static class FileMessageHeaderSerializer
{
public void serialize(FileMessageHeader header, DataOutputPlus out, int version) throws IOException
public CompressionInfo serialize(FileMessageHeader header, DataOutputPlus out, int version) throws IOException
{
UUIDSerializer.serializer.serialize(header.cfId, out, version);
out.writeInt(header.sequenceNumber);
@ -132,8 +166,13 @@ public class FileMessageHeader
out.writeLong(section.left);
out.writeLong(section.right);
}
CompressionInfo.serializer.serialize(header.compressionInfo, out, version);
// construct CompressionInfo here to avoid holding large number of Chunks on heap.
CompressionInfo compressionInfo = null;
if (header.compressionMetadata != null)
compressionInfo = new CompressionInfo(header.compressionMetadata.getChunksForSections(header.sections), header.compressionMetadata.parameters);
CompressionInfo.serializer.serialize(compressionInfo, out, version);
out.writeLong(header.repairedAt);
return compressionInfo;
}
public FileMessageHeader deserialize(DataInput in, int version) throws IOException

View File

@ -40,7 +40,7 @@ public class IncomingFileMessage extends StreamMessage
{
DataInputStream input = new DataInputStream(Channels.newInputStream(in));
FileMessageHeader header = FileMessageHeader.serializer.deserialize(input, version);
StreamReader reader = header.compressionInfo == null ? new StreamReader(header, session)
StreamReader reader = !header.isCompressed() ? new StreamReader(header, session)
: new CompressedStreamReader(header, session);
try

View File

@ -63,18 +63,12 @@ public class OutgoingFileMessage extends StreamMessage
SSTableReader sstable = ref.get();
filename = sstable.getFilename();
CompressionInfo compressionInfo = null;
if (sstable.compression)
{
CompressionMetadata meta = sstable.getCompressionMetadata();
compressionInfo = new CompressionInfo(meta.getChunksForSections(sections), meta.parameters);
}
this.header = new FileMessageHeader(sstable.metadata.cfId,
sequenceNumber,
sstable.descriptor.version.toString(),
estimatedKeys,
sections,
compressionInfo,
sstable.compression ? sstable.getCompressionMetadata() : null,
repairedAt);
}
@ -85,13 +79,12 @@ public class OutgoingFileMessage extends StreamMessage
return;
}
FileMessageHeader.serializer.serialize(header, out, version);
CompressionInfo compressionInfo = FileMessageHeader.serializer.serialize(header, out, version);
final SSTableReader reader = ref.get();
StreamWriter writer = header.compressionInfo == null ?
StreamWriter writer = compressionInfo == null ?
new StreamWriter(reader, header.sections, session) :
new CompressedStreamWriter(reader, header.sections,
header.compressionInfo, session);
new CompressedStreamWriter(reader, header.sections, compressionInfo, session);
writer.write(out.getChannel());
}

View File

@ -37,6 +37,8 @@ import org.apache.cassandra.io.sstable.Descriptor;
import org.apache.cassandra.io.sstable.metadata.MetadataCollector;
import org.apache.cassandra.utils.Pair;
import static org.junit.Assert.assertEquals;
/**
*/
public class CompressedInputStreamTest
@ -86,6 +88,11 @@ public class CompressedInputStreamTest
sections.add(Pair.create(position, position + 8));
}
CompressionMetadata.Chunk[] chunks = comp.getChunksForSections(sections);
long totalSize = comp.getTotalSizeForSections(sections);
long expectedSize = 0;
for (CompressionMetadata.Chunk c : chunks)
expectedSize += c.length + 4;
assertEquals(expectedSize, totalSize);
// buffer up only relevant parts of file
int size = 0;