diff --git a/paimon-core/src/main/java/org/apache/paimon/disk/ChannelReaderInputView.java b/paimon-core/src/main/java/org/apache/paimon/disk/ChannelReaderInputView.java index aeb6e08c2afb..9449f64e646f 100644 --- a/paimon-core/src/main/java/org/apache/paimon/disk/ChannelReaderInputView.java +++ b/paimon-core/src/main/java/org/apache/paimon/disk/ChannelReaderInputView.java @@ -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; @@ -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( @@ -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(), diff --git a/paimon-core/src/main/java/org/apache/paimon/disk/ChannelWriterOutputView.java b/paimon-core/src/main/java/org/apache/paimon/disk/ChannelWriterOutputView.java index 3f759cdcca08..6fcbbc9f78a7 100644 --- a/paimon-core/src/main/java/org/apache/paimon/disk/ChannelWriterOutputView.java +++ b/paimon-core/src/main/java/org/apache/paimon/disk/ChannelWriterOutputView.java @@ -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; @@ -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; } @@ -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)); diff --git a/paimon-core/src/test/java/org/apache/paimon/disk/ChannelWriterOutputViewTest.java b/paimon-core/src/test/java/org/apache/paimon/disk/ChannelWriterOutputViewTest.java index 9ad0c562e132..9698d9a5033a 100644 --- a/paimon-core/src/test/java/org/apache/paimon/disk/ChannelWriterOutputViewTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/disk/ChannelWriterOutputViewTest.java @@ -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; @@ -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(); + } + } + } }