diff --git a/docs/generated/core_configuration.html b/docs/generated/core_configuration.html index 143c244451a9..49500e1b61fc 100644 --- a/docs/generated/core_configuration.html +++ b/docs/generated/core_configuration.html @@ -104,6 +104,12 @@ Boolean 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. + +
blob.copy-buffer-size
+ 4 kb + MemorySize + Buffer size used when copying BLOB payloads into BLOB files. +
blob.split-by-file-size
(none) 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 4f1cb3111193..04ef7cbcb4b3 100644 --- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java +++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java @@ -830,6 +830,14 @@ public InlineElement getDescription() { "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 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 NUM_SORTED_RUNS_COMPACTION_TRIGGER = key("num-sorted-run.compaction-trigger") .intType() @@ -3286,6 +3294,21 @@ public long blobTargetFileSize() { .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 86df6b10b265..9b8897862a8b 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 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 d952018ee1f0..39d6c6eac833 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 MultipleBlobFileWriter( Set 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 1069e7c62d86..675aaf29b986 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 @@ private FileWriter createBlobFileWriter( 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 9216f7597c3e..edb9ed4a0b38 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 @@ 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 PrimaryKeyBlobExternalizer( RowType valueType, Set 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 PrimaryKeyBlobExternalizer( : new RowType(Collections.singletonList(field)), pathFactory, targetFileSize, - uncommittedPacks)); + uncommittedPacks, + copyBufferSize)); } checkArgument( unknownFields.isEmpty(), @@ -223,6 +226,7 @@ private static class ManagedBlobPackWriter { private final DataFilePathFactory pathFactory; private final long targetFileSize; private final List uncommittedPacks; + private final int copyBufferSize; private Path currentPath; private PositionOutputStream out; @@ -234,12 +238,14 @@ private ManagedBlobPackWriter( RowType blobType, DataFilePathFactory pathFactory, long targetFileSize, - List uncommittedPacks) { + List 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 @@ private void openCurrent() throws IOException { 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 547d827a73b4..dfe875120dcf 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 @@ private 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 9d79ebe6eb6c..b1b50e157b82 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 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 @@ private BlobFileContext( Set blobDescriptorFields, Set 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 static BlobFileContext create(RowType rowType, CoreOptions options) { descriptorFields, inlineFields, options.blobWriteNullOnMissingFile(), - options.blobWriteNullOnFetchFailure()); + options.blobWriteNullOnFetchFailure(), + options.blobCopyBufferSize()); } public BlobFileContext withBlobConsumer(BlobConsumer blobConsumer) { @@ -114,6 +118,10 @@ public boolean writeNullOnFetchFailure() { 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 3399990a8a4b..031dd296792f 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 void testMapStorageLayout() { 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 c8c2a7f90c4b..b86b020cd199 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.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 @@ private static DataFileMeta writeBlobFile( 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 5fe4686c9f75..da9fa84032e8 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.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 @@ public void close() throws IOException { 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 @@ void testRejectsNonBlobManagedField() throws Exception { assertThatThrownBy( () -> - new PrimaryKeyBlobExternalizer( + newExternalizer( fileIO, RowType.of(DataTypes.INT()), Collections.singleton("f0"), @@ -125,7 +126,7 @@ void testRejectsUnknownManagedField() throws Exception { assertThatThrownBy( () -> - new PrimaryKeyBlobExternalizer( + newExternalizer( fileIO, RowType.of(DataTypes.BLOB()), Collections.singleton("missing"), @@ -145,8 +146,7 @@ void testExternalizeRawBlobBeforeBuffering() throws Exception { 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 @@ void testRematerializesBlobRefInput() throws Exception { 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 @@ void testExternalizesOnlyDeclaredFields() throws Exception { }, 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 @@ void testRetractSkipsPayloadAndAbortDeletesPrivatePack() throws Exception { 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 @@ void testExternalizeBlobArrayElements() throws Exception { 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 @@ void testRematerializesBlobRefArrayElement() throws Exception { 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 @@ void testBlobArrayNullEmptyRetractAndPlaceholder() throws Exception { 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 @@ void testBlobArrayNullEmptyRetractAndPlaceholder() throws Exception { .hasMessageContaining("placeholder blob array"); assertThat(fileIO.listStatus(bucketPath)).isEmpty(); } + + private static PrimaryKeyBlobExternalizer newExternalizer( + LocalFileIO fileIO, + RowType valueType, + java.util.Set 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 5b0f359c7df1..67d0baaa4471 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 @@ 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 FormatWriter create(PositionOutputStream out, String compression) { 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 2a54d497093b..625203bfbdde 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 String identifier() { @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 9de5a169d1e2..5e2bc9f9bb24 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 BlobFormatWriter( ? 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 2b72cbb08e9f..ff5315e6684f 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 void testReadArrayBlobInlineBytes() throws IOException { @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 void testRejectArrayBlobElementLengthMismatch() throws IOException { } 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 @@ private void innerTest(boolean blobAsDescriptor) throws IOException { 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 @@ private void assertMalformedArrayPayload( @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 void testReadWithProjectedRowTypeContainingExtraFields() throws IOExcepti } 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 85e431458953..633d32e44126 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.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 void testTwoConsecutiveBlobsPreserveReadback(@TempDir java.nio.file.Path 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 void testWriteNullOnFetchFailureFallbackForHttpBadRequest( @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 void testHttpRateLimitWritesNullWhenFetchFailureEnabled( @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 void testHttpRateLimitWritesNullWhenFetchFailureEnabled( 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 void testHttpRateLimitFailsWhenFetchFailureDisabled(@TempDir java.nio.fil 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 void testWriteNullOnFetchFailureForInvalidUriDescriptor( 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 void testArrayWriteNullOnFetchFailureForInvalidUriDescriptor( 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 void testHttpNotFoundWritesNullWhenMissingFileEnabled( @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 void testBlobFetchMetricReporterForSuccessAndNullWritten( 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 void testBlobFetchMetricReporterForUnhandledFailure(@TempDir java.nio.fil 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 void testBlobFetchMetricReporterForPreCheckedMissingFile( 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 void testBlobFetchMetricReporterIgnoresUserNull(@TempDir java.nio.file.Pa 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 void testArrayBlobFetchMetricReporter(@TempDir java.nio.file.Path tempDir 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 void testArrayBlobFetchMetricReporterForUnhandledFailure( 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 void testArrayBlobFetchMetricReporterForUnhandledFailure( 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 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 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 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 singleFile(String uri, byte[] data) { + Map 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 files; + private final List opened = new ArrayList<>(); + + private RecordingUriReader(Map 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 4d2e6a7bb8bc..5af7a88a8f71 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() @@ -1111,6 +1118,16 @@ def blob_target_file_size(self, default=None): 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 9da03898e4a2..95e02047c691 100644 --- a/paimon-python/pypaimon/tests/blob_test.py +++ b/paimon-python/pypaimon/tests/blob_test.py @@ -130,6 +130,36 @@ def tearDown(self): 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 36c1044443f4..6efdf9e49fa2 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 @@ def _write_blob_data(self, blob_value: Blob, crc32: int): 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 728d86c50fab..b366dab96ea1 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 @@ def __init__(self, file_io, file_path: Path, blob_consumer: Optional[BlobConsume 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 2e0e682130c1..c5bf143728a7 100644 --- a/paimon-python/pypaimon/write/writer/blob_writer.py +++ b/paimon-python/pypaimon/write/writer/blob_writer.py @@ -44,6 +44,8 @@ def __init__(self, table, partition: Tuple, bucket: int, max_seq_number: int, bl 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 @@ def open_current_writer(self): 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: