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 91ce4d60e4 [core] Support dynamic read batch sizing (#9243)
91ce4d60e4 is described below
commit 91ce4d60e41b23fe102410ca2bb316a55ada9a74
Author: Gabriel <[email protected]>
AuthorDate: Sun Aug 16 15:13:33 2026 +0800
[core] Support dynamic read batch sizing (#9243)
---
docs/docs/program-api/java-api.mdx | 41 ++++++
.../apache/paimon/format/FormatReaderContext.java | 20 ++-
.../apache/paimon/format/FormatReaderFactory.java | 7 ++
.../paimon/format/OrcFormatReaderContext.java | 13 ++
.../paimon/reader/ReadBatchSizeController.java | 88 +++++++++++++
.../paimon/reader/ReadBatchSizeControllerTest.java | 53 ++++++++
.../paimon/io/ChainKeyValueFileReaderFactory.java | 35 +++++-
.../paimon/io/KeyValueFileReaderFactory.java | 93 ++++++++++----
.../paimon/operation/DataEvolutionSplitRead.java | 15 ++-
.../paimon/operation/MergeFileSplitRead.java | 7 ++
.../operation/PrimaryKeyIndexedSplitRead.java | 7 ++
.../apache/paimon/operation/RawFileSplitRead.java | 14 ++-
.../org/apache/paimon/operation/SplitRead.java | 11 ++
.../apache/paimon/table/ChainGroupReadTable.java | 8 ++
.../paimon/table/ChainTableFileStoreTable.java | 8 ++
.../paimon/table/FallbackReadFileStoreTable.java | 8 ++
.../paimon/table/format/FormatReadBuilder.java | 26 +++-
.../paimon/table/format/FormatTableRead.java | 15 ++-
.../paimon/table/source/AppendTableRead.java | 17 +++
.../table/source/DataEvolutionTableRead.java | 6 +
.../apache/paimon/table/source/InnerTableRead.java | 6 +
.../paimon/table/source/KeyValueTableRead.java | 15 +++
.../paimon/table/source/ReadBuilderImpl.java | 7 ++
.../org/apache/paimon/table/source/TableRead.java | 16 +++
.../source/splitread/IncrementalDiffSplitRead.java | 7 ++
.../apache/paimon/table/system/AuditLogTable.java | 8 ++
.../paimon/operation/RawFileSplitReadTest.java | 38 +++++-
.../paimon/table/format/FormatReadBuilderTest.java | 81 ++++++++++++
.../table/source/DataEvolutionTableReadTest.java | 68 ++++++++++
.../source/PrimaryKeyVectorPositionReaderTest.java | 4 +
.../paimon/table/system/AuditLogTableTest.java | 22 ++++
.../paimon/table/system/BinlogTableTest.java | 22 ++++
.../apache/paimon/utils/AsyncRecordReaderTest.java | 70 +++++++++++
.../paimon/flink/lookup/LookupCompactDiffRead.java | 9 ++
.../flink/lookup/LookupCompactDiffReadTest.java | 44 +++++++
.../apache/paimon/format/orc/OrcReaderFactory.java | 60 +++++++--
.../format/orc/reader/AbstractOrcColumnVector.java | 5 +
.../format/parquet/ParquetReaderFactory.java | 23 +++-
.../reader/VectorizedParquetRecordReader.java | 63 +++++++++-
.../paimon/format/orc/OrcReaderFactoryTest.java | 121 ++++++++++++++++++
.../format/parquet/ParquetReadWriteTest.java | 137 +++++++++++++++++++++
41 files changed, 1266 insertions(+), 52 deletions(-)
diff --git a/docs/docs/program-api/java-api.mdx
b/docs/docs/program-api/java-api.mdx
index 67979e76f0..42255dc249 100644
--- a/docs/docs/program-api/java-api.mdx
+++ b/docs/docs/program-api/java-api.mdx
@@ -212,6 +212,47 @@ public class ReadTable {
}
```
+### Adjust Read Batch Size at Runtime
+
+Parquet and ORC readers can share a `ReadBatchSizeController` to adjust the
row count and vector
+capacity of future physical batches without recreating readers:
+
+```java
+import org.apache.paimon.reader.ReadBatchSizeController;
+import org.apache.paimon.table.source.TableRead;
+
+ReadBatchSizeController controller = new ReadBatchSizeController(1024, 1024);
+TableRead read = readBuilder.newRead().withReadBatchSizeController(controller);
+RecordReader<InternalRow> reader = read.createReader(splits);
+
+controller.setRequestedBatchSize(256);
+```
+
+Configure the controller on `TableRead` before creating readers. The first
constructor argument is
+an immutable safety limit; it does not preallocate vectors of that size. The
second argument is the
+initial requested size. Every requested size must be between `1` and the
safety limit.
+
+A supporting reader snapshots the requested size before starting a physical
batch. If the size has
+changed, it replaces an idle reusable batch with vectors sized for the new
value and then starts the
+read. The allocation is reused until the requested size changes again.
Consequently, lowering the
+requested size reduces the vector capacity of future batches instead of only
changing their logical
+row count.
+
+An update never mutates a batch that has already started or is still owned by
a consumer.
+Asynchronously prefetched batches may therefore retain the previous size. For
a pooled ORC reader,
+each idle pool entry adopts the current size the next time it is acquired,
while in-flight entries
+keep their old vectors until released. During such a transition, old and new
vectors can coexist.
+
+The controller uses latest-value semantics: when updates happen faster than
physical batches start,
+readers may skip intermediate requested sizes. Engines should avoid changing
the size too frequently
+because each observed size change reallocates vectors and can add allocation
and garbage-collection
+overhead. A hysteresis interval or minimum adjustment period is recommended.
+
+For concurrent scans, estimate memory using the requested size multiplied by
the number of active
+and prefetched batches. The maximum limits an accidental or unsafe request,
but it is not reserved
+up front. This allows an engine to reduce future batch capacities under memory
pressure and grow
+them again when more memory is available.
+
## Batch Write
The writing is divided into two stages:
diff --git
a/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderContext.java
b/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderContext.java
index 053acd99a4..cd3e97dcb0 100644
---
a/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderContext.java
+++
b/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderContext.java
@@ -20,6 +20,7 @@ package org.apache.paimon.format;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.utils.RoaringBitmap32;
@@ -32,17 +33,28 @@ public class FormatReaderContext implements
FormatReaderFactory.Context {
private final Path file;
private final long fileSize;
@Nullable private final RoaringBitmap32 selection;
+ @Nullable private final ReadBatchSizeController readBatchSizeController;
public FormatReaderContext(FileIO fileIO, Path file, long fileSize) {
- this(fileIO, file, fileSize, null);
+ this(fileIO, file, fileSize, null, null);
}
public FormatReaderContext(
FileIO fileIO, Path file, long fileSize, @Nullable RoaringBitmap32
selection) {
+ this(fileIO, file, fileSize, selection, null);
+ }
+
+ public FormatReaderContext(
+ FileIO fileIO,
+ Path file,
+ long fileSize,
+ @Nullable RoaringBitmap32 selection,
+ @Nullable ReadBatchSizeController readBatchSizeController) {
this.fileIO = fileIO;
this.file = file;
this.fileSize = fileSize;
this.selection = selection;
+ this.readBatchSizeController = readBatchSizeController;
}
@Override
@@ -65,4 +77,10 @@ public class FormatReaderContext implements
FormatReaderFactory.Context {
public RoaringBitmap32 selection() {
return selection;
}
+
+ @Nullable
+ @Override
+ public ReadBatchSizeController readBatchSizeController() {
+ return readBatchSizeController;
+ }
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderFactory.java
b/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderFactory.java
index fae5f6d8a5..73ae95fa1a 100644
---
a/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderFactory.java
+++
b/paimon-common/src/main/java/org/apache/paimon/format/FormatReaderFactory.java
@@ -22,6 +22,7 @@ import org.apache.paimon.data.InternalRow;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.utils.RoaringBitmap32;
@@ -53,5 +54,11 @@ public interface FormatReaderFactory {
@Nullable
RoaringBitmap32 selection();
+
+ /** Controller shared by readers that support dynamic read batch
sizing. */
+ @Nullable
+ default ReadBatchSizeController readBatchSizeController() {
+ return null;
+ }
}
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/format/OrcFormatReaderContext.java
b/paimon-common/src/main/java/org/apache/paimon/format/OrcFormatReaderContext.java
index 8b761867fa..93b225cd16 100644
---
a/paimon-common/src/main/java/org/apache/paimon/format/OrcFormatReaderContext.java
+++
b/paimon-common/src/main/java/org/apache/paimon/format/OrcFormatReaderContext.java
@@ -20,8 +20,11 @@ package org.apache.paimon.format;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
+import javax.annotation.Nullable;
+
/** The context for creating orc {@link RecordReader}. */
public class OrcFormatReaderContext extends FormatReaderContext {
@@ -32,6 +35,16 @@ public class OrcFormatReaderContext extends
FormatReaderContext {
this.poolSize = poolSize;
}
+ public OrcFormatReaderContext(
+ FileIO fileIO,
+ Path filePath,
+ long fileSize,
+ int poolSize,
+ @Nullable ReadBatchSizeController readBatchSizeController) {
+ super(fileIO, filePath, fileSize, null, readBatchSizeController);
+ this.poolSize = poolSize;
+ }
+
public int poolSize() {
return poolSize;
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/reader/ReadBatchSizeController.java
b/paimon-common/src/main/java/org/apache/paimon/reader/ReadBatchSizeController.java
new file mode 100644
index 0000000000..90d347654a
--- /dev/null
+++
b/paimon-common/src/main/java/org/apache/paimon/reader/ReadBatchSizeController.java
@@ -0,0 +1,88 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.reader;
+
+import org.apache.paimon.annotation.Public;
+
+import javax.annotation.concurrent.ThreadSafe;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.apache.paimon.utils.Preconditions.checkArgument;
+
+/**
+ * Thread-safe controller for changing the requested read batch size within a
fixed maximum.
+ *
+ * <p>Supporting readers snapshot {@link #requestedBatchSize()} before
starting each physical batch
+ * and use that value for both the logical row count and vector capacity. When
the requested size
+ * changes, an idle reusable batch is replaced at the next safe batch
boundary. A physical batch
+ * that has already started, including an asynchronously prefetched batch,
retains its previous size
+ * and vectors.
+ *
+ * <p>{@link #maxBatchSize()} is a validation limit rather than a preallocated
vector capacity.
+ * Concurrent updates use latest-value semantics, so readers are not required
to observe every
+ * intermediate requested size.
+ */
+@Public
+@ThreadSafe
+public final class ReadBatchSizeController {
+
+ private final int maxBatchSize;
+ private final AtomicInteger requestedBatchSize;
+
+ public ReadBatchSizeController(int maxBatchSize, int requestedBatchSize) {
+ checkArgument(maxBatchSize > 0, "Maximum batch size must be
positive.");
+ checkRequestedBatchSize(maxBatchSize, requestedBatchSize);
+ this.maxBatchSize = maxBatchSize;
+ this.requestedBatchSize = new AtomicInteger(requestedBatchSize);
+ }
+
+ /** Maximum permitted requested batch size. */
+ public int maxBatchSize() {
+ return maxBatchSize;
+ }
+
+ /**
+ * Requested row count and vector capacity for a future physical batch.
+ *
+ * <p>Readers snapshot this value at a format-specific physical batch
boundary.
+ */
+ public int requestedBatchSize() {
+ return requestedBatchSize.get();
+ }
+
+ /**
+ * Set the requested size for future physical batches.
+ *
+ * <p>The value must be between {@code 1} and {@link #maxBatchSize()},
inclusive. A reader that
+ * already started or prefetched a physical batch may finish that batch
with the previous size.
+ */
+ public void setRequestedBatchSize(int requestedBatchSize) {
+ checkRequestedBatchSize(maxBatchSize, requestedBatchSize);
+ this.requestedBatchSize.set(requestedBatchSize);
+ }
+
+ private static void checkRequestedBatchSize(int maxBatchSize, int
requestedBatchSize) {
+ checkArgument(
+ requestedBatchSize > 0 && requestedBatchSize <= maxBatchSize,
+ "Requested batch size must be between 1 and %s, but was %s.",
+ maxBatchSize,
+ requestedBatchSize);
+ }
+}
diff --git
a/paimon-common/src/test/java/org/apache/paimon/reader/ReadBatchSizeControllerTest.java
b/paimon-common/src/test/java/org/apache/paimon/reader/ReadBatchSizeControllerTest.java
new file mode 100644
index 0000000000..9408a94c6d
--- /dev/null
+++
b/paimon-common/src/test/java/org/apache/paimon/reader/ReadBatchSizeControllerTest.java
@@ -0,0 +1,53 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.reader;
+
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Tests for {@link ReadBatchSizeController}. */
+class ReadBatchSizeControllerTest {
+
+ @Test
+ void testUpdateRequestedBatchSizeWithinMaximum() {
+ ReadBatchSizeController controller = new ReadBatchSizeController(1024,
64);
+
+ assertThat(controller.maxBatchSize()).isEqualTo(1024);
+ assertThat(controller.requestedBatchSize()).isEqualTo(64);
+
+ controller.setRequestedBatchSize(512);
+ assertThat(controller.requestedBatchSize()).isEqualTo(512);
+ }
+
+ @Test
+ void testRejectInvalidBatchSizes() {
+ assertThatThrownBy(() -> new ReadBatchSizeController(0, 0))
+ .isInstanceOf(IllegalArgumentException.class);
+ assertThatThrownBy(() -> new ReadBatchSizeController(1024, 1025))
+ .isInstanceOf(IllegalArgumentException.class);
+
+ ReadBatchSizeController controller = new ReadBatchSizeController(1024,
64);
+ assertThatThrownBy(() -> controller.setRequestedBatchSize(0))
+ .isInstanceOf(IllegalArgumentException.class);
+ assertThatThrownBy(() -> controller.setRequestedBatchSize(1025))
+ .isInstanceOf(IllegalArgumentException.class);
+ }
+}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java
b/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java
index 0b957bf6f2..fd02390a6c 100644
---
a/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java
+++
b/paimon-core/src/main/java/org/apache/paimon/io/ChainKeyValueFileReaderFactory.java
@@ -27,6 +27,7 @@ import
org.apache.paimon.deletionvectors.ExposeDeletionKeyValueReader;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.schema.SchemaManager;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.types.RowType;
@@ -61,6 +62,34 @@ public class ChainKeyValueFileReaderFactory extends
KeyValueFileReaderFactory {
DeletionVector.Factory dvFactory,
ChainReadContext chainReadContext,
CoreOptions coreOptions) {
+ this(
+ fileIO,
+ schemaManager,
+ schema,
+ keyType,
+ valueType,
+ formatReaderMappingBuilder,
+ pathFactory,
+ partition,
+ dvFactory,
+ chainReadContext,
+ coreOptions,
+ null);
+ }
+
+ public ChainKeyValueFileReaderFactory(
+ FileIO fileIO,
+ SchemaManager schemaManager,
+ TableSchema schema,
+ RowType keyType,
+ RowType valueType,
+ FormatReaderMapping.Builder formatReaderMappingBuilder,
+ DataFilePathFactory pathFactory,
+ BinaryRow partition,
+ DeletionVector.Factory dvFactory,
+ ChainReadContext chainReadContext,
+ CoreOptions coreOptions,
+ @Nullable ReadBatchSizeController readBatchSizeController) {
super(
fileIO,
schemaManager,
@@ -71,7 +100,8 @@ public class ChainKeyValueFileReaderFactory extends
KeyValueFileReaderFactory {
pathFactory,
partition,
dvFactory,
- coreOptions);
+ coreOptions,
+ readBatchSizeController);
this.chainReadContext = chainReadContext;
CoreOptions options = new CoreOptions(schema.options());
this.currentBranch = options.branch();
@@ -161,7 +191,8 @@ public class ChainKeyValueFileReaderFactory extends
KeyValueFileReaderFactory {
partition,
dvFactory,
chainReadContext,
- wrapped.options);
+ wrapped.options,
+ wrapped.readBatchSizeController);
}
}
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java
b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java
index 20a6fe395c..079d84c7e3 100644
---
a/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java
+++
b/paimon-core/src/main/java/org/apache/paimon/io/KeyValueFileReaderFactory.java
@@ -34,6 +34,7 @@ import org.apache.paimon.manifest.FileSource;
import org.apache.paimon.partition.PartitionUtils;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.KeyValueFieldsExtractor;
import org.apache.paimon.schema.SchemaManager;
@@ -74,6 +75,7 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
private final Map<FormatKey, FormatReaderMapping> formatReaderMappings;
private final BinaryRow partition;
protected final DeletionVector.Factory dvFactory;
+ @Nullable private final ReadBatchSizeController readBatchSizeController;
protected KeyValueFileReaderFactory(
FileIO fileIO,
@@ -86,6 +88,32 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
BinaryRow partition,
DeletionVector.Factory dvFactory,
CoreOptions coreOptions) {
+ this(
+ fileIO,
+ schemaManager,
+ schema,
+ keyType,
+ valueType,
+ formatReaderMappingBuilder,
+ pathFactory,
+ partition,
+ dvFactory,
+ coreOptions,
+ null);
+ }
+
+ protected KeyValueFileReaderFactory(
+ FileIO fileIO,
+ SchemaManager schemaManager,
+ TableSchema schema,
+ RowType keyType,
+ RowType valueType,
+ FormatReaderMapping.Builder formatReaderMappingBuilder,
+ DataFilePathFactory pathFactory,
+ BinaryRow partition,
+ DeletionVector.Factory dvFactory,
+ CoreOptions coreOptions,
+ @Nullable ReadBatchSizeController readBatchSizeController) {
this.fileIO = fileIO;
this.schemaManager = schemaManager;
this.schema = schema;
@@ -100,6 +128,7 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
this.partition = partition;
this.formatReaderMappings = new ConcurrentHashMap<>();
this.dvFactory = dvFactory;
+ this.readBatchSizeController = readBatchSizeController;
}
public TableSchema schema() {
@@ -172,9 +201,14 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
schema.logicalRowType(),
formatReaderMapping.getReaderFactory(),
orcPoolSize == null
- ? new FormatReaderContext(fileIO, filePath,
fileSize)
+ ? new FormatReaderContext(
+ fileIO, filePath, fileSize, null,
readBatchSizeController)
: new OrcFormatReaderContext(
- fileIO, filePath, fileSize,
orcPoolSize),
+ fileIO,
+ filePath,
+ fileSize,
+ orcPoolSize,
+ readBatchSizeController),
ignoreCorruptFiles,
ignoreLostFiles,
formatReaderMapping.getIndexMapping(),
@@ -240,6 +274,7 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
protected RowType readKeyType;
protected RowType readValueType;
+ @Nullable protected ReadBatchSizeController readBatchSizeController;
private Builder(
FileIO fileIO,
@@ -266,29 +301,35 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
}
public Builder copyWithoutProjection() {
- return new Builder(
- fileIO,
- schemaManager,
- schema,
- keyType,
- valueType,
- formatDiscover,
- pathFactory,
- extractor,
- options);
+ Builder copy =
+ new Builder(
+ fileIO,
+ schemaManager,
+ schema,
+ keyType,
+ valueType,
+ formatDiscover,
+ pathFactory,
+ extractor,
+ options);
+ copy.readBatchSizeController = readBatchSizeController;
+ return copy;
}
public Builder copyWithoutValue() {
- return new Builder(
- fileIO,
- schemaManager,
- schema,
- keyType,
- RowType.of(),
- formatDiscover,
- pathFactory,
- extractor,
- options);
+ Builder copy =
+ new Builder(
+ fileIO,
+ schemaManager,
+ schema,
+ keyType,
+ RowType.of(),
+ formatDiscover,
+ pathFactory,
+ extractor,
+ options);
+ copy.readBatchSizeController = readBatchSizeController;
+ return copy;
}
public Builder withReadKeyType(RowType readKeyType) {
@@ -301,6 +342,11 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
return this;
}
+ public Builder withReadBatchSizeController(ReadBatchSizeController
controller) {
+ this.readBatchSizeController = controller;
+ return this;
+ }
+
public RowType keyType() {
return keyType;
}
@@ -341,7 +387,8 @@ public class KeyValueFileReaderFactory implements
FileReaderFactory<KeyValue> {
pathFactory.createDataFilePathFactory(partition, bucket),
partition,
dvFactory,
- options);
+ options,
+ readBatchSizeController);
}
protected FormatReaderMapping.Builder formatReaderMappingBuilder(
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
index 9cde51bccf..4a0524cb3e 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionSplitRead.java
@@ -46,6 +46,7 @@ import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.reader.DataEvolutionFileReader;
import org.apache.paimon.reader.EmptyFileRecordReader;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.ReaderSupplier;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.SchemaEvolutionUtil;
@@ -128,6 +129,7 @@ public class DataEvolutionSplitRead implements
SplitRead<InternalRow> {
protected RowType readRowType;
@Nullable private List<Predicate> filters;
+ @Nullable private ReadBatchSizeController readBatchSizeController;
public DataEvolutionSplitRead(
FileIO fileIO,
@@ -177,6 +179,12 @@ public class DataEvolutionSplitRead implements
SplitRead<InternalRow> {
return this;
}
+ @Override
+ public SplitRead<InternalRow>
withReadBatchSizeController(ReadBatchSizeController controller) {
+ this.readBatchSizeController = controller;
+ return this;
+ }
+
/**
* Row tracking fields are assigned from the manifest entry instead of
being read from the file,
* and data evolution may reassign row ids, so a physical copy in the file
can be stale. Never
@@ -598,7 +606,12 @@ public class DataEvolutionSplitRead implements
SplitRead<InternalRow> {
}
FormatReaderContext formatReaderContext =
- new FormatReaderContext(fileIO, readTarget.path,
readTarget.fileSize, selection);
+ new FormatReaderContext(
+ fileIO,
+ readTarget.path,
+ readTarget.fileSize,
+ selection,
+ readBatchSizeController);
FileRecordReader<InternalRow> fileRecordReader =
new DataFileRecordReader(
readRowType,
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java
b/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java
index 2714177959..730d72bc89 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/MergeFileSplitRead.java
@@ -43,6 +43,7 @@ import
org.apache.paimon.mergetree.compact.MergeFunctionWrapper;
import org.apache.paimon.mergetree.compact.ReducerMergeFunctionWrapper;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.reader.EmptyRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.ReaderSupplier;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
@@ -187,6 +188,12 @@ public class MergeFileSplitRead implements
SplitRead<KeyValue> {
return this;
}
+ @Override
+ public MergeFileSplitRead
withReadBatchSizeController(ReadBatchSizeController controller) {
+ readerFactoryBuilder.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public MergeFileSplitRead forceKeepDelete() {
this.forceKeepDelete = true;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/PrimaryKeyIndexedSplitRead.java
b/paimon-core/src/main/java/org/apache/paimon/operation/PrimaryKeyIndexedSplitRead.java
index 519d2c4a94..8ce3066126 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/PrimaryKeyIndexedSplitRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/PrimaryKeyIndexedSplitRead.java
@@ -24,6 +24,7 @@ import org.apache.paimon.globalindex.IndexedSplit;
import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.table.source.DataSplit;
import org.apache.paimon.table.source.PrimaryKeyIndexPositionReader;
@@ -74,6 +75,12 @@ public class PrimaryKeyIndexedSplitRead implements
SplitRead<InternalRow> {
return this;
}
+ @Override
+ public SplitRead<InternalRow>
withReadBatchSizeController(ReadBatchSizeController controller) {
+ rawRead.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
IndexedSplit indexedSplit = (IndexedSplit) split;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
b/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
index 87f7063a81..57755e0ece 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java
@@ -41,6 +41,7 @@ import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.TopN;
import org.apache.paimon.reader.EmptyFileRecordReader;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.ReaderSupplier;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.SchemaManager;
@@ -90,6 +91,7 @@ public class RawFileSplitRead implements
SplitRead<InternalRow> {
@Nullable private List<Predicate> filters;
@Nullable private TopN topN;
@Nullable private Integer limit;
+ @Nullable private ReadBatchSizeController readBatchSizeController;
public RawFileSplitRead(
FileIO fileIO,
@@ -151,6 +153,12 @@ public class RawFileSplitRead implements
SplitRead<InternalRow> {
return this;
}
+ @Override
+ public SplitRead<InternalRow>
withReadBatchSizeController(ReadBatchSizeController controller) {
+ this.readBatchSizeController = controller;
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split s) throws IOException {
if (s instanceof DataSplit) {
@@ -338,7 +346,11 @@ public class RawFileSplitRead implements
SplitRead<InternalRow> {
FormatReaderContext formatReaderContext =
new FormatReaderContext(
- fileIO, dataFilePathFactory.toPath(file),
file.fileSize(), selection);
+ fileIO,
+ dataFilePathFactory.toPath(file),
+ file.fileSize(),
+ selection,
+ readBatchSizeController);
FileRecordReader<InternalRow> fileRecordReader =
new DataFileRecordReader(
outputRowType,
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/SplitRead.java
b/paimon-core/src/main/java/org/apache/paimon/operation/SplitRead.java
index f2cb5fdba0..75b2188b9e 100644
--- a/paimon-core/src/main/java/org/apache/paimon/operation/SplitRead.java
+++ b/paimon-core/src/main/java/org/apache/paimon/operation/SplitRead.java
@@ -21,6 +21,7 @@ package org.apache.paimon.operation;
import org.apache.paimon.disk.IOManager;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.TopN;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.table.source.Split;
import org.apache.paimon.types.RowType;
@@ -53,6 +54,10 @@ public interface SplitRead<T> {
return this;
}
+ default SplitRead<T> withReadBatchSizeController(ReadBatchSizeController
controller) {
+ return this;
+ }
+
/** Create a {@link RecordReader} from split. */
RecordReader<T> createReader(Split split) throws IOException;
@@ -83,6 +88,12 @@ public interface SplitRead<T> {
return this;
}
+ @Override
+ public SplitRead<R>
withReadBatchSizeController(ReadBatchSizeController controller) {
+ read.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<R> createReader(Split split) throws
IOException {
return splitConvert.apply(split);
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
index cb6541662a..2077e255e9 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/ChainGroupReadTable.java
@@ -28,6 +28,7 @@ import org.apache.paimon.manifest.PartitionEntry;
import org.apache.paimon.partition.PartitionPredicate;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.source.ChainSplit;
@@ -593,6 +594,13 @@ public class ChainGroupReadTable extends
FallbackReadFileStoreTable {
return this;
}
+ @Override
+ public InnerTableRead
withReadBatchSizeController(ReadBatchSizeController controller) {
+ mainRead.withReadBatchSizeController(controller);
+ fallbackRead.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
if (split instanceof ChainSplit || split instanceof DataSplit) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/ChainTableFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/ChainTableFileStoreTable.java
index 863cd969bc..3d91a46e61 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/ChainTableFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/ChainTableFileStoreTable.java
@@ -23,6 +23,7 @@ import org.apache.paimon.CoreOptions.StartupMode;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.disk.IOManager;
import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.source.ChainSplit;
@@ -203,6 +204,13 @@ public class ChainTableFileStoreTable extends
FallbackReadFileStoreTable {
return this;
}
+ @Override
+ public InnerTableRead
withReadBatchSizeController(ReadBatchSizeController controller) {
+ chainGroupRead.withReadBatchSizeController(controller);
+ fallbackRead.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
if (split instanceof FallbackSplit) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
index 0d970f2b99..09abec537d 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/FallbackReadFileStoreTable.java
@@ -35,6 +35,7 @@ import org.apache.paimon.options.Options;
import org.apache.paimon.partition.PartitionPredicate;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.TopN;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.source.DataFilePlan;
@@ -695,6 +696,13 @@ public class FallbackReadFileStoreTable extends
DelegatedFileStoreTable {
return this;
}
+ @Override
+ public InnerTableRead
withReadBatchSizeController(ReadBatchSizeController controller) {
+ mainRead.withReadBatchSizeController(controller);
+ fallbackRead.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
if (split instanceof FallbackSplit) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatReadBuilder.java
b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatReadBuilder.java
index d05f0ebc12..10653df1ff 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatReadBuilder.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatReadBuilder.java
@@ -35,6 +35,7 @@ import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.predicate.TopN;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.ReaderSupplier;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.table.FormatTable;
@@ -182,6 +183,12 @@ public class FormatReadBuilder implements ReadBuilder {
}
protected RecordReader<InternalRow> createReader(FormatDataSplit
dataSplit) throws IOException {
+ return createReader(dataSplit, null);
+ }
+
+ protected RecordReader<InternalRow> createReader(
+ FormatDataSplit dataSplit, @Nullable ReadBatchSizeController
readBatchSizeController)
+ throws IOException {
// Skip pushing down partition filters to reader.
List<Predicate> readFilters =
excludePredicateWithFields(
@@ -200,7 +207,14 @@ public class FormatReadBuilder implements ReadBuilder {
BinaryRow partition = dataSplit.partition();
List<ReaderSupplier<InternalRow>> suppliers = new ArrayList<>();
for (FormatDataSplit.FileMeta file : dataSplit.files()) {
- suppliers.add(() -> createFileReader(file, partition,
readerFactory, partitionMapping));
+ suppliers.add(
+ () ->
+ createFileReader(
+ file,
+ partition,
+ readerFactory,
+ partitionMapping,
+ readBatchSizeController));
}
return ConcatRecordReader.create(suppliers);
}
@@ -209,10 +223,16 @@ public class FormatReadBuilder implements ReadBuilder {
FormatDataSplit.FileMeta file,
@Nullable BinaryRow partition,
FormatReaderFactory readerFactory,
- Pair<int[], RowType> partitionMapping)
+ Pair<int[], RowType> partitionMapping,
+ @Nullable ReadBatchSizeController readBatchSizeController)
throws IOException {
FormatReaderContext formatReaderContext =
- new FormatReaderContext(table.fileIO(), file.filePath(),
file.fileSize(), null);
+ new FormatReaderContext(
+ table.fileIO(),
+ file.filePath(),
+ file.fileSize(),
+ null,
+ readBatchSizeController);
try {
FileRecordReader<InternalRow> reader;
Long length = file.length();
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableRead.java
index 9939b1b750..e32a2f3d6b 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/format/FormatTableRead.java
@@ -24,12 +24,15 @@ import org.apache.paimon.metrics.MetricRegistry;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateProjectionConverter;
import org.apache.paimon.reader.LimitRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.table.FormatTable;
import org.apache.paimon.table.source.Split;
import org.apache.paimon.table.source.TableRead;
import org.apache.paimon.types.RowType;
+import javax.annotation.Nullable;
+
import java.io.IOException;
import java.util.Optional;
@@ -43,6 +46,7 @@ public class FormatTableRead implements TableRead {
private final Integer limit;
private boolean executeFilter = false;
+ @Nullable private ReadBatchSizeController readBatchSizeController;
public FormatTableRead(
RowType readType,
@@ -73,10 +77,19 @@ public class FormatTableRead implements TableRead {
return this;
}
+ @Override
+ public TableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ this.readBatchSizeController = controller;
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
FormatDataSplit dataSplit = (FormatDataSplit) split;
- RecordReader<InternalRow> reader = read.createReader(dataSplit);
+ // Capture the binding per TableRead so lazy file suppliers cannot
observe another read's
+ // controller.
+ ReadBatchSizeController controller = this.readBatchSizeController;
+ RecordReader<InternalRow> reader = read.createReader(dataSplit,
controller);
if (executeFilter) {
reader = executeFilter(reader);
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/AppendTableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/AppendTableRead.java
index 07b365407e..e4d7657146 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/AppendTableRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/AppendTableRead.java
@@ -23,6 +23,7 @@ import org.apache.paimon.operation.MergeFileSplitRead;
import org.apache.paimon.operation.SplitRead;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.TopN;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.source.splitread.SplitReadConfig;
@@ -48,6 +49,7 @@ public class AppendTableRead extends AbstractDataTableRead {
private Predicate predicate = null;
protected TopN topN = null;
protected Integer limit = null;
+ @Nullable private ReadBatchSizeController readBatchSizeController;
public AppendTableRead(
List<Function<SplitReadConfig, SplitReadProvider>>
providerFactories,
@@ -76,6 +78,9 @@ public class AppendTableRead extends AbstractDataTableRead {
read.withFilter(predicate);
read.withTopN(topN);
read.withLimit(limit);
+ if (readBatchSizeController != null) {
+ read.withReadBatchSizeController(readBatchSizeController);
+ }
}
@Override
@@ -105,6 +110,18 @@ public class AppendTableRead extends AbstractDataTableRead
{
return this;
}
+ @Override
+ public InnerTableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ initialized().forEach(r -> r.withReadBatchSizeController(controller));
+ this.readBatchSizeController = controller;
+ return this;
+ }
+
+ @Nullable
+ protected ReadBatchSizeController readBatchSizeController() {
+ return readBatchSizeController;
+ }
+
@Override
public RecordReader<InternalRow> reader(Split split) throws IOException {
for (SplitReadProvider readProvider : readProviders) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionTableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionTableRead.java
index 1efc9352cf..584e597428 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionTableRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionTableRead.java
@@ -21,6 +21,7 @@ package org.apache.paimon.table.source;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.catalog.CatalogContext;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.source.splitread.SplitReadConfig;
@@ -57,6 +58,7 @@ public class DataEvolutionTableRead extends AppendTableRead {
QueryAuthContext queryAuthContext = unwrapQueryAuthSplit(split);
int[] blobViewFields =
BlobViewTableReadSupport.blobViewFieldIndexes(currentReadType(), options);
+ ReadBatchSizeController controller = readBatchSizeController();
if (catalogContext != null && blobViewFields.length > 0) {
if (readFactory == null) {
throw new IllegalStateException(
@@ -75,6 +77,10 @@ public class DataEvolutionTableRead extends AppendTableRead {
() -> createDataReader(queryAuthContext.split(),
queryAuthContext.authResult()),
() -> {
InnerTableRead prescanRead = readFactory.get();
+ if (controller != null) {
+ // Blob-view prescan is a separate physical read
under the same budget.
+
prescanRead.withReadBatchSizeController(controller);
+ }
if (executeFilter) {
prescanRead.executeFilter();
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/InnerTableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/InnerTableRead.java
index b4da78ef6f..f570c180f4 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/InnerTableRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/InnerTableRead.java
@@ -22,6 +22,7 @@ import org.apache.paimon.metrics.MetricRegistry;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.predicate.TopN;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.types.RowType;
import java.util.List;
@@ -63,6 +64,11 @@ public interface InnerTableRead extends TableRead {
return this;
}
+ @Override
+ default InnerTableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ return this;
+ }
+
@Override
default TableRead executeFilter() {
return this;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/KeyValueTableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/KeyValueTableRead.java
index 1c6de402d0..b4e4abc60b 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/KeyValueTableRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/KeyValueTableRead.java
@@ -30,6 +30,7 @@ import org.apache.paimon.operation.SplitRead;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.TopN;
import org.apache.paimon.reader.LimitRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import
org.apache.paimon.table.source.splitread.IncrementalChangelogReadProvider;
@@ -68,6 +69,7 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
private IOManager ioManager = null;
@Nullable private TopN topN = null;
@Nullable private Integer limit = null;
+ @Nullable private ReadBatchSizeController readBatchSizeController;
public KeyValueTableRead(
Supplier<MergeFileSplitRead> mergeReadSupplier,
@@ -111,6 +113,9 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
read = read.withTopN(topN);
}
read.withFilter(predicate).withIOManager(ioManager);
+ if (readBatchSizeController != null) {
+ read.withReadBatchSizeController(readBatchSizeController);
+ }
}
@Override
@@ -198,6 +203,9 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
if (executeFilter) {
read.executeFilter();
}
+ if (readBatchSizeController != null) {
+ read.withReadBatchSizeController(readBatchSizeController);
+ }
return read;
}
@@ -208,6 +216,13 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
return this;
}
+ @Override
+ public InnerTableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ initialized().forEach(r -> r.withReadBatchSizeController(controller));
+ this.readBatchSizeController = controller;
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> reader(Split split) throws IOException {
for (SplitReadProvider readProvider : readProviders) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/ReadBuilderImpl.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/ReadBuilderImpl.java
index 6dbc364002..08241b9e37 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/ReadBuilderImpl.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/ReadBuilderImpl.java
@@ -26,6 +26,7 @@ import org.apache.paimon.partition.PartitionPredicate;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.predicate.TopN;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.reader.RecordReader.RecordIterator;
import org.apache.paimon.table.InnerTable;
@@ -314,6 +315,12 @@ public class ReadBuilderImpl implements ReadBuilder {
return this;
}
+ @Override
+ public TableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ delegate.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
return limit(delegate.createReader(split));
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/TableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/TableRead.java
index 68d5a0f1f2..3202c9cb06 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/source/TableRead.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/source/TableRead.java
@@ -24,6 +24,7 @@ import org.apache.paimon.disk.IOManager;
import org.apache.paimon.mergetree.compact.ConcatRecordReader;
import org.apache.paimon.metrics.MetricRegistry;
import org.apache.paimon.operation.SplitRead;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.ReaderSupplier;
import org.apache.paimon.reader.RecordReader;
@@ -46,6 +47,21 @@ public interface TableRead {
TableRead withIOManager(IOManager ioManager);
+ /**
+ * Configure a controller shared by all physical readers created by this
table read.
+ *
+ * <p>The controller must be configured before creating readers. Reader
creation binds the
+ * controller instance, not its current integer value, so later updates
through the same
+ * controller remain visible. Replacing the controller on this table read
after reader creation
+ * is unsupported and is not required to affect existing readers.
+ *
+ * <p>Formats that support dynamic sizing snapshot the requested size when
the next physical
+ * batch starts; already started or asynchronously prefetched batches may
use the previous size.
+ */
+ default TableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ return this;
+ }
+
RecordReader<InternalRow> createReader(Split split) throws IOException;
default RecordReader<InternalRow> createReader(List<Split> splits) throws
IOException {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/splitread/IncrementalDiffSplitRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/splitread/IncrementalDiffSplitRead.java
index c5032158fa..6538d14453 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/splitread/IncrementalDiffSplitRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/splitread/IncrementalDiffSplitRead.java
@@ -28,6 +28,7 @@ import
org.apache.paimon.mergetree.compact.MergeFunctionWrapper;
import org.apache.paimon.operation.MergeFileSplitRead;
import org.apache.paimon.operation.SplitRead;
import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.table.source.IncrementalSplit;
import org.apache.paimon.table.source.KeyValueTableRead;
@@ -84,6 +85,12 @@ public class IncrementalDiffSplitRead implements
SplitRead<InternalRow> {
return this;
}
+ @Override
+ public SplitRead<InternalRow>
withReadBatchSizeController(ReadBatchSizeController controller) {
+ mergeRead.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split s) throws IOException {
IncrementalSplit split = (IncrementalSplit) s;
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java
index e13e32ce60..36eeb0657b 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/system/AuditLogTable.java
@@ -41,6 +41,7 @@ import org.apache.paimon.predicate.LeafPredicate;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.predicate.PredicateReplaceVisitor;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.SchemaManager;
import org.apache.paimon.table.DataTable;
@@ -785,6 +786,13 @@ public class AuditLogTable implements DataTable,
ReadonlyTable {
return this;
}
+ @Override
+ public InnerTableRead
withReadBatchSizeController(ReadBatchSizeController controller) {
+ // System-table wrappers must preserve memory control on the
physical data read.
+ dataRead.withReadBatchSizeController(controller);
+ return this;
+ }
+
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
return dataRead.createReader(split).transform(this::convertRow);
diff --git
a/paimon-core/src/test/java/org/apache/paimon/operation/RawFileSplitReadTest.java
b/paimon-core/src/test/java/org/apache/paimon/operation/RawFileSplitReadTest.java
index fd4fcb2d48..666df6ca4a 100644
---
a/paimon-core/src/test/java/org/apache/paimon/operation/RawFileSplitReadTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/operation/RawFileSplitReadTest.java
@@ -29,6 +29,7 @@ import org.apache.paimon.fs.Path;
import org.apache.paimon.fs.local.LocalFileIO;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.Schema;
import org.apache.paimon.schema.SchemaManager;
@@ -133,12 +134,45 @@ class RawFileSplitReadTest {
}
}
+ @Test
+ void testTableReadSharesDynamicBatchSizeController() throws Exception {
+ FileStoreTable table = createTable("dynamic-batch-size", 20);
+ ReadBatchSizeController controller = new ReadBatchSizeController(8, 5);
+ InnerTableRead read =
table.newRead().withReadBatchSizeController(controller);
+
+ try (RecordReader<InternalRow> reader =
read.createReader(singleSplit(table))) {
+ assertThat(readBatchSize(reader)).isEqualTo(5);
+
+ controller.setRequestedBatchSize(2);
+ assertThat(readBatchSize(reader)).isEqualTo(2);
+
+ controller.setRequestedBatchSize(8);
+ assertThat(readBatchSize(reader)).isEqualTo(8);
+ }
+ }
+
+ private static int readBatchSize(RecordReader<InternalRow> reader) throws
Exception {
+ RecordReader.RecordIterator<InternalRow> batch = reader.readBatch();
+ assertThat(batch).isNotNull();
+ int count = 0;
+ while (batch.next() != null) {
+ count++;
+ }
+ batch.releaseBatch();
+ return count;
+ }
+
private FileStoreTable createTable(String directory) throws Exception {
+ return createTable(directory, 1);
+ }
+
+ private FileStoreTable createTable(String directory, int rowCount) throws
Exception {
Path tablePath = new Path(tempDir.resolve(directory).toUri());
Options options = new Options();
options.set(CoreOptions.PATH, tablePath.toString());
options.set(CoreOptions.BUCKET, 1);
options.set(CoreOptions.BUCKET_KEY, "first");
+ options.set(CoreOptions.READ_BATCH_SIZE, 8);
Schema schema =
Schema.newBuilder()
.column("first", DataTypes.STRING())
@@ -153,7 +187,9 @@ class RawFileSplitReadTest {
BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
try (BatchTableWrite write = writeBuilder.newWrite();
BatchTableCommit commit = writeBuilder.newCommit()) {
- write.write(GenericRow.of(BinaryString.fromString("value"), 42));
+ for (int i = 0; i < rowCount; i++) {
+ write.write(GenericRow.of(BinaryString.fromString("value"), i
+ 42));
+ }
commit.commit(write.prepareCommit());
}
return table;
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/format/FormatReadBuilderTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/format/FormatReadBuilderTest.java
index 19bc1d74fc..18cfb7373c 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/format/FormatReadBuilderTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/format/FormatReadBuilderTest.java
@@ -33,8 +33,15 @@ import org.apache.paimon.fs.local.LocalFileIO;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.table.FormatTable;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchTableWrite;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.source.TableRead;
+import org.apache.paimon.table.source.TableScan;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.InstantiationUtil;
@@ -136,6 +143,35 @@ public class FormatReadBuilderTest {
assertThat(deserialized.newRead()).isNotNull();
}
+ @Test
+ public void testControllerDoesNotBreakBuilderSerialization() {
+ FormatReadBuilder readBuilder = new
FormatReadBuilder(createOrcTable("serializable"));
+ readBuilder.newRead().withReadBatchSizeController(new
ReadBatchSizeController(16, 4));
+
+ assertThatNoException().isThrownBy(() ->
InstantiationUtil.serializeObject(readBuilder));
+ }
+
+ @Test
+ public void testControllersAreIsolatedBetweenReads() throws Exception {
+ FormatTable table = createOrcTable("isolated");
+ writeRows(table, 20);
+ FormatReadBuilder readBuilder = new FormatReadBuilder(table);
+ TableScan.Plan plan = readBuilder.newScan().plan();
+ ReadBatchSizeController firstController = new
ReadBatchSizeController(16, 3);
+ ReadBatchSizeController secondController = new
ReadBatchSizeController(16, 5);
+ TableRead firstRead =
readBuilder.newRead().withReadBatchSizeController(firstController);
+ TableRead secondRead =
readBuilder.newRead().withReadBatchSizeController(secondController);
+
+ try (RecordReader<InternalRow> firstReader =
firstRead.createReader(plan);
+ RecordReader<InternalRow> secondReader =
secondRead.createReader(plan)) {
+ assertThat(readBatchSize(firstReader)).isEqualTo(3);
+ assertThat(readBatchSize(secondReader)).isEqualTo(5);
+
+ firstController.setRequestedBatchSize(2);
+ assertThat(readBatchSize(firstReader)).isEqualTo(2);
+ }
+ }
+
@Test
public void testCreateReaderWithCsvSplit() throws IOException {
RowType rowType =
@@ -228,4 +264,49 @@ public class FormatReadBuilderTest {
}
return result;
}
+
+ private FormatTable createOrcTable(String name) {
+ Path tablePath = new Path(tempPath.resolve(name).toUri());
+ Map<String, String> options = new HashMap<>();
+ options.put("path", tablePath.toString());
+ options.put("file.format", "orc");
+ options.put("file.compression", "zstd");
+ return FormatTable.builder()
+ .fileIO(LocalFileIO.create())
+ .identifier(Identifier.create("test_db", name))
+ .rowType(RowType.of(DataTypes.INT()))
+ .partitionKeys(new ArrayList<>())
+ .location(tablePath.toString())
+ .format(FormatTable.Format.ORC)
+ .options(options)
+ .build();
+ }
+
+ private static void writeRows(FormatTable table, int count) throws
Exception {
+ BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
+ List<CommitMessage> messages;
+ try (BatchTableWrite write = writeBuilder.newWrite()) {
+ for (int i = 0; i < count; i++) {
+ write.write(GenericRow.of(i));
+ }
+ messages = write.prepareCommit();
+ }
+ try (BatchTableCommit commit = writeBuilder.newCommit()) {
+ commit.commit(messages);
+ }
+ }
+
+ private static int readBatchSize(RecordReader<InternalRow> reader) throws
IOException {
+ RecordReader.RecordIterator<InternalRow> batch = reader.readBatch();
+ assertThat(batch).isNotNull();
+ int size = 0;
+ try {
+ while (batch.next() != null) {
+ size++;
+ }
+ } finally {
+ batch.releaseBatch();
+ }
+ return size;
+ }
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/DataEvolutionTableReadTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/DataEvolutionTableReadTest.java
new file mode 100644
index 0000000000..67da6a4a0d
--- /dev/null
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/DataEvolutionTableReadTest.java
@@ -0,0 +1,68 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.table.source;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.catalog.CatalogContext;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.reader.ReadBatchSizeController;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.util.Collections;
+
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests for {@link DataEvolutionTableRead}. */
+class DataEvolutionTableReadTest {
+
+ @Test
+ void testReadBatchSizeControllerPropagatesToBlobPrescan() throws
IOException {
+ Options options = new Options();
+ options.set(CoreOptions.BLOB_VIEW_FIELD, "blob");
+ TableSchema schema = mock(TableSchema.class);
+ when(schema.logicalRowType())
+ .thenReturn(RowType.builder().field("blob",
DataTypes.BLOB()).build());
+ InnerTableRead prescanRead = mock(InnerTableRead.class);
+ when(prescanRead.createReader(any(Split.class)))
+ .thenThrow(new IOException("expected prescan stop"));
+ DataEvolutionTableRead read =
+ new DataEvolutionTableRead(
+ Collections.emptyList(),
+ schema,
+ new CoreOptions(options),
+ CatalogContext.create(new Options()),
+ () -> prescanRead);
+ ReadBatchSizeController controller = new ReadBatchSizeController(16,
4);
+ read.withReadBatchSizeController(controller);
+
+ assertThatThrownBy(() -> read.createReader(mock(Split.class)))
+ .isInstanceOf(IOException.class)
+ .hasMessage("expected prescan stop");
+ verify(prescanRead).withReadBatchSizeController(controller);
+ }
+}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/PrimaryKeyVectorPositionReaderTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/PrimaryKeyVectorPositionReaderTest.java
index c74771b2fd..2b610389df 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/source/PrimaryKeyVectorPositionReaderTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/PrimaryKeyVectorPositionReaderTest.java
@@ -30,6 +30,7 @@ import org.apache.paimon.operation.MergeFileSplitRead;
import org.apache.paimon.operation.RawFileSplitRead;
import org.apache.paimon.reader.FileRecordIterator;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.ScoreRecordIterator;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.stats.SimpleStats;
@@ -99,8 +100,11 @@ class PrimaryKeyVectorPositionReaderTest {
mock(TableSchema.class),
CoreOptions.fromMap(Collections.emptyMap()),
null);
+ ReadBatchSizeController controller = new ReadBatchSizeController(8, 5);
+ tableRead.withReadBatchSizeController(controller);
assertThat(tableRead.createReader(split)).isInstanceOf(PrimaryKeyIndexPositionReader.class);
+ verify(rawRead).withReadBatchSizeController(controller);
verify(rawRead, never()).createReader(split);
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/system/AuditLogTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/system/AuditLogTableTest.java
index 00b42c2841..c06def6fad 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/system/AuditLogTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/system/AuditLogTableTest.java
@@ -31,6 +31,7 @@ import org.apache.paimon.fs.local.LocalFileIO;
import org.apache.paimon.globalindex.IndexedSplit;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.Schema;
import org.apache.paimon.schema.SchemaChange;
@@ -43,23 +44,44 @@ import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.FileStoreTableFactory;
import org.apache.paimon.table.TableTestBase;
import org.apache.paimon.table.source.ChainSplit;
+import org.apache.paimon.table.source.InnerTableRead;
import org.apache.paimon.table.source.ReadBuilder;
import org.apache.paimon.table.source.TableScan;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowKind;
+import org.apache.paimon.types.RowType;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
import java.util.List;
import static org.apache.paimon.catalog.Identifier.SYSTEM_TABLE_SPLITTER;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/** Unit tests for {@link AuditLogTable}. */
public class AuditLogTableTest extends TableTestBase {
+ @Test
+ public void testReadBatchSizeControllerPropagatesToDataRead() {
+ FileStoreTable wrapped = mock(FileStoreTable.class);
+ InnerTableRead dataRead = mock(InnerTableRead.class);
+ when(wrapped.options()).thenReturn(Collections.emptyMap());
+ when(wrapped.rowType()).thenReturn(RowType.of(DataTypes.INT()));
+ when(wrapped.newRead()).thenReturn(dataRead);
+ when(dataRead.forceKeepDelete()).thenReturn(dataRead);
+ ReadBatchSizeController controller = new ReadBatchSizeController(16,
4);
+
+ new
AuditLogTable(wrapped).newRead().withReadBatchSizeController(controller);
+
+ verify(dataRead).withReadBatchSizeController(controller);
+ }
+
@Test
public void testReadAuditLogFromLatest() throws Exception {
AuditLogTable auditLogTable = createAuditLogTable("audit_table",
false);
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/system/BinlogTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/system/BinlogTableTest.java
index 6f187fbffd..0f860f6a8a 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/system/BinlogTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/system/BinlogTableTest.java
@@ -27,6 +27,7 @@ import org.apache.paimon.data.InternalRow;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.schema.Schema;
import org.apache.paimon.schema.SchemaChange;
import org.apache.paimon.schema.SchemaManager;
@@ -35,20 +36,41 @@ import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.FileStoreTableFactory;
import org.apache.paimon.table.TableTestBase;
+import org.apache.paimon.table.source.InnerTableRead;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowKind;
+import org.apache.paimon.types.RowType;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
import static org.apache.paimon.catalog.Identifier.SYSTEM_TABLE_SPLITTER;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/** Unit tests for {@link BinlogTable}. */
public class BinlogTableTest extends TableTestBase {
+ @Test
+ public void testReadBatchSizeControllerPropagatesToDataRead() {
+ FileStoreTable wrapped = mock(FileStoreTable.class);
+ InnerTableRead dataRead = mock(InnerTableRead.class);
+ when(wrapped.options()).thenReturn(Collections.emptyMap());
+ when(wrapped.rowType()).thenReturn(RowType.of(DataTypes.INT()));
+ when(wrapped.newRead()).thenReturn(dataRead);
+ when(dataRead.forceKeepDelete()).thenReturn(dataRead);
+ ReadBatchSizeController controller = new ReadBatchSizeController(16,
4);
+
+ new
BinlogTable(wrapped).newRead().withReadBatchSizeController(controller);
+
+ verify(dataRead).withReadBatchSizeController(controller);
+ }
+
@Test
public void testReadBinlogFromLatest() throws Exception {
BinlogTable binlogTable = createBinlogTable("binlog_table", false);
diff --git
a/paimon-core/src/test/java/org/apache/paimon/utils/AsyncRecordReaderTest.java
b/paimon-core/src/test/java/org/apache/paimon/utils/AsyncRecordReaderTest.java
index aea412fed9..5eecb7dd5e 100644
---
a/paimon-core/src/test/java/org/apache/paimon/utils/AsyncRecordReaderTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/utils/AsyncRecordReaderTest.java
@@ -18,6 +18,7 @@
package org.apache.paimon.utils;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.junit.jupiter.api.Test;
@@ -30,6 +31,8 @@ import java.util.Arrays;
import java.util.LinkedList;
import java.util.List;
import java.util.Queue;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.assertj.core.api.Assertions.assertThat;
@@ -100,6 +103,73 @@ public class AsyncRecordReaderTest {
.hasMessageContaining(message);
}
+ @Test
+ public void testPrefetchedBatchesMayUsePreviousRequestedSize() throws
Exception {
+ ReadBatchSizeController controller = new ReadBatchSizeController(8, 5);
+ CountDownLatch twoBatchesPrefetched = new CountDownLatch(1);
+ CountDownLatch continueReading = new CountDownLatch(1);
+ AtomicInteger batchNumber = new AtomicInteger();
+ RecordReader<Integer> physicalReader =
+ new RecordReader<Integer>() {
+ @Nullable
+ @Override
+ public RecordIterator<Integer> readBatch() throws
IOException {
+ int current = batchNumber.getAndIncrement();
+ if (current == 2) {
+ twoBatchesPrefetched.countDown();
+ try {
+ continueReading.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IOException(e);
+ }
+ } else if (current > 2) {
+ return null;
+ }
+
+ int size = controller.requestedBatchSize();
+ AtomicInteger remaining = new AtomicInteger(size);
+ return new RecordIterator<Integer>() {
+ @Nullable
+ @Override
+ public Integer next() {
+ return remaining.getAndDecrement() > 0 ?
current : null;
+ }
+
+ @Override
+ public void releaseBatch() {}
+ };
+ }
+
+ @Override
+ public void close() {}
+ };
+
+ try (AsyncRecordReader<Integer> asyncReader =
+ new AsyncRecordReader<>(() -> physicalReader)) {
+ assertThat(twoBatchesPrefetched.await(30,
TimeUnit.SECONDS)).isTrue();
+ controller.setRequestedBatchSize(2);
+ continueReading.countDown();
+
+ assertThat(readBatchSize(asyncReader)).isEqualTo(5);
+ assertThat(readBatchSize(asyncReader)).isEqualTo(5);
+ assertThat(readBatchSize(asyncReader)).isEqualTo(2);
+ } finally {
+ continueReading.countDown();
+ }
+ }
+
+ private static int readBatchSize(RecordReader<Integer> reader) throws
IOException {
+ RecordReader.RecordIterator<Integer> batch = reader.readBatch();
+ assertThat(batch).isNotNull();
+ int size = 0;
+ while (batch.next() != null) {
+ size++;
+ }
+ batch.releaseBatch();
+ return size;
+ }
+
@Test
public void testClassLoader() throws IOException {
ClassLoader goodClassLoader =
Thread.currentThread().getContextClassLoader();
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupCompactDiffRead.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupCompactDiffRead.java
index e4870de583..a4a6cdae86 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupCompactDiffRead.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupCompactDiffRead.java
@@ -23,6 +23,7 @@ import org.apache.paimon.disk.IOManager;
import org.apache.paimon.operation.MergeFileSplitRead;
import org.apache.paimon.operation.SplitRead;
import org.apache.paimon.predicate.Predicate;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.source.AbstractDataTableRead;
@@ -87,4 +88,12 @@ public class LookupCompactDiffRead extends
AbstractDataTableRead {
incrementalDiffRead.withIOManager(ioManager);
return this;
}
+
+ @Override
+ public InnerTableRead withReadBatchSizeController(ReadBatchSizeController
controller) {
+ // Both lookup phases can reach physical files and must share memory
control.
+ fullPhaseMergeRead.withReadBatchSizeController(controller);
+ incrementalDiffRead.withReadBatchSizeController(controller);
+ return this;
+ }
}
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupCompactDiffReadTest.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupCompactDiffReadTest.java
new file mode 100644
index 0000000000..7a851e63ff
--- /dev/null
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupCompactDiffReadTest.java
@@ -0,0 +1,44 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.flink.lookup;
+
+import org.apache.paimon.operation.MergeFileSplitRead;
+import org.apache.paimon.reader.ReadBatchSizeController;
+import org.apache.paimon.schema.TableSchema;
+
+import org.junit.jupiter.api.Test;
+
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
+/** Tests for {@link LookupCompactDiffRead}. */
+class LookupCompactDiffReadTest {
+
+ @Test
+ void testReadBatchSizeControllerPropagatesToMergeRead() {
+ MergeFileSplitRead mergeRead = mock(MergeFileSplitRead.class);
+ LookupCompactDiffRead read = new LookupCompactDiffRead(mergeRead,
mock(TableSchema.class));
+ ReadBatchSizeController controller = new ReadBatchSizeController(16,
4);
+
+ read.withReadBatchSizeController(controller);
+
+ verify(mergeRead, times(2)).withReadBatchSizeController(controller);
+ }
+}
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
b/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
index 05ebc5da22..4ea340512a 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/orc/OrcReaderFactory.java
@@ -37,6 +37,7 @@ import
org.apache.paimon.format.shredding.ShreddingReadPlanFactory;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader.RecordIterator;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.RowType;
@@ -64,6 +65,7 @@ import java.nio.charset.StandardCharsets;
import java.util.Collections;
import java.util.List;
import java.util.Map;
+import java.util.function.IntFunction;
import static org.apache.paimon.format.orc.OrcTypeUtil.convertToOrcSchema;
import static
org.apache.paimon.format.orc.reader.AbstractOrcColumnVector.createPaimonVector;
@@ -123,13 +125,23 @@ public class OrcReaderFactory implements
FormatReaderFactory {
RowType physicalReadType = readPlan.physicalRowType();
TypeDescription physicalReadSchema =
readPlan.isIdentity() ? schema :
convertToOrcSchema(physicalReadType);
+ ReadBatchSizeController readBatchSizeController =
context.readBatchSizeController();
Pool<OrcReaderBatch> poolOfBatches =
createPoolOfBatches(
context.filePath(),
poolSize,
context.fileIO(),
physicalReadSchema,
- physicalReadType);
+ physicalReadType,
+ readBatchSizeController);
+ IntFunction<OrcReaderBatch> batchFactory =
+ size ->
+ createReaderBatch(
+ context.filePath(),
+ createBatchWrapper(physicalReadSchema,
size),
+ poolOfBatches.recycler(),
+ context.fileIO(),
+ physicalReadType);
OrcRecordReader orcReader =
createRecordReader(
@@ -142,7 +154,8 @@ public class OrcReaderFactory implements
FormatReaderFactory {
context.selection(),
deletionVectorsEnabled);
OrcVectorizedReader orcVectorizedReader =
- new OrcVectorizedReader(orcReader, poolOfBatches);
+ new OrcVectorizedReader(
+ orcReader, poolOfBatches, readBatchSizeController,
batchFactory);
return readPlan.isIdentity()
? orcVectorizedReader
: new ShreddingFormatReader(orcVectorizedReader, readPlan);
@@ -200,7 +213,7 @@ public class OrcReaderFactory implements
FormatReaderFactory {
// ------------------------------------------------------------------------
private Pool<OrcReaderBatch> createPoolOfBatches(Path filePath, int
numBatches, FileIO fileIO) {
- return createPoolOfBatches(filePath, numBatches, fileIO, schema,
tableType);
+ return createPoolOfBatches(filePath, numBatches, fileIO, schema,
tableType, null);
}
private Pool<OrcReaderBatch> createPoolOfBatches(
@@ -208,12 +221,16 @@ public class OrcReaderFactory implements
FormatReaderFactory {
int numBatches,
FileIO fileIO,
TypeDescription readSchema,
- RowType readType) {
+ RowType readType,
+ @Nullable ReadBatchSizeController readBatchSizeController) {
final Pool<OrcReaderBatch> pool = new Pool<>(numBatches);
+ int allocatedBatchSize =
+ readBatchSizeController == null
+ ? Math.max(1, batchSize / numBatches)
+ : readBatchSizeController.requestedBatchSize();
for (int i = 0; i < numBatches; i++) {
- final VectorizedRowBatch orcBatch =
- createBatchWrapper(readSchema, Math.max(1, batchSize /
numBatches));
+ final VectorizedRowBatch orcBatch = createBatchWrapper(readSchema,
allocatedBatchSize);
final OrcReaderBatch batch =
createReaderBatch(filePath, orcBatch, pool.recycler(),
fileIO, readType);
pool.add(batch);
@@ -259,6 +276,10 @@ public class OrcReaderFactory implements
FormatReaderFactory {
return orcVectorizedRowBatch;
}
+ public int batchSize() {
+ return orcVectorizedRowBatch.getMaxSize();
+ }
+
private ColumnarRowIterator convertAndGetIterator(
VectorizedRowBatch orcBatch, long rowNumber) {
// no copying from the ORC column vectors to the Paimon columns
vectors necessary,
@@ -288,17 +309,40 @@ public class OrcReaderFactory implements
FormatReaderFactory {
private final OrcRecordReader orcReader;
private final Pool<OrcReaderBatch> pool;
+ @Nullable private final ReadBatchSizeController
readBatchSizeController;
+ private final IntFunction<OrcReaderBatch> batchFactory;
private OrcVectorizedReader(
- final OrcRecordReader orcReader, final Pool<OrcReaderBatch>
pool) {
+ final OrcRecordReader orcReader,
+ final Pool<OrcReaderBatch> pool,
+ @Nullable final ReadBatchSizeController
readBatchSizeController,
+ final IntFunction<OrcReaderBatch> batchFactory) {
this.orcReader = checkNotNull(orcReader, "orcReader");
this.pool = checkNotNull(pool, "pool");
+ this.readBatchSizeController = readBatchSizeController;
+ this.batchFactory = batchFactory;
}
@Nullable
@Override
public ColumnarRowIterator readBatch() throws IOException {
- final OrcReaderBatch batch = getCachedEntry();
+ OrcReaderBatch batch = getCachedEntry();
+ // Snapshot after acquiring a reusable batch and before the
physical read starts.
+ int requestedBatchSize =
+ readBatchSizeController == null
+ ? batch.batchSize()
+ : readBatchSizeController.requestedBatchSize();
+ if (batch.batchSize() != requestedBatchSize) {
+ // Only an acquired idle entry can be replaced. In-flight
pooled batches retain
+ // their old vectors until consumers release them.
+ try {
+ batch = batchFactory.apply(requestedBatchSize);
+ } catch (RuntimeException | Error e) {
+ // Preserve the pool size invariant if allocating the
replacement fails.
+ batch.recycle();
+ throw e;
+ }
+ }
final VectorizedRowBatch orcVectorBatch =
batch.orcVectorizedRowBatch();
long rowNumber = orcReader.recordReader.getRowNumber();
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/orc/reader/AbstractOrcColumnVector.java
b/paimon-format/src/main/java/org/apache/paimon/format/orc/reader/AbstractOrcColumnVector.java
index a77f4ebe3d..0a96606d4e 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/orc/reader/AbstractOrcColumnVector.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/orc/reader/AbstractOrcColumnVector.java
@@ -60,6 +60,11 @@ public abstract class AbstractOrcColumnVector
return !vector.noNulls && vector.isNull[rowMapper(i)];
}
+ @Override
+ public int getCapacity() {
+ return vector.isNull.length;
+ }
+
public static org.apache.paimon.data.columnar.ColumnVector
createPaimonVector(
ColumnVector vector,
VectorizedRowBatch orcBatch,
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
index dd425a23a7..2fe7f7b928 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/ParquetReaderFactory.java
@@ -35,6 +35,7 @@ import org.apache.paimon.options.CatalogOptions;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.DataType;
@@ -70,6 +71,7 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.function.IntFunction;
import static
org.apache.paimon.data.columnar.ColumnVectorUtils.createParquetWritableColumnVector;
import static
org.apache.paimon.format.parquet.ParquetSchemaConverter.PAIMON_SCHEMA;
@@ -167,14 +169,21 @@ public class ParquetReaderFactory implements
FormatReaderFactory {
requestedSchema.messageType);
}
- int actualBatchSize = computeBatchSize(reader,
requestedSchema.messageType);
+ int configuredBatchSize = computeBatchSize(reader,
requestedSchema.messageType);
Preconditions.checkArgument(
- actualBatchSize > 0,
+ configuredBatchSize > 0,
"Parquet read batch size should be positive: %s",
- actualBatchSize);
+ configuredBatchSize);
+ ReadBatchSizeController readBatchSizeController =
context.readBatchSizeController();
+ int initialBatchSize =
+ readBatchSizeController == null
+ ? configuredBatchSize
+ : readBatchSizeController.requestedBatchSize();
reader.setRequestedSchema(requestedSchema.messageType);
WritableColumnVector[] writableVectors =
- createWritableVectors(actualBatchSize, physicalReadFields);
+ createWritableVectors(initialBatchSize, physicalReadFields);
+ IntFunction<WritableColumnVector[]> vectorFactory =
+ size -> createWritableVectors(size, physicalReadFields);
VectorizedParquetRecordReader parquetReader =
new VectorizedParquetRecordReader(
@@ -183,8 +192,10 @@ public class ParquetReaderFactory implements
FormatReaderFactory {
fileSchema,
requestedSchema.fields,
writableVectors,
- actualBatchSize,
- context.fileIO());
+ initialBatchSize,
+ context.fileIO(),
+ readBatchSizeController,
+ vectorFactory);
return readPlan.isIdentity()
? parquetReader
: new ShreddingFormatReader(parquetReader, readPlan);
diff --git
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
index 1b652c2f68..32f570f8f3 100644
---
a/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
+++
b/paimon-format/src/main/java/org/apache/paimon/format/parquet/reader/VectorizedParquetRecordReader.java
@@ -26,6 +26,7 @@ import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
import org.apache.paimon.reader.FileRecordIterator;
import org.apache.paimon.reader.FileRecordReader;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.parquet.VersionParser;
import org.apache.parquet.column.ColumnDescriptor;
@@ -42,6 +43,7 @@ import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
+import java.util.function.IntFunction;
import java.util.stream.Collectors;
import static java.lang.String.format;
@@ -53,7 +55,7 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
private ParquetFileReader reader;
// The capacity of vectorized batch.
- private final int batchSize;
+ private int batchSize;
/**
* The total number of rows this RecordReader will eventually read. The
sum of the rows of all
@@ -80,6 +82,8 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
private final MessageType fileSchema;
private final List<ParquetField> fields;
private final RowIndexGenerator rowIndexGenerator;
+ @Nullable private final ReadBatchSizeController readBatchSizeController;
+ @Nullable private final IntFunction<WritableColumnVector[]> vectorFactory;
private Set<ParquetField> missingColumns;
private VersionParser.ParsedVersion writerVersion;
@@ -93,6 +97,20 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
int batchSize,
FileIO fileIO)
throws IOException {
+ this(filePath, reader, fileSchema, fields, vectors, batchSize, fileIO,
null, null);
+ }
+
+ public VectorizedParquetRecordReader(
+ Path filePath,
+ ParquetFileReader reader,
+ MessageType fileSchema,
+ List<ParquetField> fields,
+ WritableColumnVector[] vectors,
+ int batchSize,
+ FileIO fileIO,
+ @Nullable ReadBatchSizeController readBatchSizeController,
+ @Nullable IntFunction<WritableColumnVector[]> vectorFactory)
+ throws IOException {
this.filePath = filePath;
this.reader = reader;
this.fileSchema = fileSchema;
@@ -101,6 +119,8 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
this.batchSize = batchSize;
this.fileIO = fileIO;
this.rowIndexGenerator = new RowIndexGenerator();
+ this.readBatchSizeController = readBatchSizeController;
+ this.vectorFactory = vectorFactory;
// fetch writer version from file metadata
try {
@@ -113,10 +133,10 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
// Check if all the required columns are present in the file.
checkMissingColumns();
// Initialize the columnarBatch and columnVectors,
- initBatch(vectors);
+ initBatch(vectors, batchSize);
}
- private void initBatch(WritableColumnVector[] vectors) {
+ private void initBatch(WritableColumnVector[] vectors, int capacity) {
columnarBatch =
new ColumnarBatch(
filePath,
@@ -130,7 +150,7 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
for (int i = 0; i < columnVectors.length; i++) {
columnVectors[i] =
new ParquetColumnVector(
- fields.get(i), vectors[i], batchSize,
missingColumns, true);
+ fields.get(i), vectors[i], capacity,
missingColumns, true);
}
}
@@ -180,13 +200,19 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
if (rowsReturned >= totalRowCount) {
return false;
}
+ // Snapshot once so a concurrent update only affects the next
physical batch.
+ int requestedBatchSize =
+ readBatchSizeController == null
+ ? batchSize
+ : readBatchSizeController.requestedBatchSize();
+ resizeBatchIfNeeded(requestedBatchSize);
for (ParquetColumnVector vector : columnVectors) {
vector.reset();
}
columnarBatch.setNumRows(0);
checkEndOfRowGroup();
- int num = (int) Math.min(batchSize, totalCountLoadedSoFar -
rowsReturned);
+ int num = (int) Math.min(requestedBatchSize, totalCountLoadedSoFar
- rowsReturned);
for (ParquetColumnVector cv : columnVectors) {
for (ParquetColumnVector leafCv : cv.getLeaves()) {
VectorizedColumnReader columnReader =
leafCv.getColumnReader();
@@ -220,6 +246,33 @@ public class VectorizedParquetRecordReader implements
FileRecordReader<InternalR
}
}
+ private void resizeBatchIfNeeded(int requestedBatchSize) {
+ if (requestedBatchSize == batchSize) {
+ return;
+ }
+
+ ParquetColumnVector[] previousVectors = columnVectors;
+ // A new physical batch starts only after the prior iterator is
released, so replacing the
+ // wrappers here cannot mutate vectors still visible to the consumer.
+ initBatch(vectorFactory.apply(requestedBatchSize), requestedBatchSize);
+ for (int i = 0; i < columnVectors.length; i++) {
+ copyColumnReaders(previousVectors[i], columnVectors[i]);
+ }
+ batchSize = requestedBatchSize;
+ }
+
+ private static void copyColumnReaders(
+ ParquetColumnVector previous, ParquetColumnVector replacement) {
+ if (previous.getColumn().isPrimitive()) {
+ replacement.setColumnReader(previous.getColumnReader());
+ return;
+ }
+
+ for (int i = 0; i < previous.getChildren().size(); i++) {
+ copyColumnReaders(previous.getChildren().get(i),
replacement.getChildren().get(i));
+ }
+ }
+
private void checkEndOfRowGroup() throws IOException {
if (rowsReturned != totalCountLoadedSoFar) {
return;
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcReaderFactoryTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcReaderFactoryTest.java
index 237066943b..fb754d3900 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcReaderFactoryTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/orc/OrcReaderFactoryTest.java
@@ -19,11 +19,13 @@
package org.apache.paimon.format.orc;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.columnar.ColumnarRow;
import org.apache.paimon.format.FormatReaderContext;
import org.apache.paimon.format.OrcFormatReaderContext;
import org.apache.paimon.format.orc.filter.OrcFilters;
import org.apache.paimon.fs.Path;
import org.apache.paimon.fs.local.LocalFileIO;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
@@ -171,6 +173,125 @@ class OrcReaderFactoryTest {
assertThat(totalF0.get()).isEqualTo(1844737280400L);
}
+ @Test
+ void testDynamicReadBatchSize() throws IOException {
+ OrcReaderFactory format = createFormat(FLAT_FILE_TYPE, new int[] {0});
+ ReadBatchSizeController controller = new
ReadBatchSizeController(BATCH_SIZE, 5);
+ LocalFileIO fileIO = new LocalFileIO();
+
+ try (RecordReader<InternalRow> reader =
+ format.createReader(
+ new FormatReaderContext(
+ fileIO,
+ flatFile,
+ fileIO.getFileSize(flatFile),
+ null,
+ controller))) {
+ assertThat(readBatch(reader)).isEqualTo(new BatchResult(5, 5));
+
+ controller.setRequestedBatchSize(2);
+ assertThat(readBatch(reader)).isEqualTo(new BatchResult(2, 2));
+
+ controller.setRequestedBatchSize(BATCH_SIZE);
+ assertThat(readBatch(reader)).isEqualTo(new
BatchResult(BATCH_SIZE, BATCH_SIZE));
+ }
+ }
+
+ @Test
+ void testDynamicReadBatchSizeWithPooledBatches() throws IOException {
+ OrcReaderFactory format = createFormat(FLAT_FILE_TYPE, new int[] {0});
+ ReadBatchSizeController controller = new
ReadBatchSizeController(BATCH_SIZE, 5);
+ LocalFileIO fileIO = new LocalFileIO();
+
+ try (RecordReader<InternalRow> reader =
+ format.createReader(
+ new OrcFormatReaderContext(
+ fileIO, flatFile,
fileIO.getFileSize(flatFile), 2, controller))) {
+ RecordReader.RecordIterator<InternalRow> first =
reader.readBatch();
+ assertThat(first).isNotNull();
+ assertThat(consumeBatch(first)).isEqualTo(new BatchResult(5, 5));
+
+ controller.setRequestedBatchSize(2);
+ RecordReader.RecordIterator<InternalRow> second =
reader.readBatch();
+ assertThat(second).isNotNull();
+ assertThat(consumeBatch(second)).isEqualTo(new BatchResult(2, 2));
+
+ first.releaseBatch();
+ second.releaseBatch();
+
+ controller.setRequestedBatchSize(BATCH_SIZE);
+ assertThat(readBatch(reader)).isEqualTo(new
BatchResult(BATCH_SIZE, BATCH_SIZE));
+ }
+ }
+
+ @Test
+ void testStaticReadBatchSizeKeepsPoolBudget() throws IOException {
+ OrcReaderFactory format = createFormat(FLAT_FILE_TYPE, new int[] {0});
+ LocalFileIO fileIO = new LocalFileIO();
+
+ try (RecordReader<InternalRow> reader =
+ format.createReader(
+ new OrcFormatReaderContext(
+ fileIO, flatFile,
fileIO.getFileSize(flatFile), 3))) {
+ assertThat(readBatch(reader)).isEqualTo(new BatchResult(3, 3));
+ }
+ }
+
+ private static BatchResult readBatch(RecordReader<InternalRow> reader)
throws IOException {
+ RecordReader.RecordIterator<InternalRow> batch = reader.readBatch();
+ assertThat(batch).isNotNull();
+ BatchResult result = consumeBatch(batch);
+ batch.releaseBatch();
+ return result;
+ }
+
+ private static BatchResult
consumeBatch(RecordReader.RecordIterator<InternalRow> batch)
+ throws IOException {
+ int count = 0;
+ int capacity = -1;
+ InternalRow row;
+ while ((row = batch.next()) != null) {
+ if (capacity < 0) {
+ capacity = ((ColumnarRow)
row).batch().columns[0].getCapacity();
+ }
+ count++;
+ }
+ return new BatchResult(count, capacity);
+ }
+
+ private static class BatchResult {
+
+ private final int size;
+ private final int capacity;
+
+ private BatchResult(int size, int capacity) {
+ this.size = size;
+ this.capacity = capacity;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof BatchResult)) {
+ return false;
+ }
+ BatchResult that = (BatchResult) o;
+ return size == that.size && capacity == that.capacity;
+ }
+
+ @Override
+ public int hashCode() {
+ return 31 * size + capacity;
+ }
+
+ @Override
+ public String toString() {
+ return "BatchResult{" + "size=" + size + ", capacity=" + capacity
+ '}';
+ }
+ }
+
@RepeatedTest(10)
void testReadRowPositionWithRandomFilterAndPool() throws IOException {
ArrayList<OrcFilters.Predicate> predicates = new ArrayList<>();
diff --git
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
index f585a0cfdd..b048ccc3c1 100644
---
a/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
+++
b/paimon-format/src/test/java/org/apache/paimon/format/parquet/ParquetReadWriteTest.java
@@ -28,6 +28,7 @@ import org.apache.paimon.data.InternalMap;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.InternalVector;
import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.data.columnar.ColumnarRow;
import org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.format.FormatReaderContext;
import org.apache.paimon.format.FormatWriter;
@@ -36,6 +37,7 @@ import org.apache.paimon.fs.Path;
import org.apache.paimon.fs.local.LocalFileIO;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.reader.ReadBatchSizeController;
import org.apache.paimon.reader.RecordReader;
import org.apache.paimon.types.ArrayType;
import org.apache.paimon.types.BigIntType;
@@ -204,6 +206,141 @@ public class ParquetReadWriteTest {
return Arrays.asList(10, 1000);
}
+ @Test
+ void testDynamicReadBatchSize() throws IOException {
+ List<InternalRow> records = new ArrayList<>();
+ for (int i = 0; i < 20; i++) {
+ records.add(newRow(i));
+ }
+ Path path = createTempParquetFileByPaimon(folder, records, 10_000,
ROW_TYPE);
+
+ ParquetReaderFactory factory =
+ new ParquetReaderFactory(
+ new Options(),
+ RowType.builder().field("f4", new IntType()).build(),
+ 4,
+ null);
+ ReadBatchSizeController controller = new ReadBatchSizeController(8, 5);
+ LocalFileIO fileIO = new LocalFileIO();
+ try (RecordReader<InternalRow> reader =
+ factory.createReader(
+ new FormatReaderContext(
+ fileIO, path, fileIO.getFileSize(path), null,
controller))) {
+ BatchResult firstBatch = readIntBatch(reader);
+ assertThat(firstBatch.values).containsExactly(0, 1, 2, 3, 4);
+ assertThat(firstBatch.capacity).isEqualTo(5);
+
+ controller.setRequestedBatchSize(3);
+ controller.setRequestedBatchSize(2);
+ BatchResult secondBatch = readIntBatch(reader);
+ assertThat(secondBatch.values).containsExactly(5, 6);
+ assertThat(secondBatch.capacity).isEqualTo(2);
+
+ controller.setRequestedBatchSize(8);
+ BatchResult thirdBatch = readIntBatch(reader);
+ assertThat(thirdBatch.values).containsExactly(7, 8, 9, 10, 11, 12,
13, 14);
+ assertThat(thirdBatch.capacity).isEqualTo(8);
+ }
+ }
+
+ @Test
+ void testDynamicReadBatchSizeForNestedTypes() throws IOException {
+ List<InternalRow> records = prepareNestedData(20);
+ Path path = createTempParquetFileByPaimon(folder, records, 10,
NESTED_ARRAY_MAP_TYPE);
+ ParquetReaderFactory factory =
+ new ParquetReaderFactory(new Options(), NESTED_ARRAY_MAP_TYPE,
4, null);
+ ReadBatchSizeController controller = new ReadBatchSizeController(8, 5);
+ InternalRowSerializer serializer = new
InternalRowSerializer(NESTED_ARRAY_MAP_TYPE);
+ List<InternalRow> results = new ArrayList<>();
+ LocalFileIO fileIO = new LocalFileIO();
+
+ try (RecordReader<InternalRow> reader =
+ factory.createReader(
+ new FormatReaderContext(
+ fileIO, path, fileIO.getFileSize(path), null,
controller))) {
+ assertThat(readNestedBatch(reader, serializer,
results)).isEqualTo(5);
+
+ controller.setRequestedBatchSize(2);
+ assertThat(readNestedBatch(reader, serializer,
results)).isEqualTo(2);
+
+ controller.setRequestedBatchSize(8);
+ assertThat(readNestedBatch(reader, serializer,
results)).isEqualTo(8);
+ assertThat(readNestedBatch(reader, serializer,
results)).isEqualTo(8);
+ }
+
+ compareNestedRow(records, results);
+ }
+
+ @Test
+ void testStaticReadBatchSizeKeepsConfiguredVectorCapacity() throws
IOException {
+ List<InternalRow> records = new ArrayList<>();
+ for (int i = 0; i < 10; i++) {
+ records.add(newRow(i));
+ }
+ Path path = createTempParquetFileByPaimon(folder, records, 10_000,
ROW_TYPE);
+ ParquetReaderFactory factory =
+ new ParquetReaderFactory(
+ new Options(),
+ RowType.builder().field("f4", new IntType()).build(),
+ 4,
+ null);
+ LocalFileIO fileIO = new LocalFileIO();
+
+ try (RecordReader<InternalRow> reader =
+ factory.createReader(
+ new FormatReaderContext(fileIO, path,
fileIO.getFileSize(path)))) {
+ BatchResult batch = readIntBatch(reader);
+ assertThat(batch.values).containsExactly(0, 1, 2, 3);
+ assertThat(batch.capacity).isEqualTo(4);
+ }
+ }
+
+ private static BatchResult readIntBatch(RecordReader<InternalRow> reader)
throws IOException {
+ RecordReader.RecordIterator<InternalRow> batch = reader.readBatch();
+ assertThat(batch).isNotNull();
+ List<Integer> values = new ArrayList<>();
+ int capacity = -1;
+ InternalRow row;
+ while ((row = batch.next()) != null) {
+ if (capacity < 0) {
+ capacity = ((ColumnarRow)
row).batch().columns[0].getCapacity();
+ }
+ values.add(row.getInt(0));
+ }
+ batch.releaseBatch();
+ return new BatchResult(values, capacity);
+ }
+
+ private static int readNestedBatch(
+ RecordReader<InternalRow> reader,
+ InternalRowSerializer serializer,
+ List<InternalRow> results)
+ throws IOException {
+ RecordReader.RecordIterator<InternalRow> batch = reader.readBatch();
+ assertThat(batch).isNotNull();
+ int capacity = -1;
+ InternalRow row;
+ while ((row = batch.next()) != null) {
+ if (capacity < 0) {
+ capacity = ((ColumnarRow)
row).batch().columns[0].getCapacity();
+ }
+ results.add(serializer.copy(row));
+ }
+ batch.releaseBatch();
+ return capacity;
+ }
+
+ private static class BatchResult {
+
+ private final List<Integer> values;
+ private final int capacity;
+
+ private BatchResult(List<Integer> values, int capacity) {
+ this.values = values;
+ this.capacity = capacity;
+ }
+ }
+
@ParameterizedTest
@MethodSource("parameters")
void testTypesReadWithSplits(int rowGroupSize) throws IOException {