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 @@ -28,6 +28,8 @@
import org.apache.paimon.memory.MemorySegment;
import org.apache.paimon.utils.MutableObjectIterator;

import javax.annotation.Nullable;

import java.io.EOFException;
import java.io.IOException;
import java.util.Collections;
Expand All @@ -53,13 +55,19 @@ public class ChannelReaderInputView extends AbstractPagedInputView {
public ChannelReaderInputView(
FileIOChannel.ID id,
IOManager ioManager,
BlockCompressionFactory compressionCodecFactory,
@Nullable BlockCompressionFactory compressionCodecFactory,
int compressionBlockSize,
int numBlocks)
throws IOException {
this.numBlocksRemaining = numBlocks;
this.reader = ioManager.createBufferFileReader(id);
uncompressedBuffer = MemorySegment.wrap(new byte[compressionBlockSize]);
// spill-compression 'none' maps the factory to null: read plain blocks
if (compressionCodecFactory == null) {
decompressor = null;
compressedBuffer = uncompressedBuffer;
return;
}
decompressor = compressionCodecFactory.getDecompressor();
compressedBuffer =
MemorySegment.wrap(
Expand All @@ -79,6 +87,11 @@ protected MemorySegment nextSegment(MemorySegment current) throws IOException {

Buffer buffer = Buffer.create(compressedBuffer);
reader.readInto(buffer);
if (decompressor == null) {
this.currentSegmentLimit = buffer.getSize();
this.numBlocksRemaining--;
return compressedBuffer;
}
this.currentSegmentLimit =
decompressor.decompress(
buffer.getMemorySegment().getArray(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@
import org.apache.paimon.memory.Buffer;
import org.apache.paimon.memory.MemorySegment;

import javax.annotation.Nullable;

import java.io.Closeable;
import java.io.IOException;

Expand All @@ -47,13 +49,18 @@ public final class ChannelWriterOutputView extends AbstractPagedOutputView imple

public ChannelWriterOutputView(
BufferFileWriter writer,
BlockCompressionFactory compressionCodecFactory,
@Nullable BlockCompressionFactory compressionCodecFactory,
int compressionBlockSize) {
super(MemorySegment.wrap(new byte[compressionBlockSize]), compressionBlockSize);

compressor = compressionCodecFactory.getCompressor();
// spill-compression 'none' maps the factory to null: write plain blocks
compressor =
compressionCodecFactory == null ? null : compressionCodecFactory.getCompressor();
compressedBuffer =
MemorySegment.wrap(new byte[compressor.getMaxCompressedSize(compressionBlockSize)]);
compressor == null
? null
: MemorySegment.wrap(
new byte[compressor.getMaxCompressedSize(compressionBlockSize)]);
this.writer = writer;
}

Expand Down Expand Up @@ -90,6 +97,13 @@ protected MemorySegment nextSegment(MemorySegment current, int positionInCurrent
}

private void writeCompressed(MemorySegment current, int size) throws IOException {
if (compressor == null) {
writer.writeBlock(Buffer.create(current, size));
blockCount++;
numBytes += size;
numCompressedBytes += size;
return;
}
int compressedLen =
compressor.compress(current.getArray(), 0, size, compressedBuffer.getArray(), 0);
writer.writeBlock(Buffer.create(compressedBuffer, compressedLen));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@

import org.apache.paimon.compression.BlockCompressionFactory;
import org.apache.paimon.compression.BlockCompressionType;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.BinaryRowWriter;
import org.apache.paimon.data.serializer.BinaryRowSerializer;

import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -66,4 +68,39 @@ input, null, new BinaryRowSerializer(1))
}
}
}

@Test
public void testSpillCompressionNoneRoundTrip() throws Exception {
// spill-compression 'none' maps the factory to null: writing and reading must
// fall back to plain blocks instead of crashing on the missing codec
BinaryRowSerializer serializer = new BinaryRowSerializer(1);
BinaryRow row = new BinaryRow(1);
BinaryRowWriter writer = new BinaryRowWriter(row);
writer.writeInt(0, 42);
writer.complete();
try (IOManager ioManager = IOManager.create(tempDir.toString())) {
FileIOChannel.ID channel = ioManager.createChannel();
ChannelWriterOutputView output =
FileChannelUtil.createOutputView(ioManager, channel, null, BLOCK_SIZE);
BinaryRowSerializer ser = serializer.duplicate();
for (int i = 0; i < 100; i++) {
ser.serializeToPages(row, output);
}
output.close();

ChannelReaderInputView input =
new ChannelReaderInputView(
channel, ioManager, null, BLOCK_SIZE, output.getBlockCount());
try {
ChannelReaderInputViewIterator iterator =
new ChannelReaderInputViewIterator(input, null, serializer);
for (int i = 0; i < 100; i++) {
assertThat(iterator.next().getInt(0)).isEqualTo(42);
}
assertThat(iterator.next()).isNull();
} finally {
input.getChannel().closeAndDelete();
}
}
}
}
Loading