This is an automated email from the ASF dual-hosted git repository. discivigour pushed a commit to branch feat/blobRefRead in repository https://gitbox.apache.org/repos/asf/paimon.git
commit 4c5cfd8747a892473b8e0deb7e2c7dd93925d764 Author: umi <[email protected]> AuthorDate: Mon Jul 13 22:01:04 2026 +0800 [common] Optimize BlobRef data materialization --- .../main/java/org/apache/paimon/data/BlobRef.java | 13 ++++++- .../test/java/org/apache/paimon/data/BlobTest.java | 43 ++++++++++++++++++++++ 2 files changed, 54 insertions(+), 2 deletions(-) diff --git a/paimon-common/src/main/java/org/apache/paimon/data/BlobRef.java b/paimon-common/src/main/java/org/apache/paimon/data/BlobRef.java index 0248454ee9..87546031f4 100644 --- a/paimon-common/src/main/java/org/apache/paimon/data/BlobRef.java +++ b/paimon-common/src/main/java/org/apache/paimon/data/BlobRef.java @@ -45,8 +45,17 @@ public class BlobRef implements Blob { @Override public byte[] toData() { - try { - return IOUtils.readFully(newInputStream(), true); + try (SeekableInputStream inputStream = newInputStream()) { + long length = descriptor.length(); + if (length >= 0) { + if (length > Integer.MAX_VALUE) { + throw new IOException("Blob is too large to materialize as byte[]: " + length); + } + byte[] data = new byte[(int) length]; + IOUtils.readFully(inputStream, data); + return data; + } + return IOUtils.readFully(inputStream, false); } catch (IOException e) { throw new RuntimeException(e); } diff --git a/paimon-common/src/test/java/org/apache/paimon/data/BlobTest.java b/paimon-common/src/test/java/org/apache/paimon/data/BlobTest.java index 27703662da..8eb82ee689 100644 --- a/paimon-common/src/test/java/org/apache/paimon/data/BlobTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/data/BlobTest.java @@ -18,7 +18,9 @@ package org.apache.paimon.data; +import org.apache.paimon.fs.ByteArraySeekableStream; import org.apache.paimon.fs.local.LocalFileIO; +import org.apache.paimon.utils.UriReader; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -28,6 +30,7 @@ import java.io.File; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Paths; +import java.util.Arrays; import static org.assertj.core.api.Assertions.assertThat; @@ -79,4 +82,44 @@ public class BlobTest { Blob blob = Blob.fromHttp(uri); assertThat(blob).isInstanceOf(BlobRef.class); } + + @Test + public void testBlobRefReadsKnownLengthDirectly() { + byte[] data = new byte[8194]; + for (int i = 0; i < data.length; i++) { + data[i] = (byte) i; + } + TrackingSeekableInputStream inputStream = new TrackingSeekableInputStream(data); + UriReader uriReader = uri -> inputStream; + Blob blob = Blob.fromDescriptor(uriReader, new BlobDescriptor("test", 2, data.length - 2)); + + assertThat(blob.toData()).isEqualTo(Arrays.copyOfRange(data, 2, data.length)); + assertThat(inputStream.readCalls).isEqualTo(1); + assertThat(inputStream.maxRequestedBytes).isEqualTo(data.length - 2); + assertThat(inputStream.closed).isTrue(); + } + + private static class TrackingSeekableInputStream extends ByteArraySeekableStream { + + private int readCalls; + private int maxRequestedBytes; + private boolean closed; + + private TrackingSeekableInputStream(byte[] data) { + super(data); + } + + @Override + public int read(byte[] bytes, int offset, int length) throws IOException { + readCalls++; + maxRequestedBytes = Math.max(maxRequestedBytes, length); + return super.read(bytes, offset, length); + } + + @Override + public void close() throws IOException { + closed = true; + super.close(); + } + } }
