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 2d7aa36f29 [core] Refactor blob view support in KeyValueTableRead
2d7aa36f29 is described below
commit 2d7aa36f296d02aad091d096fd6b8796d4eb6bfa
Author: JingsongLi <[email protected]>
AuthorDate: Wed Aug 5 15:23:10 2026 +0800
[core] Refactor blob view support in KeyValueTableRead
---
.../paimon/table/AppendOnlyFileStoreTable.java | 1 +
.../paimon/table/PrimaryKeyFileStoreTable.java | 1 +
.../paimon/table/source/AbstractDataTableRead.java | 6 +-
.../table/source/DataEvolutionTableRead.java | 8 ++-
.../paimon/table/source/KeyValueTableRead.java | 65 +++++++++++-----------
.../source/PrimaryKeyVectorPositionReaderTest.java | 5 +-
.../flink/source/TestChangelogDataReadWrite.java | 7 ++-
7 files changed, 50 insertions(+), 43 deletions(-)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/AppendOnlyFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/AppendOnlyFileStoreTable.java
index f5f59ee70c..63710698ca 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/AppendOnlyFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/AppendOnlyFileStoreTable.java
@@ -131,6 +131,7 @@ public class AppendOnlyFileStoreTable extends
AbstractFileStoreTable {
? new DataEvolutionTableRead(
providerFactories,
schema(),
+ coreOptions(),
catalogEnvironment.dependencyReadContext(),
() -> new AppendTableRead(providerFactories, schema()))
: new AppendTableRead(providerFactories, schema());
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
index ca482ca095..f521f2e9b1 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/PrimaryKeyFileStoreTable.java
@@ -153,6 +153,7 @@ public class PrimaryKeyFileStoreTable extends
AbstractFileStoreTable {
() -> store().newRead(),
() -> store().newBatchRawFileRead(),
schema(),
+ coreOptions(),
catalogEnvironment.dependencyReadContext());
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableRead.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableRead.java
index c46e1f299f..59e8cc0666 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableRead.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableRead.java
@@ -44,7 +44,7 @@ import static
org.apache.paimon.predicate.PredicateVisitor.collectFieldNames;
public abstract class AbstractDataTableRead implements InnerTableRead {
private RowType readType;
- private boolean executeFilter = false;
+ protected boolean executeFilter = false;
private Predicate predicate;
private final TableSchema schema;
@@ -103,10 +103,6 @@ public abstract class AbstractDataTableRead implements
InnerTableRead {
return predicate;
}
- protected boolean shouldExecuteFilter() {
- return executeFilter;
- }
-
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
QueryAuthContext queryAuthContext = unwrapQueryAuthSplit(split);
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 60bd2b99fe..1efc9352cf 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
@@ -36,15 +36,18 @@ import java.util.function.Supplier;
/** A {@link TableRead} for data-evolution enabled append-only tables. */
public class DataEvolutionTableRead extends AppendTableRead {
+ private final CoreOptions options;
@Nullable private final CatalogContext catalogContext;
@Nullable private final Supplier<InnerTableRead> readFactory;
public DataEvolutionTableRead(
List<Function<SplitReadConfig, SplitReadProvider>>
providerFactories,
TableSchema schema,
+ CoreOptions options,
@Nullable CatalogContext catalogContext,
@Nullable Supplier<InnerTableRead> readFactory) {
super(providerFactories, schema);
+ this.options = options;
this.catalogContext = catalogContext;
this.readFactory = readFactory;
}
@@ -52,7 +55,6 @@ public class DataEvolutionTableRead extends AppendTableRead {
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
QueryAuthContext queryAuthContext = unwrapQueryAuthSplit(split);
- CoreOptions options = CoreOptions.fromMap(schema().options());
int[] blobViewFields =
BlobViewTableReadSupport.blobViewFieldIndexes(currentReadType(), options);
if (catalogContext != null && blobViewFields.length > 0) {
@@ -69,11 +71,11 @@ public class DataEvolutionTableRead extends AppendTableRead
{
predicate(),
topN,
limit,
- shouldExecuteFilter(),
+ executeFilter,
() -> createDataReader(queryAuthContext.split(),
queryAuthContext.authResult()),
() -> {
InnerTableRead prescanRead = readFactory.get();
- if (shouldExecuteFilter()) {
+ if (executeFilter) {
prescanRead.executeFilter();
}
return prescanRead;
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 9c02bb470f..1c6de402d0 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
@@ -49,6 +49,8 @@ import java.util.List;
import java.util.Map;
import java.util.function.Supplier;
+import static
org.apache.paimon.table.source.BlobViewTableReadSupport.blobViewFieldIndexes;
+
/**
* An abstraction layer above {@link MergeFileSplitRead} to provide reading of
{@link InternalRow}.
*/
@@ -57,6 +59,7 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
private final Supplier<MergeFileSplitRead> mergeReadSupplier;
private final Supplier<RawFileSplitRead> batchRawReadSupplier;
private final List<SplitReadProvider> readProviders;
+ private final CoreOptions options;
@Nullable private final CatalogContext catalogContext;
@Nullable private RowType readType = null;
@@ -66,21 +69,16 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
@Nullable private TopN topN = null;
@Nullable private Integer limit = null;
- public KeyValueTableRead(
- Supplier<MergeFileSplitRead> mergeReadSupplier,
- Supplier<RawFileSplitRead> batchRawReadSupplier,
- TableSchema schema) {
- this(mergeReadSupplier, batchRawReadSupplier, schema, null);
- }
-
public KeyValueTableRead(
Supplier<MergeFileSplitRead> mergeReadSupplier,
Supplier<RawFileSplitRead> batchRawReadSupplier,
TableSchema schema,
+ CoreOptions options,
@Nullable CatalogContext catalogContext) {
super(schema);
this.mergeReadSupplier = mergeReadSupplier;
this.batchRawReadSupplier = batchRawReadSupplier;
+ this.options = options;
this.catalogContext = catalogContext;
this.readProviders =
Arrays.asList(
@@ -157,46 +155,47 @@ public final class KeyValueTableRead extends
AbstractDataTableRead {
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
QueryAuthContext queryAuthContext = unwrapQueryAuthSplit(split);
RecordReader<InternalRow> reader;
- if (catalogContext != null) {
- CoreOptions options = CoreOptions.fromMap(schema().options());
- int[] blobViewFields =
-
BlobViewTableReadSupport.blobViewFieldIndexes(currentReadType(), options);
- if (blobViewFields.length > 0) {
- reader =
- BlobViewTableReadSupport.createBlobViewReader(
- catalogContext,
- queryAuthContext.split(),
- queryAuthContext.authResult(),
- blobViewFields,
- currentReadType(),
- predicate(),
- topN,
- limit,
- shouldExecuteFilter(),
- () ->
- createDataReader(
- queryAuthContext.split(),
- queryAuthContext.authResult()),
- this::createBlobViewPrescanRead);
- } else {
- reader = createDataReader(queryAuthContext.split(),
queryAuthContext.authResult());
- }
+ int[] blobViewFields = blobViewFieldIndexes(currentReadType(),
options);
+ if (catalogContext != null && blobViewFields.length > 0) {
+ reader = createReaderWithBlobView(queryAuthContext,
blobViewFields);
} else {
reader = createDataReader(queryAuthContext.split(),
queryAuthContext.authResult());
}
return LimitRecordReader.limit(reader, limit);
}
+ private RecordReader<InternalRow> createReaderWithBlobView(
+ QueryAuthContext queryAuthContext, int[] blobViewFields) throws
IOException {
+ RecordReader<InternalRow> reader;
+ reader =
+ BlobViewTableReadSupport.createBlobViewReader(
+ catalogContext,
+ queryAuthContext.split(),
+ queryAuthContext.authResult(),
+ blobViewFields,
+ currentReadType(),
+ predicate(),
+ topN,
+ limit,
+ executeFilter,
+ () ->
+ createDataReader(
+ queryAuthContext.split(),
queryAuthContext.authResult()),
+ this::createBlobViewPrescanRead);
+ return reader;
+ }
+
private InnerTableRead createBlobViewPrescanRead() {
KeyValueTableRead read =
- new KeyValueTableRead(mergeReadSupplier, batchRawReadSupplier,
schema(), null);
+ new KeyValueTableRead(
+ mergeReadSupplier, batchRawReadSupplier, schema(),
options, null);
if (ioManager != null) {
read.withIOManager(ioManager);
}
if (forceKeepDelete) {
read.forceKeepDelete();
}
- if (shouldExecuteFilter()) {
+ if (executeFilter) {
read.executeFilter();
}
return read;
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 cd2f4bd65c..c74771b2fd 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
@@ -18,6 +18,7 @@
package org.apache.paimon.table.source;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
@@ -95,7 +96,9 @@ class PrimaryKeyVectorPositionReaderTest {
new KeyValueTableRead(
() -> mock(MergeFileSplitRead.class),
() -> rawRead,
- mock(TableSchema.class));
+ mock(TableSchema.class),
+ CoreOptions.fromMap(Collections.emptyMap()),
+ null);
assertThat(tableRead.createReader(split)).isInstanceOf(PrimaryKeyIndexPositionReader.class);
verify(rawRead, never()).createReader(split);
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/TestChangelogDataReadWrite.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/TestChangelogDataReadWrite.java
index 21edabbb42..2a98e34cc9 100644
---
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/TestChangelogDataReadWrite.java
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/TestChangelogDataReadWrite.java
@@ -152,7 +152,12 @@ public class TestChangelogDataReadWrite {
FileFormatDiscover.of(options),
pathFactory,
options);
- return new KeyValueTableRead(() -> read, () -> rawFileRead, null);
+ return new KeyValueTableRead(
+ () -> read,
+ () -> rawFileRead,
+ null,
+ CoreOptions.fromMap(Collections.emptyMap()),
+ null);
}
public <T> List<DataFileMeta> writeFiles(