This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new b52d83bb88 [core] Support configurable BLOB copy buffer size (#8722)
b52d83bb88 is described below

commit b52d83bb8888be492096f3178caaf95c9004ee7d
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Jul 20 16:20:51 2026 +0800

    [core] Support configurable BLOB copy buffer size (#8722)
    
    `BlobFormatWriter` copies BLOB payloads using a hard-coded 4 KiB buffer
    (`new byte[4096]`), which cannot be tuned.
    
    This PR adds a new option `blob.copy-buffer-size` to make it
    configurable, and threads it through all `BlobFormatWriter` construction
    paths (append / primary-key writes and BLOB compaction). The default
    stays **4 KiB**, so behavior is unchanged unless the option is set.
---
 docs/generated/core_configuration.html             |   6 +
 .../main/java/org/apache/paimon/CoreOptions.java   |  23 ++
 .../append/DedicatedFormatRollingFileWriter.java   |   3 +-
 .../paimon/append/MultipleBlobFileWriter.java      |   5 +-
 .../DataEvolutionBlobCompactTask.java              |   2 +-
 .../paimon/blob/PrimaryKeyBlobExternalizer.java    |  18 +-
 .../paimon/io/KeyValueFileWriterFactory.java       |   3 +-
 .../apache/paimon/operation/BlobFileContext.java   |  12 +-
 .../java/org/apache/paimon/CoreOptionsTest.java    |  27 ++
 .../org/apache/paimon/append/BlobUpdateTest.java   |   3 +-
 .../blob/PrimaryKeyBlobExternalizerTest.java       |  37 ++-
 .../apache/paimon/format/blob/BlobFileFormat.java  |  11 +-
 .../paimon/format/blob/BlobFileFormatFactory.java  |   5 +-
 .../paimon/format/blob/BlobFormatWriter.java       |  38 +--
 .../paimon/format/blob/BlobFileFormatTest.java     |  15 +-
 .../paimon/format/blob/BlobFormatWriterTest.java   | 338 ++++++++++++++-------
 .../pypaimon/common/options/core_options.py        |  17 ++
 paimon-python/pypaimon/tests/blob_test.py          |  30 ++
 paimon-python/pypaimon/write/blob_format_writer.py |  11 +-
 .../pypaimon/write/writer/blob_file_writer.py      |   4 +-
 paimon-python/pypaimon/write/writer/blob_writer.py |   9 +-
 21 files changed, 443 insertions(+), 174 deletions(-)

diff --git a/docs/generated/core_configuration.html 
b/docs/generated/core_configuration.html
index 29947a5e44..8875140929 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -104,6 +104,12 @@ under the License.
             <td>Boolean</td>
             <td>Whether to write NULL for a descriptor BLOB value when the 
referenced file or HTTP resource does not exist during Flink writes. When 
false, the write fails when the descriptor is read.</td>
         </tr>
+        <tr>
+            <td><h5>blob.copy-buffer-size</h5></td>
+            <td style="word-wrap: break-word;">4 kb</td>
+            <td>MemorySize</td>
+            <td>Buffer size used when copying BLOB payloads into BLOB 
files.</td>
+        </tr>
         <tr>
             <td><h5>blob.split-by-file-size</h5></td>
             <td style="word-wrap: break-word;">(none)</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java 
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 459f4a1545..0f26a2a79b 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -846,6 +846,14 @@ public class CoreOptions implements Serializable {
                                             "Whether to consider blob file 
size as a factor when performing scan splitting.")
                                     .build());
 
+    // Keep this default in sync with 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE (= 4 * 1024).
+    public static final ConfigOption<MemorySize> BLOB_COPY_BUFFER_SIZE =
+            key("blob.copy-buffer-size")
+                    .memoryType()
+                    .defaultValue(MemorySize.parse("4 kb"))
+                    .withDescription(
+                            "Buffer size used when copying BLOB payloads into 
BLOB files.");
+
     public static final ConfigOption<Integer> 
NUM_SORTED_RUNS_COMPACTION_TRIGGER =
             key("num-sorted-run.compaction-trigger")
                     .intType()
@@ -3302,6 +3310,21 @@ public class CoreOptions implements Serializable {
                 .orElse(targetFileSize(false));
     }
 
+    public int blobCopyBufferSize() {
+        return 
checkedBlobCopyBufferSize(options.get(BLOB_COPY_BUFFER_SIZE).getBytes());
+    }
+
+    /** Validates {@link #BLOB_COPY_BUFFER_SIZE} bytes and narrows to a 
positive int. */
+    public static int checkedBlobCopyBufferSize(long bytes) {
+        checkArgument(
+                bytes > 0 && bytes <= Integer.MAX_VALUE,
+                "'%s' must be between 1 byte and %s bytes, but was %s bytes.",
+                BLOB_COPY_BUFFER_SIZE.key(),
+                Integer.MAX_VALUE,
+                bytes);
+        return (int) bytes;
+    }
+
     public boolean blobSplitByFileSize() {
         return options.getOptional(BLOB_SPLIT_BY_FILE_SIZE)
                 .orElse(!options.get(BLOB_AS_DESCRIPTOR));
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
 
b/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
index 86df6b10b2..9b8897862a 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
@@ -201,7 +201,8 @@ public class DedicatedFormatRollingFileWriter
                                     context.blobInlineFields(),
                                     context.writeNullOnMissingFile(),
                                     context.writeNullOnFetchFailure(),
-                                    context.blobFetchMetricReporter());
+                                    context.blobFetchMetricReporter(),
+                                    context.copyBufferSize());
         } else {
             this.blobWriterFactory = null;
         }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/MultipleBlobFileWriter.java
 
b/paimon-core/src/main/java/org/apache/paimon/append/MultipleBlobFileWriter.java
index d952018ee1..39d6c6eac8 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/append/MultipleBlobFileWriter.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/append/MultipleBlobFileWriter.java
@@ -67,11 +67,12 @@ public class MultipleBlobFileWriter implements Closeable {
             Set<String> blobInlineFields,
             boolean writeNullOnMissingFile,
             boolean writeNullOnFetchFailure,
-            BlobFetchMetricReporter blobFetchMetricReporter) {
+            BlobFetchMetricReporter blobFetchMetricReporter,
+            int copyBufferSize) {
         RowType blobRowType = new RowType(fieldsInBlobFile(writeSchema, 
blobInlineFields));
         this.blobWriters = new ArrayList<>();
         for (String blobFieldName : blobRowType.getFieldNames()) {
-            BlobFileFormat blobFileFormat = new BlobFileFormat();
+            BlobFileFormat blobFileFormat = new BlobFileFormat(false, 
copyBufferSize);
             blobFileFormat.setWriteConsumer(blobConsumer);
             blobFileFormat.setWriteNullOnMissingFile(writeNullOnMissingFile);
             blobFileFormat.setWriteNullOnFetchFailure(writeNullOnFetchFailure);
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionBlobCompactTask.java
 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionBlobCompactTask.java
index 1069e7c62d..675aaf29b9 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionBlobCompactTask.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionBlobCompactTask.java
@@ -128,7 +128,7 @@ public class DataEvolutionBlobCompactTask extends 
DataEvolutionCompactTask {
             RowType blobWriteType,
             String blobFieldName,
             DataFilePathFactory pathFactory) {
-        BlobFileFormat blobFileFormat = new BlobFileFormat();
+        BlobFileFormat blobFileFormat = new BlobFileFormat(false, 
options.blobCopyBufferSize());
         return new RowDataFileWriter(
                 table.fileIO(),
                 RollingFileWriter.createFileWriterContext(
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
 
b/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
index 9216f7597c..edb9ed4a0b 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizer.java
@@ -20,6 +20,7 @@ package org.apache.paimon.blob;
 
 import org.apache.paimon.data.Blob;
 import org.apache.paimon.data.BlobDescriptor;
+import org.apache.paimon.data.BlobFetchMetricReporter;
 import org.apache.paimon.data.GenericArray;
 import org.apache.paimon.data.GenericRow;
 import org.apache.paimon.data.InternalArray;
@@ -61,7 +62,8 @@ public class PrimaryKeyBlobExternalizer {
             RowType valueType,
             Set<String> managedBlobFields,
             DataFilePathFactory pathFactory,
-            long targetFileSize) {
+            long targetFileSize,
+            int copyBufferSize) {
         checkArgument(targetFileSize > 0, "Managed BLOB target file size must 
be positive.");
         this.fileIO = fileIO;
         this.rowConverter = new RowDataToObjectArrayConverter(valueType);
@@ -97,7 +99,8 @@ public class PrimaryKeyBlobExternalizer {
                                     : new 
RowType(Collections.singletonList(field)),
                             pathFactory,
                             targetFileSize,
-                            uncommittedPacks));
+                            uncommittedPacks,
+                            copyBufferSize));
         }
         checkArgument(
                 unknownFields.isEmpty(),
@@ -223,6 +226,7 @@ public class PrimaryKeyBlobExternalizer {
         private final DataFilePathFactory pathFactory;
         private final long targetFileSize;
         private final List<Path> uncommittedPacks;
+        private final int copyBufferSize;
 
         private Path currentPath;
         private PositionOutputStream out;
@@ -234,12 +238,14 @@ public class PrimaryKeyBlobExternalizer {
                 RowType blobType,
                 DataFilePathFactory pathFactory,
                 long targetFileSize,
-                List<Path> uncommittedPacks) {
+                List<Path> uncommittedPacks,
+                int copyBufferSize) {
             this.fileIO = fileIO;
             this.blobType = blobType;
             this.pathFactory = pathFactory;
             this.targetFileSize = targetFileSize;
             this.uncommittedPacks = uncommittedPacks;
+            this.copyBufferSize = copyBufferSize;
         }
 
         private BlobDescriptor write(Blob blob) throws IOException {
@@ -271,7 +277,11 @@ public class PrimaryKeyBlobExternalizer {
                                 lastDescriptor = descriptor;
                                 return false;
                             },
-                            blobType);
+                            blobType,
+                            false,
+                            false,
+                            BlobFetchMetricReporter.NOOP,
+                            copyBufferSize);
             writer.setFile(currentPath);
         }
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java 
b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java
index 547d827a73..dfe875120d 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileWriterFactory.java
@@ -96,7 +96,8 @@ public class KeyValueFileWriterFactory {
                                 valueType,
                                 managedBlobFields,
                                 formatContext.pathFactory(new 
WriteFormatKey(0, false)),
-                                options.blobTargetFileSize());
+                                options.blobTargetFileSize(),
+                                options.blobCopyBufferSize());
     }
 
     public RowType keyType() {
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/BlobFileContext.java 
b/paimon-core/src/main/java/org/apache/paimon/operation/BlobFileContext.java
index 9d79ebe6eb..b1b50e157b 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/BlobFileContext.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/BlobFileContext.java
@@ -36,6 +36,7 @@ public class BlobFileContext {
     private final Set<String> blobInlineFields;
     private final boolean writeNullOnMissingFile;
     private final boolean writeNullOnFetchFailure;
+    private final int copyBufferSize;
 
     private @Nullable BlobConsumer blobConsumer;
     private BlobFetchMetricReporter blobFetchMetricReporter = 
BlobFetchMetricReporter.NOOP;
@@ -44,11 +45,13 @@ public class BlobFileContext {
             Set<String> blobDescriptorFields,
             Set<String> blobInlineFields,
             boolean writeNullOnMissingFile,
-            boolean writeNullOnFetchFailure) {
+            boolean writeNullOnFetchFailure,
+            int copyBufferSize) {
         this.blobDescriptorFields = blobDescriptorFields;
         this.blobInlineFields = blobInlineFields;
         this.writeNullOnMissingFile = writeNullOnMissingFile;
         this.writeNullOnFetchFailure = writeNullOnFetchFailure;
+        this.copyBufferSize = copyBufferSize;
     }
 
     @Nullable
@@ -72,7 +75,8 @@ public class BlobFileContext {
                 descriptorFields,
                 inlineFields,
                 options.blobWriteNullOnMissingFile(),
-                options.blobWriteNullOnFetchFailure());
+                options.blobWriteNullOnFetchFailure(),
+                options.blobCopyBufferSize());
     }
 
     public BlobFileContext withBlobConsumer(BlobConsumer blobConsumer) {
@@ -114,6 +118,10 @@ public class BlobFileContext {
         return writeNullOnFetchFailure;
     }
 
+    public int copyBufferSize() {
+        return copyBufferSize;
+    }
+
     public BlobFetchMetricReporter blobFetchMetricReporter() {
         return blobFetchMetricReporter;
     }
diff --git a/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java 
b/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java
index 3399990a8a..031dd29679 100644
--- a/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java
@@ -18,6 +18,7 @@
 
 package org.apache.paimon;
 
+import org.apache.paimon.options.MemorySize;
 import org.apache.paimon.options.Options;
 
 import org.junit.jupiter.api.Test;
@@ -157,4 +158,30 @@ public class CoreOptionsTest {
         assertThatThrownBy(() -> 
negativeMaxColumnsOptions.mapSharedShreddingMaxColumns("metrics"))
                 .hasMessageContaining("options 
map.shared-shredding.max-columns must > 0");
     }
+
+    @Test
+    public void testBlobCopyBufferSize() {
+        Options conf = new Options();
+        // default preserves the historical 4 KiB buffer.
+        assertThat(new CoreOptions(conf).blobCopyBufferSize()).isEqualTo(4 * 
1024);
+
+        conf.set(CoreOptions.BLOB_COPY_BUFFER_SIZE, MemorySize.parse("64 kb"));
+        assertThat(new CoreOptions(conf).blobCopyBufferSize()).isEqualTo(64 * 
1024);
+
+        // zero is rejected early with an option-named message.
+        conf.set(CoreOptions.BLOB_COPY_BUFFER_SIZE, MemorySize.parse("0 
bytes"));
+        assertThatThrownBy(() -> new CoreOptions(conf).blobCopyBufferSize())
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining("blob.copy-buffer-size");
+
+        // There is no arbitrary memory ceiling; only the Java int-sized array 
limit applies.
+        conf.set(CoreOptions.BLOB_COPY_BUFFER_SIZE, MemorySize.parse("512 
mb"));
+        assertThat(new CoreOptions(conf).blobCopyBufferSize()).isEqualTo(512 * 
1024 * 1024);
+        conf.set(CoreOptions.BLOB_COPY_BUFFER_SIZE, 
MemorySize.parse(Integer.MAX_VALUE + " bytes"));
+        assertThat(new 
CoreOptions(conf).blobCopyBufferSize()).isEqualTo(Integer.MAX_VALUE);
+        conf.set(CoreOptions.BLOB_COPY_BUFFER_SIZE, MemorySize.parse("2 gb"));
+        assertThatThrownBy(() -> new CoreOptions(conf).blobCopyBufferSize())
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining("blob.copy-buffer-size");
+    }
 }
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java 
b/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java
index c8c2a7f90c..b86b020cd1 100644
--- a/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/append/BlobUpdateTest.java
@@ -29,6 +29,7 @@ import org.apache.paimon.data.InternalRow;
 import org.apache.paimon.data.Timestamp;
 import org.apache.paimon.format.FormatWriter;
 import org.apache.paimon.format.blob.BlobFileFormat;
+import org.apache.paimon.format.blob.BlobFormatWriter;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.fs.PositionOutputStream;
@@ -473,7 +474,7 @@ public class BlobUpdateTest extends TableTestBase {
             throws IOException {
         try (PositionOutputStream out = fileIO.newOutputStream(path, false)) {
             FormatWriter writer =
-                    new BlobFileFormat()
+                    new BlobFileFormat(false, 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE)
                             .createWriterFactory(RowType.of(DataTypes.BLOB()))
                             .create(out, "none");
             for (Blob blob : blobs) {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
index 5fe4686c9f..da9fa84032 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/blob/PrimaryKeyBlobExternalizerTest.java
@@ -25,6 +25,7 @@ import org.apache.paimon.data.GenericArray;
 import org.apache.paimon.data.GenericRow;
 import org.apache.paimon.data.InternalArray;
 import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.format.blob.BlobFormatWriter;
 import org.apache.paimon.fs.Path;
 import org.apache.paimon.fs.PositionOutputStream;
 import org.apache.paimon.fs.PositionOutputStreamWrapper;
@@ -79,7 +80,7 @@ class PrimaryKeyBlobExternalizerTest {
                 new DataFilePathFactory(
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO,
                         RowType.of(DataTypes.BLOB()),
                         Collections.singleton("f0"),
@@ -104,7 +105,7 @@ class PrimaryKeyBlobExternalizerTest {
 
         assertThatThrownBy(
                         () ->
-                                new PrimaryKeyBlobExternalizer(
+                                newExternalizer(
                                         fileIO,
                                         RowType.of(DataTypes.INT()),
                                         Collections.singleton("f0"),
@@ -125,7 +126,7 @@ class PrimaryKeyBlobExternalizerTest {
 
         assertThatThrownBy(
                         () ->
-                                new PrimaryKeyBlobExternalizer(
+                                newExternalizer(
                                         fileIO,
                                         RowType.of(DataTypes.BLOB()),
                                         Collections.singleton("missing"),
@@ -145,8 +146,7 @@ class PrimaryKeyBlobExternalizerTest {
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         RowType valueType = RowType.of(DataTypes.INT(), DataTypes.BLOB());
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
-                        fileIO, valueType, Collections.singleton("f1"), 
pathFactory, 1024L);
+                newExternalizer(fileIO, valueType, 
Collections.singleton("f1"), pathFactory, 1024L);
         byte[] expected = "managed-blob".getBytes(StandardCharsets.UTF_8);
 
         InternalRow result =
@@ -172,7 +172,7 @@ class PrimaryKeyBlobExternalizerTest {
                 new DataFilePathFactory(
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO,
                         RowType.of(DataTypes.INT(), DataTypes.BLOB()),
                         Collections.singleton("f1"),
@@ -212,7 +212,7 @@ class PrimaryKeyBlobExternalizerTest {
                         },
                         new String[] {"id", "managed", "unmanaged"});
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO, valueType, Collections.singleton("managed"), 
pathFactory, 1024L);
         Blob unmanaged = Blob.fromData(new byte[] {2});
 
@@ -233,7 +233,7 @@ class PrimaryKeyBlobExternalizerTest {
                 new DataFilePathFactory(
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO,
                         RowType.of(DataTypes.INT(), DataTypes.BLOB()),
                         Collections.singleton("f1"),
@@ -268,7 +268,7 @@ class PrimaryKeyBlobExternalizerTest {
                 new DataFilePathFactory(
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO,
                         RowType.of(DataTypes.INT(), 
DataTypes.ARRAY(DataTypes.BLOB())),
                         Collections.singleton("f1"),
@@ -309,7 +309,7 @@ class PrimaryKeyBlobExternalizerTest {
                 new DataFilePathFactory(
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO,
                         RowType.of(DataTypes.INT(), 
DataTypes.ARRAY(DataTypes.BLOB())),
                         Collections.singleton("f1"),
@@ -352,7 +352,7 @@ class PrimaryKeyBlobExternalizerTest {
                 new DataFilePathFactory(
                         bucketPath, "avro", "data-", "changelog-", false, 
null, null);
         PrimaryKeyBlobExternalizer externalizer =
-                new PrimaryKeyBlobExternalizer(
+                newExternalizer(
                         fileIO,
                         RowType.of(DataTypes.INT(), 
DataTypes.ARRAY(DataTypes.BLOB())),
                         Collections.singleton("f1"),
@@ -380,4 +380,19 @@ class PrimaryKeyBlobExternalizerTest {
                 .hasMessageContaining("placeholder blob array");
         assertThat(fileIO.listStatus(bucketPath)).isEmpty();
     }
+
+    private static PrimaryKeyBlobExternalizer newExternalizer(
+            LocalFileIO fileIO,
+            RowType valueType,
+            java.util.Set<String> managedBlobFields,
+            DataFilePathFactory pathFactory,
+            long targetFileSize) {
+        return new PrimaryKeyBlobExternalizer(
+                fileIO,
+                valueType,
+                managedBlobFields,
+                pathFactory,
+                targetFileSize,
+                BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
+    }
 }
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormat.java 
b/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormat.java
index 5b0f359c7d..67d0baaa44 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormat.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormat.java
@@ -53,19 +53,17 @@ import static 
org.apache.paimon.utils.Preconditions.checkArgument;
 public class BlobFileFormat extends FileFormat {
 
     private final boolean blobAsDescriptor;
+    private final int copyBufferSize;
     private boolean writeNullOnMissingFile;
     private boolean writeNullOnFetchFailure;
     private BlobFetchMetricReporter blobFetchMetricReporter = 
BlobFetchMetricReporter.NOOP;
 
     @Nullable public BlobConsumer writeConsumer;
 
-    public BlobFileFormat() {
-        this(false);
-    }
-
-    public BlobFileFormat(boolean blobAsDescriptor) {
+    public BlobFileFormat(boolean blobAsDescriptor, int copyBufferSize) {
         super(BlobFileFormatFactory.IDENTIFIER);
         this.blobAsDescriptor = blobAsDescriptor;
+        this.copyBufferSize = copyBufferSize;
     }
 
     public static boolean isBlobFile(String fileName) {
@@ -132,7 +130,8 @@ public class BlobFileFormat extends FileFormat {
                     type,
                     writeNullOnMissingFile,
                     writeNullOnFetchFailure,
-                    blobFetchMetricReporter);
+                    blobFetchMetricReporter,
+                    copyBufferSize);
         }
     }
 
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormatFactory.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormatFactory.java
index 2a54d49709..625203bfbd 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormatFactory.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFileFormatFactory.java
@@ -35,6 +35,9 @@ public class BlobFileFormatFactory implements 
FileFormatFactory {
     @Override
     public FileFormat create(FormatContext formatContext) {
         boolean blobAsDescriptor = 
formatContext.options().get(CoreOptions.BLOB_AS_DESCRIPTOR);
-        return new BlobFileFormat(blobAsDescriptor);
+        int copyBufferSize =
+                CoreOptions.checkedBlobCopyBufferSize(
+                        
formatContext.options().get(CoreOptions.BLOB_COPY_BUFFER_SIZE).getBytes());
+        return new BlobFileFormat(blobAsDescriptor, copyBufferSize);
     }
 }
diff --git 
a/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFormatWriter.java
 
b/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFormatWriter.java
index 9de5a169d1..5e2bc9f9bb 100644
--- 
a/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFormatWriter.java
+++ 
b/paimon-format/src/main/java/org/apache/paimon/format/blob/BlobFormatWriter.java
@@ -64,6 +64,7 @@ public class BlobFormatWriter implements 
FileAwareFormatWriter {
     public static final long NULL_LENGTH = -1L;
     public static final long PLACE_HOLDER_LENGTH = -2L;
     public static final long ARRAY_NULL_ELEMENT_LENGTH = -1L;
+    public static final int DEFAULT_COPY_BUFFER_SIZE = 4 * 1024;
 
     private final PositionOutputStream out;
     @Nullable private final BlobConsumer writeConsumer;
@@ -78,41 +79,18 @@ public class BlobFormatWriter implements 
FileAwareFormatWriter {
 
     private String pathString;
 
-    public BlobFormatWriter(
-            PositionOutputStream out, @Nullable BlobConsumer writeConsumer, 
RowType type) {
-        this(out, writeConsumer, type, false, false);
-    }
-
-    public BlobFormatWriter(
-            PositionOutputStream out,
-            @Nullable BlobConsumer writeConsumer,
-            RowType type,
-            boolean writeNullOnMissingFile) {
-        this(out, writeConsumer, type, writeNullOnMissingFile, false);
-    }
-
-    public BlobFormatWriter(
-            PositionOutputStream out,
-            @Nullable BlobConsumer writeConsumer,
-            RowType type,
-            boolean writeNullOnMissingFile,
-            boolean writeNullOnFetchFailure) {
-        this(
-                out,
-                writeConsumer,
-                type,
-                writeNullOnMissingFile,
-                writeNullOnFetchFailure,
-                BlobFetchMetricReporter.NOOP);
-    }
-
     public BlobFormatWriter(
             PositionOutputStream out,
             @Nullable BlobConsumer writeConsumer,
             RowType type,
             boolean writeNullOnMissingFile,
             boolean writeNullOnFetchFailure,
-            BlobFetchMetricReporter blobFetchMetricReporter) {
+            BlobFetchMetricReporter blobFetchMetricReporter,
+            int copyBufferSize) {
+        checkArgument(
+                copyBufferSize > 0,
+                "BLOB copy buffer size must be positive, but was %s.",
+                copyBufferSize);
         this.out = out;
         this.writeConsumer = writeConsumer;
         this.blobFetchMetricReporter = blobFetchMetricReporter;
@@ -125,7 +103,7 @@ public class BlobFormatWriter implements 
FileAwareFormatWriter {
                         ? new ArrayBlobElementWriter()
                         : new RawBlobElementWriter();
         this.crc32 = new CRC32();
-        this.tmpBuffer = new byte[4096];
+        this.tmpBuffer = new byte[copyBufferSize];
         this.lengths = new LongArrayList(16);
     }
 
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
index 2b72cbb08e..ff5315e668 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFileFormatTest.java
@@ -95,7 +95,8 @@ public class BlobFileFormatTest {
 
     @Test
     public void testWriteArrayBlobPlaceholderWithProjectedRow() throws 
IOException {
-        BlobFileFormat format = new BlobFileFormat();
+        BlobFileFormat format =
+                new BlobFileFormat(false, 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
         RowType rowType = RowType.of(DataTypes.ARRAY(DataTypes.BLOB()));
 
         try (PositionOutputStream out = fileIO.newOutputStream(file, false)) {
@@ -162,7 +163,8 @@ public class BlobFileFormatTest {
     }
 
     private void innerTest(boolean blobAsDescriptor) throws IOException {
-        BlobFileFormat format = new BlobFileFormat(blobAsDescriptor);
+        BlobFileFormat format =
+                new BlobFileFormat(blobAsDescriptor, 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
         RowType rowType = RowType.of(DataTypes.BLOB());
 
         // write
@@ -235,7 +237,8 @@ public class BlobFileFormatTest {
 
     private void assertMalformedArrayPayload(
             ArrayPayloadCorruptor corruptor, String expectedMessage) throws 
IOException {
-        BlobFileFormat format = new BlobFileFormat();
+        BlobFileFormat format =
+                new BlobFileFormat(false, 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
         RowType rowType = RowType.of(DataTypes.ARRAY(DataTypes.BLOB()));
         try (PositionOutputStream out = fileIO.newOutputStream(file, false)) {
             FormatWriter writer = 
format.createWriterFactory(rowType).create(out, null);
@@ -272,7 +275,8 @@ public class BlobFileFormatTest {
 
     @Test
     public void testReadWithProjectedRowTypeContainingExtraFields() throws 
IOException {
-        BlobFileFormat format = new BlobFileFormat(false);
+        BlobFileFormat format =
+                new BlobFileFormat(false, 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
         RowType writeRowType = RowType.of(DataTypes.BLOB());
 
         // write blob data
@@ -310,7 +314,8 @@ public class BlobFileFormatTest {
     }
 
     private void innerTestArray(boolean blobAsDescriptor) throws IOException {
-        BlobFileFormat format = new BlobFileFormat(blobAsDescriptor);
+        BlobFileFormat format =
+                new BlobFileFormat(blobAsDescriptor, 
BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
         RowType rowType = RowType.of(DataTypes.ARRAY(DataTypes.BLOB()));
 
         GenericArray first =
diff --git 
a/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFormatWriterTest.java
 
b/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFormatWriterTest.java
index 85e4314589..633d32e441 100644
--- 
a/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFormatWriterTest.java
+++ 
b/paimon-format/src/test/java/org/apache/paimon/format/blob/BlobFormatWriterTest.java
@@ -43,6 +43,10 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 
 import java.nio.file.Files;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -58,11 +62,7 @@ public class BlobFormatWriterTest {
         byte[] firstPayload = "first-blob".getBytes();
         byte[] secondPayload = "second-blob-payload".getBytes();
 
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new 
LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
-                        null,
-                        rowType);
+        BlobFormatWriter writer = newWriter(outputFile, rowType);
         writer.addElement(GenericRow.of(Blob.fromData(firstPayload)));
         writer.addElement(GenericRow.of(Blob.fromData(secondPayload)));
         writer.close();
@@ -88,13 +88,7 @@ public class BlobFormatWriterTest {
             @TempDir java.nio.file.Path tempDir) throws Exception {
         RowType rowType = RowType.of(DataTypes.BLOB());
         java.nio.file.Path outputFile = tempDir.resolve("blob.out");
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new 
LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
-                        null,
-                        rowType,
-                        false,
-                        true);
+        BlobFormatWriter writer = newWriter(outputFile, rowType, false, true);
 
         writer.addElement(
                 GenericRow.of(
@@ -112,13 +106,7 @@ public class BlobFormatWriterTest {
             @TempDir java.nio.file.Path tempDir) throws Exception {
         RowType rowType = RowType.of(DataTypes.BLOB());
         java.nio.file.Path outputFile = tempDir.resolve("blob.out");
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new 
LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
-                        null,
-                        rowType,
-                        false,
-                        true);
+        BlobFormatWriter writer = newWriter(outputFile, rowType, false, true);
 
         writer.addElement(
                 GenericRow.of(
@@ -135,14 +123,7 @@ public class BlobFormatWriterTest {
     public void testHttpRateLimitFailsWhenFetchFailureDisabled(@TempDir 
java.nio.file.Path tempDir)
             throws Exception {
         RowType rowType = RowType.of(DataTypes.BLOB());
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        false,
-                        false);
+        BlobFormatWriter writer = newWriter(tempDir.resolve("blob.out"), 
rowType, false, false);
 
         assertThatThrownBy(
                         () ->
@@ -162,14 +143,7 @@ public class BlobFormatWriterTest {
     public void testHttpNotFoundPropagatesWhenFetchFailureDisabled(
             @TempDir java.nio.file.Path tempDir) throws Exception {
         RowType rowType = RowType.of(DataTypes.BLOB());
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        false,
-                        false);
+        BlobFormatWriter writer = newWriter(tempDir.resolve("blob.out"), 
rowType, false, false);
 
         assertThatThrownBy(
                         () ->
@@ -196,13 +170,7 @@ public class BlobFormatWriterTest {
                 new 
BlobDescriptor("https://img.alicdn.com/imgextra/##1304008055350781673";, 0, -1)
                         .serialize();
 
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new 
LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
-                        null,
-                        rowType,
-                        false,
-                        true);
+        BlobFormatWriter writer = newWriter(outputFile, rowType, false, true);
 
         writer.addElement(new DescriptorBytesRow(descriptorBytes, 
uriReaderFactory));
         writer.close();
@@ -228,13 +196,7 @@ public class BlobFormatWriterTest {
                 new 
BlobDescriptor("https://img.alicdn.com/imgextra/##1304008055350781673";, 0, -1)
                         .serialize();
 
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new 
LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
-                        null,
-                        rowType,
-                        false,
-                        true);
+        BlobFormatWriter writer = newWriter(outputFile, rowType, false, true);
 
         writer.addElement(
                 GenericRow.of(new DescriptorBytesArray(descriptorBytes, 
uriReaderFactory)));
@@ -269,13 +231,7 @@ public class BlobFormatWriterTest {
             @TempDir java.nio.file.Path tempDir) throws Exception {
         RowType rowType = RowType.of(DataTypes.BLOB());
         java.nio.file.Path outputFile = tempDir.resolve("blob.out");
-        BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new 
LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
-                        null,
-                        rowType,
-                        true,
-                        false);
+        BlobFormatWriter writer = newWriter(outputFile, rowType, true, false);
 
         writer.addElement(
                 GenericRow.of(
@@ -300,14 +256,7 @@ public class BlobFormatWriterTest {
         RowType rowType = RowType.of(DataTypes.BLOB());
         TestingBlobFetchMetricReporter metricReporter = new 
TestingBlobFetchMetricReporter();
         BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        false,
-                        true,
-                        metricReporter);
+                newWriter(tempDir.resolve("blob.out"), rowType, false, true, 
metricReporter);
 
         writer.addElement(GenericRow.of(Blob.fromData("image".getBytes())));
         writer.addElement(
@@ -329,14 +278,7 @@ public class BlobFormatWriterTest {
         RowType rowType = RowType.of(DataTypes.BLOB());
         TestingBlobFetchMetricReporter metricReporter = new 
TestingBlobFetchMetricReporter();
         BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        false,
-                        false,
-                        metricReporter);
+                newWriter(tempDir.resolve("blob.out"), rowType, false, false, 
metricReporter);
 
         assertThatThrownBy(
                         () ->
@@ -364,14 +306,7 @@ public class BlobFormatWriterTest {
                 new BlobDescriptor("https://example.com/missing.jpg";, 0, 
-1).serialize();
         TestingBlobFetchMetricReporter metricReporter = new 
TestingBlobFetchMetricReporter();
         BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        true,
-                        false,
-                        metricReporter);
+                newWriter(tempDir.resolve("blob.out"), rowType, true, false, 
metricReporter);
 
         writer.addElement(new DescriptorBytesRow(descriptorBytes, 
uriReaderFactory, true));
         writer.close();
@@ -386,14 +321,7 @@ public class BlobFormatWriterTest {
         RowType rowType = RowType.of(DataTypes.BLOB());
         TestingBlobFetchMetricReporter metricReporter = new 
TestingBlobFetchMetricReporter();
         BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        true,
-                        false,
-                        metricReporter);
+                newWriter(tempDir.resolve("blob.out"), rowType, true, false, 
metricReporter);
 
         writer.addElement(GenericRow.of((Object) null));
         writer.close();
@@ -408,14 +336,7 @@ public class BlobFormatWriterTest {
         RowType rowType = RowType.of(DataTypes.ARRAY(DataTypes.BLOB()));
         TestingBlobFetchMetricReporter metricReporter = new 
TestingBlobFetchMetricReporter();
         BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        true,
-                        true,
-                        metricReporter);
+                newWriter(tempDir.resolve("blob.out"), rowType, true, true, 
metricReporter);
 
         writer.addElement(
                 GenericRow.of(
@@ -449,14 +370,7 @@ public class BlobFormatWriterTest {
         RowType rowType = RowType.of(DataTypes.ARRAY(DataTypes.BLOB()));
         TestingBlobFetchMetricReporter metricReporter = new 
TestingBlobFetchMetricReporter();
         BlobFormatWriter writer =
-                new BlobFormatWriter(
-                        new LocalFileIO.LocalPositionOutputStream(
-                                tempDir.resolve("blob.out").toFile()),
-                        null,
-                        rowType,
-                        false,
-                        false,
-                        metricReporter);
+                newWriter(tempDir.resolve("blob.out"), rowType, false, false, 
metricReporter);
 
         assertThatThrownBy(
                         () ->
@@ -477,6 +391,222 @@ public class BlobFormatWriterTest {
         assertThat(metricReporter.fetchFailureNullWritten).isEqualTo(0);
     }
 
+    @Test
+    public void testCopyBufferSizeIsRespectedForBlobRef(@TempDir 
java.nio.file.Path tempDir)
+            throws Exception {
+        String uri = "mem://file";
+        byte[] source = sequentialBytes(20);
+        RecordingUriReader reader = new RecordingUriReader(singleFile(uri, 
source));
+        java.nio.file.Path outputFile = tempDir.resolve("blob.out");
+
+        BlobFormatWriter writer = newWriter(outputFile, 
RowType.of(DataTypes.BLOB()), 8);
+        writer.addElement(GenericRow.of(new BlobRef(reader, new 
BlobDescriptor(uri, 0, 20))));
+        writer.close();
+
+        // With an 8-byte copy buffer, no single read request exceeds 8 bytes.
+        assertThat(reader.opened).hasSize(1);
+        assertThat(reader.opened.get(0).maxReadRequest).isEqualTo(8);
+        assertThat(readBackBlobs(outputFile, 1)).containsExactly(source);
+    }
+
+    @Test
+    public void testDefaultCopyBufferSize(@TempDir java.nio.file.Path tempDir) 
throws Exception {
+        // The configured default preserves the historical 4 KiB copy buffer.
+        assertThat(BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE).isEqualTo(4 * 
1024);
+
+        String uri = "mem://file";
+        byte[] source = sequentialBytes(5000);
+        RecordingUriReader reader = new RecordingUriReader(singleFile(uri, 
source));
+        java.nio.file.Path outputFile = tempDir.resolve("blob.out");
+
+        BlobFormatWriter writer = newWriter(outputFile, 
RowType.of(DataTypes.BLOB()));
+        writer.addElement(GenericRow.of(new BlobRef(reader, new 
BlobDescriptor(uri, 0, 5000))));
+        writer.close();
+
+        assertThat(reader.opened.get(0).maxReadRequest).isEqualTo(4 * 1024);
+        assertThat(readBackBlobs(outputFile, 1)).containsExactly(source);
+    }
+
+    private static BlobFormatWriter newWriter(java.nio.file.Path outputFile, 
RowType rowType)
+            throws java.io.FileNotFoundException {
+        return newWriter(
+                outputFile,
+                rowType,
+                false,
+                false,
+                BlobFetchMetricReporter.NOOP,
+                BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
+    }
+
+    private static BlobFormatWriter newWriter(
+            java.nio.file.Path outputFile,
+            RowType rowType,
+            boolean writeNullOnMissingFile,
+            boolean writeNullOnFetchFailure)
+            throws java.io.FileNotFoundException {
+        return newWriter(
+                outputFile,
+                rowType,
+                writeNullOnMissingFile,
+                writeNullOnFetchFailure,
+                BlobFetchMetricReporter.NOOP,
+                BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
+    }
+
+    private static BlobFormatWriter newWriter(
+            java.nio.file.Path outputFile,
+            RowType rowType,
+            boolean writeNullOnMissingFile,
+            boolean writeNullOnFetchFailure,
+            BlobFetchMetricReporter blobFetchMetricReporter)
+            throws java.io.FileNotFoundException {
+        return newWriter(
+                outputFile,
+                rowType,
+                writeNullOnMissingFile,
+                writeNullOnFetchFailure,
+                blobFetchMetricReporter,
+                BlobFormatWriter.DEFAULT_COPY_BUFFER_SIZE);
+    }
+
+    private static BlobFormatWriter newWriter(
+            java.nio.file.Path outputFile, RowType rowType, int copyBufferSize)
+            throws java.io.FileNotFoundException {
+        return newWriter(
+                outputFile, rowType, false, false, 
BlobFetchMetricReporter.NOOP, copyBufferSize);
+    }
+
+    private static BlobFormatWriter newWriter(
+            java.nio.file.Path outputFile,
+            RowType rowType,
+            boolean writeNullOnMissingFile,
+            boolean writeNullOnFetchFailure,
+            BlobFetchMetricReporter blobFetchMetricReporter,
+            int copyBufferSize)
+            throws java.io.FileNotFoundException {
+        return new BlobFormatWriter(
+                new LocalFileIO.LocalPositionOutputStream(outputFile.toFile()),
+                null,
+                rowType,
+                writeNullOnMissingFile,
+                writeNullOnFetchFailure,
+                blobFetchMetricReporter,
+                copyBufferSize);
+    }
+
+    private static List<byte[]> readBackBlobs(java.nio.file.Path outputFile, 
int expectedCount)
+            throws Exception {
+        LocalFileIO fileIO = new LocalFileIO();
+        Path filePath = new Path(outputFile.toUri());
+        long fileSize = Files.size(outputFile);
+        List<byte[]> result = new ArrayList<>();
+        try (SeekableInputStream in = fileIO.newInputStream(filePath)) {
+            BlobFileMeta fileMeta = new BlobFileMeta(in, fileSize, null);
+            assertThat(fileMeta.recordNumber()).isEqualTo(expectedCount);
+            BlobFormatReader reader =
+                    new BlobFormatReader(
+                            fileIO, filePath, fileMeta, in, 1, 0, 
DataTypes.BLOB(), false);
+            FileRecordIterator<InternalRow> iterator = reader.readBatch();
+            for (int i = 0; i < expectedCount; i++) {
+                InternalRow row = iterator.next();
+                assertThat(row).isNotNull();
+                result.add(readAll(row.getBlob(0)));
+            }
+        }
+        return result;
+    }
+
+    private static byte[] readAll(Blob blob) throws Exception {
+        try (SeekableInputStream in = blob.newInputStream()) {
+            return org.apache.paimon.utils.IOUtils.readFully(in, false);
+        }
+    }
+
+    private static byte[] sequentialBytes(int length) {
+        byte[] bytes = new byte[length];
+        for (int i = 0; i < length; i++) {
+            bytes[i] = (byte) i;
+        }
+        return bytes;
+    }
+
+    private static Map<String, byte[]> singleFile(String uri, byte[] data) {
+        Map<String, byte[]> files = new LinkedHashMap<>();
+        files.put(uri, data);
+        return files;
+    }
+
+    /** A {@link UriReader} over in-memory files that records opened streams. 
*/
+    private static final class RecordingUriReader implements UriReader {
+
+        private final Map<String, byte[]> files;
+        private final List<CountingSeekableInputStream> opened = new 
ArrayList<>();
+
+        private RecordingUriReader(Map<String, byte[]> files) {
+            this.files = files;
+        }
+
+        @Override
+        public SeekableInputStream newInputStream(String uri) {
+            byte[] data = files.get(uri);
+            if (data == null) {
+                throw new IllegalArgumentException("Unknown uri: " + uri);
+            }
+            CountingSeekableInputStream stream = new 
CountingSeekableInputStream(data);
+            opened.add(stream);
+            return stream;
+        }
+    }
+
+    /** A seekable stream over a byte array that records close count and max 
read request size. */
+    private static final class CountingSeekableInputStream extends 
SeekableInputStream {
+
+        private final byte[] data;
+        private int pos;
+        private int maxReadRequest;
+
+        private CountingSeekableInputStream(byte[] data) {
+            this.data = data;
+        }
+
+        @Override
+        public void seek(long desired) {
+            this.pos = (int) desired;
+        }
+
+        @Override
+        public long getPos() {
+            return pos;
+        }
+
+        @Override
+        public int read() {
+            maxReadRequest = Math.max(maxReadRequest, 1);
+            if (pos >= data.length) {
+                return -1;
+            }
+            return data[pos++] & 0xFF;
+        }
+
+        @Override
+        public int read(byte[] b, int off, int len) {
+            if (len == 0) {
+                return 0;
+            }
+            maxReadRequest = Math.max(maxReadRequest, len);
+            if (pos >= data.length) {
+                return -1;
+            }
+            int n = Math.min(len, data.length - pos);
+            System.arraycopy(data, pos, b, off, n);
+            pos += n;
+            return n;
+        }
+
+        @Override
+        public void close() {}
+    }
+
     private static void assertBlobPayload(Blob blob, byte[] expected) throws 
Exception {
         try (SeekableInputStream blobIn = blob.newInputStream()) {
             byte[] actual = new byte[expected.length];
diff --git a/paimon-python/pypaimon/common/options/core_options.py 
b/paimon-python/pypaimon/common/options/core_options.py
index 353c140ff5..a6aa5f96b5 100644
--- a/paimon-python/pypaimon/common/options/core_options.py
+++ b/paimon-python/pypaimon/common/options/core_options.py
@@ -395,6 +395,13 @@ class CoreOptions:
         .with_description("The target file size for blob files.")
     )
 
+    BLOB_COPY_BUFFER_SIZE: ConfigOption[MemorySize] = (
+        ConfigOptions.key("blob.copy-buffer-size")
+        .memory_type()
+        .default_value(MemorySize.of_kibi_bytes(4))
+        .with_description("Buffer size used when copying BLOB payloads into 
BLOB files.")
+    )
+
     VECTOR_FILE_FORMAT: ConfigOption[str] = (
         ConfigOptions.key("vector.file.format")
         .string_type()
@@ -1120,6 +1127,16 @@ class CoreOptions:
         else:
             return self.target_file_size(has_primary_key=False)
 
+    def blob_copy_buffer_size(self):
+        size = self.options.get(CoreOptions.BLOB_COPY_BUFFER_SIZE, 
None).get_bytes()
+        # Java BlobFormatWriter stores the byte-array size in an int.
+        max_size = (1 << 31) - 1
+        if not 1 <= size <= max_size:
+            raise ValueError(
+                f"'{CoreOptions.BLOB_COPY_BUFFER_SIZE.key()}' must be between 
1 byte and "
+                f"{max_size} bytes, but was {size} bytes.")
+        return size
+
     def vector_file_format(self, default=None):
         return self.options.get(CoreOptions.VECTOR_FILE_FORMAT, default)
 
diff --git a/paimon-python/pypaimon/tests/blob_test.py 
b/paimon-python/pypaimon/tests/blob_test.py
index 89373495af..efa317fa71 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -130,6 +130,36 @@ class BlobTest(unittest.TestCase):
         except OSError:
             pass  # Ignore cleanup errors
 
+    def test_blob_copy_buffer_size_validation(self):
+        """blob.copy-buffer-size must be positive and fit Java's int-sized 
array."""
+        from pypaimon.write.blob_format_writer import BlobFormatWriter
+        from pypaimon.common.options.core_options import CoreOptions
+
+        # The internal writer already receives an int-sized value and only 
rejects non-positive
+        # buffers.
+        for bad in [0, -1]:
+            with self.assertRaises(ValueError):
+                BlobFormatWriter(io.BytesIO(), copy_buffer_size=bad)
+        BlobFormatWriter(io.BytesIO(), copy_buffer_size=8).close()
+
+        # CoreOptions defaults to 4 KiB, accepts values above the old 256 MiB 
ceiling,
+        # and retains only the Java int technical limit for cross-language 
consistency.
+        self.assertEqual(CoreOptions(Options({})).blob_copy_buffer_size(), 
4096)
+        self.assertEqual(
+            CoreOptions(Options({'blob.copy-buffer-size': '256 
kb'})).blob_copy_buffer_size(),
+            256 * 1024)
+        self.assertEqual(
+            CoreOptions(Options({'blob.copy-buffer-size': '512 
mb'})).blob_copy_buffer_size(),
+            512 * 1024 * 1024)
+        self.assertEqual(
+            CoreOptions(Options({
+                'blob.copy-buffer-size': f'{(1 << 31) - 1} bytes'
+            })).blob_copy_buffer_size(),
+            (1 << 31) - 1)
+        for bad in ['0 bytes', '2 gb', '3 gb']:
+            with self.assertRaises(ValueError):
+                CoreOptions(Options({'blob.copy-buffer-size': 
bad})).blob_copy_buffer_size()
+
     def test_from_data(self):
         """Test Blob.from_data() method."""
         test_data = b"test data"
diff --git a/paimon-python/pypaimon/write/blob_format_writer.py 
b/paimon-python/pypaimon/write/blob_format_writer.py
index 36c1044443..6efdf9e49f 100644
--- a/paimon-python/pypaimon/write/blob_format_writer.py
+++ b/paimon-python/pypaimon/write/blob_format_writer.py
@@ -41,10 +41,15 @@ class BlobFormatWriter:
 
     def __init__(self, output_stream: BinaryIO,
                  blob_consumer: Optional[BlobConsumer] = None,
-                 file_path: Optional[str] = None):
+                 file_path: Optional[str] = None,
+                 copy_buffer_size: int = BUFFER_SIZE):
+        if copy_buffer_size <= 0:
+            raise ValueError(
+                f"BLOB copy buffer size must be positive, but was 
{copy_buffer_size}.")
         self.output_stream = output_stream
         self._blob_consumer = blob_consumer
         self._file_path = file_path
+        self.copy_buffer_size = copy_buffer_size
         self.lengths: List[int] = []
         self.position = 0
 
@@ -176,10 +181,10 @@ class BlobFormatWriter:
         else:
             stream = blob_value.new_input_stream()
             try:
-                chunk = stream.read(self.BUFFER_SIZE)
+                chunk = stream.read(self.copy_buffer_size)
                 while chunk:
                     crc32 = self._write_with_crc(chunk, crc32)
-                    chunk = stream.read(self.BUFFER_SIZE)
+                    chunk = stream.read(self.copy_buffer_size)
             finally:
                 stream.close()
 
diff --git a/paimon-python/pypaimon/write/writer/blob_file_writer.py 
b/paimon-python/pypaimon/write/writer/blob_file_writer.py
index 728d86c50f..b366dab96e 100644
--- a/paimon-python/pypaimon/write/writer/blob_file_writer.py
+++ b/paimon-python/pypaimon/write/writer/blob_file_writer.py
@@ -36,7 +36,8 @@ class BlobFileWriter:
     Writes rows one by one and tracks file size.
     """
 
-    def __init__(self, file_io, file_path: Path, blob_consumer: 
Optional[BlobConsumer] = None):
+    def __init__(self, file_io, file_path: Path, blob_consumer: 
Optional[BlobConsumer] = None,
+                 copy_buffer_size: int = BlobFormatWriter.BUFFER_SIZE):
         self.file_io = file_io
         self.file_path = file_path
         self._blob_consumer = blob_consumer
@@ -45,6 +46,7 @@ class BlobFileWriter:
             self.output_stream,
             blob_consumer=blob_consumer,
             file_path=str(file_path),
+            copy_buffer_size=copy_buffer_size,
         )
         self.row_count = 0
         self.closed = False
diff --git a/paimon-python/pypaimon/write/writer/blob_writer.py 
b/paimon-python/pypaimon/write/writer/blob_writer.py
index 2e0e682130..c5bf143728 100644
--- a/paimon-python/pypaimon/write/writer/blob_writer.py
+++ b/paimon-python/pypaimon/write/writer/blob_writer.py
@@ -44,6 +44,8 @@ class BlobWriter(AppendOnlyDataWriter):
 
         options = self.table.options
         self.blob_target_file_size = CoreOptions.blob_target_file_size(options)
+        # Use the effective options (constructor-provided), consistent with 
the other accessors.
+        self.blob_copy_buffer_size = self.options.blob_copy_buffer_size()
 
         self._blob_consumer = blob_consumer
         self.current_writer: Optional[BlobFileWriter] = None
@@ -100,7 +102,12 @@ class BlobWriter(AppendOnlyDataWriter):
         self.file_count += 1  # Increment counter for next file
         file_path = self._generate_file_path(file_name)
         self.current_file_path = file_path
-        self.current_writer = BlobFileWriter(self.file_io, file_path, 
blob_consumer=self._blob_consumer)
+        self.current_writer = BlobFileWriter(
+            self.file_io,
+            file_path,
+            blob_consumer=self._blob_consumer,
+            copy_buffer_size=self.blob_copy_buffer_size,
+        )
 
     def rolling_file(self) -> bool:
         if self.current_writer is None:

Reply via email to