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 0ce6a13fe0 [core][python][spark] Optimize data evolution write column
metadata (#9574)
0ce6a13fe0 is described below
commit 0ce6a13fe0de2f31373ab40ec913184afcd3bfd7
Author: Jingsong Lee <[email protected]>
AuthorDate: Thu Sep 3 17:20:54 2026 +0800
[core][python][spark] Optimize data evolution write column metadata (#9574)
---
docs/docs/multimodal-table/data-evolution.mdx | 18 ++++
.../main/java/org/apache/paimon/CoreOptions.java | 14 +++
.../java/org/apache/paimon/schema/TableSchema.java | 37 ++++++++
.../org/apache/paimon/schema/TableSchemaTest.java | 26 ++++++
.../org/apache/paimon/append/AppendOnlyWriter.java | 71 ++------------
.../append/DedicatedFormatRollingFileWriter.java | 11 ++-
.../DataEvolutionNormalCompactTask.java | 5 +-
.../DataEvolutionGlobalIndexRefreshPlanner.java | 6 +-
.../java/org/apache/paimon/io/DataFileMeta.java | 9 +-
.../paimon/operation/BaseAppendFileStoreWrite.java | 39 +++++++-
.../operation/DataEvolutionFileStoreScan.java | 5 +-
.../paimon/operation/DataEvolutionSplitRead.java | 8 +-
.../paimon/operation/EvolutionStatsCache.java | 3 +-
.../commit/RowIdColumnConflictChecker.java | 28 ++++--
.../org/apache/paimon/table/system/FilesTable.java | 23 ++++-
.../apache/paimon/utils/DataEvolutionUtils.java | 104 ++++++++++++++++++++-
.../apache/paimon/append/AppendOnlyWriterTest.java | 8 +-
.../org/apache/paimon/append/BlobTableTest.java | 88 +++++++++++++++++
.../DedicatedFormatRollingFileWriterTest.java | 27 ++++--
...DedicatedFormatRollingFileWriterVectorTest.java | 15 ++-
.../apache/paimon/append/VectorStoreTableTest.java | 45 +++++++++
.../paimon/io/KeyValueFileReadWriteTest.java | 4 +-
.../commit/RowIdColumnConflictCheckerTest.java | 54 +++++++++++
.../apache/paimon/table/system/FilesTableTest.java | 38 ++++++++
.../pypaimon/common/options/core_options.py | 18 ++++
paimon-python/pypaimon/daft/daft_datasource.py | 19 +++-
.../pypaimon/manifest/manifest_file_manager.py | 22 ++++-
.../pypaimon/ray/row_id_conflict_rewriter.py | 21 ++++-
.../pypaimon/read/scanner/data_evolution_stats.py | 12 ++-
.../pypaimon/read/scanner/file_scanner.py | 8 +-
paimon-python/pypaimon/read/split_read.py | 7 +-
paimon-python/pypaimon/schema/table_schema.py | 59 +++++++++++-
paimon-python/pypaimon/tests/blob_table_test.py | 97 +++++++++++++++++++
.../tests/ray_row_id_conflict_rewriter_test.py | 23 +++++
paimon-python/pypaimon/tests/table_schema_test.py | 36 +++++++
paimon-python/pypaimon/tests/vector_table_test.py | 17 +++-
.../tests/write/conflict_detection_test.py | 39 ++++++++
.../pypaimon/write/commit/conflict_detection.py | 20 +++-
.../write/commit/row_id_conflict_rewriter.py | 20 +++-
paimon-python/pypaimon/write/file_store_commit.py | 7 +-
.../pypaimon/write/table_update_by_row_id.py | 2 +-
.../pypaimon/write/writer/data_vector_writer.py | 11 ++-
.../write/writer/dedicated_format_writer.py | 11 ++-
.../apache/paimon/spark/copy/CopyFilesUtil.java | 21 ++++-
.../paimon/spark/copy/ListDataFilesOperator.java | 20 +++-
...DataEvolutionCompactMergeConflictRewriter.scala | 31 ++++--
.../DataEvolutionRowIdConflictRewriter.scala | 31 ++++--
.../spark/sources/PaimonMicroBatchStream.scala | 5 +-
.../paimon/spark/copy/CopyFilesUtilTest.java | 78 ++++++++++++++++
49 files changed, 1142 insertions(+), 179 deletions(-)
diff --git a/docs/docs/multimodal-table/data-evolution.mdx
b/docs/docs/multimodal-table/data-evolution.mdx
index 6217737f45..e00f0785a3 100644
--- a/docs/docs/multimodal-table/data-evolution.mdx
+++ b/docs/docs/multimodal-table/data-evolution.mdx
@@ -117,6 +117,24 @@ commit.close()
</Tabs>
+### Compact Write-Column Metadata
+
+Tables with dedicated BLOB or vector files normally repeat the complete list of
+non-dedicated columns in every normal file's metadata. For wide tables, you can
+omit this redundant list by setting
+`data-evolution.write-cols-optimization.enabled` to `true`. A missing
+`writeCols` value then means that the file contains every non-dedicated column;
+dedicated files continue to record their columns explicitly.
+
+The option affects only new writes and is disabled by default. Readers support
+both representations without a read option. Before enabling it, upgrade every
+reader and maintenance job that accesses the table to a version that supports
+the compact representation, because older readers interpret a missing
+`writeCols` value as all table columns. Persist the option with `ALTER TABLE`
(or
+set it when creating the table) so that file schema versions record the compact
+representation. Before rolling readers back to an older version, rewrite files
+produced with this option or keep those readers upgraded.
+
## Partial Updates
You can update selected columns with Spark SQL `UPDATE` or `MERGE INTO`, the
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 666ea8fd6b..ba549628dc 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -2521,6 +2521,16 @@ public class CoreOptions implements Serializable {
.defaultValue(false)
.withDescription("Whether enable data evolution for row
tracking table.");
+ public static final ConfigOption<Boolean>
DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED =
+ key("data-evolution.write-cols-optimization.enabled")
+ .booleanType()
+ .defaultValue(false)
+ .withDescription(
+ "Whether to omit write columns from data file
metadata when a data "
+ + "evolution file contains all
non-dedicated columns. Readers "
+ + "always support the omitted metadata,
but writing it is "
+ + "disabled by default for compatibility
with older readers.");
+
public static final ConfigOption<Boolean>
DATA_EVOLUTION_NESTED_FIELD_ENABLED =
key("data-evolution.nested-field.enabled")
.booleanType()
@@ -4407,6 +4417,10 @@ public class CoreOptions implements Serializable {
return options.get(DATA_EVOLUTION_ENABLED);
}
+ public boolean dataEvolutionWriteColsOptimizationEnabled() {
+ return options.get(DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED);
+ }
+
public boolean dataEvolutionNestedFieldEnabled() {
return options.get(DATA_EVOLUTION_NESTED_FIELD_ENABLED);
}
diff --git a/paimon-api/src/main/java/org/apache/paimon/schema/TableSchema.java
b/paimon-api/src/main/java/org/apache/paimon/schema/TableSchema.java
index 4539b3470f..a4f385a396 100644
--- a/paimon-api/src/main/java/org/apache/paimon/schema/TableSchema.java
+++ b/paimon-api/src/main/java/org/apache/paimon/schema/TableSchema.java
@@ -35,9 +35,12 @@ import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
+import java.util.Set;
import java.util.stream.Collectors;
import static org.apache.paimon.CoreOptions.BUCKET_KEY;
+import static org.apache.paimon.types.BlobType.fieldNamesInBlobFile;
+import static org.apache.paimon.types.VectorType.fieldNamesInVectorFile;
/**
* Schema of a table. Unlike schema, it has more information than {@link
Schema}, including schemaId
@@ -290,6 +293,40 @@ public class TableSchema implements Serializable {
new CoreOptions(options).dataEvolutionNestedFieldEnabled()
? rowType.projectByPaths(writeCols).getFields()
: rowType.project(writeCols).getFields();
+ return copy(projectedFields);
+ }
+
+ public TableSchema dataFileSchema(@Nullable List<String> writeCols) {
+ if (writeCols != null) {
+ return project(writeCols);
+ }
+
+ // A null write-cols value normally means the full schema. For
data-evolution tables with
+ // dedicated BLOB or vector files, however, a schema which enables the
compact metadata can
+ // omit the redundant list when a normal file contains every
non-dedicated field. Dedicated
+ // files always keep their explicit write columns. Checking the schema
option preserves the
+ // legacy meaning for older schema versions which used null for a true
full-schema file.
+ CoreOptions coreOptions = CoreOptions.fromMap(options);
+ if (!coreOptions.dataEvolutionEnabled()
+ || !coreOptions.dataEvolutionWriteColsOptimizationEnabled()) {
+ return this;
+ }
+ RowType rowType = new RowType(fields);
+ Set<String> dedicatedFields =
+ new HashSet<>(fieldNamesInBlobFile(rowType,
coreOptions.blobInlineField()));
+ dedicatedFields.addAll(fieldNamesInVectorFile(rowType,
coreOptions.withVectorFormat()));
+ if (dedicatedFields.isEmpty()) {
+ return this;
+ }
+
+ List<DataField> nonDedicatedFields =
+ fields.stream()
+ .filter(field ->
!dedicatedFields.contains(field.name()))
+ .collect(Collectors.toList());
+ return copy(nonDedicatedFields);
+ }
+
+ private TableSchema copy(List<DataField> projectedFields) {
return new TableSchema(
version,
id,
diff --git
a/paimon-api/src/test/java/org/apache/paimon/schema/TableSchemaTest.java
b/paimon-api/src/test/java/org/apache/paimon/schema/TableSchemaTest.java
index 1a1fa4562a..a13c4afaac 100644
--- a/paimon-api/src/test/java/org/apache/paimon/schema/TableSchemaTest.java
+++ b/paimon-api/src/test/java/org/apache/paimon/schema/TableSchemaTest.java
@@ -73,6 +73,18 @@ class TableSchemaTest {
.containsExactly(1);
}
+ @Test
+ void testCompactNullWriteColsRequiresSchemaOption() {
+ Map<String, String> options = new HashMap<>();
+ options.put(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true");
+ TableSchema legacy = dedicatedSchema(options);
+
assertThat(legacy.dataFileSchema(null).fieldNames()).containsExactly("id",
"blob", "name");
+
+
options.put(CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED.key(),
"true");
+ TableSchema optimized = dedicatedSchema(options);
+
assertThat(optimized.dataFileSchema(null).fieldNames()).containsExactly("id",
"name");
+ }
+
private static TableSchema nestedSchema(Map<String, String> options) {
return new TableSchema(
1L,
@@ -89,4 +101,18 @@ class TableSchemaTest {
options,
"");
}
+
+ private static TableSchema dedicatedSchema(Map<String, String> options) {
+ return new TableSchema(
+ 1L,
+ Arrays.asList(
+ new DataField(1, "id", DataTypes.INT()),
+ new DataField(2, "blob", DataTypes.BLOB()),
+ new DataField(3, "name", DataTypes.STRING())),
+ 3,
+ Collections.emptyList(),
+ Collections.emptyList(),
+ options,
+ "");
+ }
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
b/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
index 98bbc618d6..750bf0b918 100644
--- a/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
+++ b/paimon-core/src/main/java/org/apache/paimon/append/AppendOnlyWriter.java
@@ -89,6 +89,7 @@ public class AppendOnlyWriter implements BatchRecordWriter,
MemoryOwner {
private final FileSource fileSource;
@Nullable private final FileFormat rowSidecarFileFormat;
@Nullable private final BlobFileContext blobContext;
+ private final boolean omitAllNonDedicatedWriteCols;
private final List<DataFileMeta> newFiles;
private final List<DataFileMeta> deletedFiles;
private final List<DataFileMeta> compactBefore;
@@ -105,69 +106,6 @@ public class AppendOnlyWriter implements
BatchRecordWriter, MemoryOwner {
private SinkWriter<InternalRow> sinkWriter;
private MemorySegmentPool memorySegmentPool;
- public AppendOnlyWriter(
- FileIO fileIO,
- @Nullable IOManager ioManager,
- long schemaId,
- FileFormat fileFormat,
- @Nullable FileFormat vectorFileFormat,
- long targetFileSize,
- long blobTargetFileSize,
- long vectorTargetFileSize,
- long targetFileRowNum,
- RowType writeSchema,
- @Nullable List<String> writeCols,
- long maxSequenceNumber,
- CompactManager compactManager,
- IOFunction<List<DataFileMeta>, RecordReaderIterator<InternalRow>>
dataFileRead,
- boolean forceCompact,
- DataFilePathFactory pathFactory,
- @Nullable CommitIncrement increment,
- boolean useWriteBuffer,
- boolean spillable,
- String fileCompression,
- CompressOptions spillCompression,
- StatsCollectorFactories statsCollectorFactories,
- MemorySize maxDiskSize,
- FileIndexOptions fileIndexOptions,
- boolean asyncFileWrite,
- boolean statsDenseStore,
- boolean dataEvolutionEnabled,
- @Nullable FileFormat rowSidecarFileFormat,
- @Nullable BlobFileContext blobContext) {
- this(
- fileIO,
- ioManager,
- schemaId,
- fileFormat,
- vectorFileFormat,
- targetFileSize,
- blobTargetFileSize,
- vectorTargetFileSize,
- targetFileRowNum,
- writeSchema,
- writeCols,
- maxSequenceNumber,
- compactManager,
- dataFileRead,
- forceCompact,
- pathFactory,
- increment,
- useWriteBuffer,
- spillable,
- fileCompression,
- spillCompression,
- statsCollectorFactories,
- maxDiskSize,
- fileIndexOptions,
- asyncFileWrite,
- statsDenseStore,
- dataEvolutionEnabled,
- rowSidecarFileFormat,
- blobContext,
- FileSource.APPEND);
- }
-
public AppendOnlyWriter(
FileIO fileIO,
@Nullable IOManager ioManager,
@@ -198,7 +136,8 @@ public class AppendOnlyWriter implements BatchRecordWriter,
MemoryOwner {
boolean dataEvolutionEnabled,
@Nullable FileFormat rowSidecarFileFormat,
@Nullable BlobFileContext blobContext,
- FileSource fileSource) {
+ FileSource fileSource,
+ boolean omitAllNonDedicatedWriteCols) {
this.fileIO = fileIO;
this.schemaId = schemaId;
this.fileFormat = fileFormat;
@@ -218,6 +157,7 @@ public class AppendOnlyWriter implements BatchRecordWriter,
MemoryOwner {
this.fileSource = fileSource;
this.rowSidecarFileFormat = dataEvolutionEnabled ?
rowSidecarFileFormat : null;
this.blobContext = blobContext;
+ this.omitAllNonDedicatedWriteCols = omitAllNonDedicatedWriteCols;
this.newFiles = new ArrayList<>();
this.deletedFiles = new ArrayList<>();
this.compactBefore = new ArrayList<>();
@@ -403,7 +343,8 @@ public class AppendOnlyWriter implements BatchRecordWriter,
MemoryOwner {
fileIndexOptions,
fileSource,
statsDenseStore,
- blobContext);
+ blobContext,
+ omitAllNonDedicatedWriteCols);
}
return new RowDataRollingFileWriter(
fileIO,
diff --git
a/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
b/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
index 9abf3635c2..aface9879a 100644
---
a/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
+++
b/paimon-core/src/main/java/org/apache/paimon/append/DedicatedFormatRollingFileWriter.java
@@ -148,7 +148,8 @@ public class DedicatedFormatRollingFileWriter
FileIndexOptions fileIndexOptions,
FileSource fileSource,
boolean statsDenseStore,
- @Nullable BlobFileContext context) {
+ @Nullable BlobFileContext context,
+ boolean omitAllNonDedicatedWriteCols) {
// Initialize basic fields
Preconditions.checkArgument(
targetFileRowNum > 0,
@@ -196,7 +197,8 @@ public class DedicatedFormatRollingFileWriter
fileIndexOptions,
fileSource,
asyncFileWrite,
- statsDenseStore);
+ statsDenseStore,
+ omitAllNonDedicatedWriteCols);
}
if (context != null) {
@@ -260,7 +262,8 @@ public class DedicatedFormatRollingFileWriter
FileIndexOptions fileIndexOptions,
FileSource fileSource,
boolean asyncFileWrite,
- boolean statsDenseStore) {
+ boolean statsDenseStore,
+ boolean omitAllNonDedicatedWriteCols) {
RowType normalRowType = new RowType(fieldsInNormalFile);
List<String> normalColumnNames = normalRowType.getFieldNames();
int[] projectionNormalFields =
writeSchema.projectIndexes(normalColumnNames);
@@ -283,7 +286,7 @@ public class DedicatedFormatRollingFileWriter
asyncFileWrite,
statsDenseStore,
pathFactory.isExternalPath(),
- normalColumnNames,
+ omitAllNonDedicatedWriteCols ? null :
normalColumnNames,
null,
null);
return new ProjectedFileWriter<>(rowDataFileWriter,
projectionNormalFields);
diff --git
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
index 2d4987d680..7fe5739744 100644
---
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
+++
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionNormalCompactTask.java
@@ -200,9 +200,6 @@ public class DataEvolutionNormalCompactTask extends
DataEvolutionCompactTask {
private static List<DataField> fileFields(
Function<Long, TableSchema> schemaLoader, DataFileMeta file) {
TableSchema fileSchema = schemaLoader.apply(file.schemaId());
- boolean nestedFieldEnabled =
- new
CoreOptions(fileSchema.options()).dataEvolutionNestedFieldEnabled();
- return org.apache.paimon.utils.DataEvolutionUtils.fileFields(
- fileSchema.fields(), file, nestedFieldEnabled);
+ return
org.apache.paimon.utils.DataEvolutionUtils.fileFields(fileSchema, file);
}
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java
index 19ce8da18e..c1d888587a 100644
---
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexRefreshPlanner.java
@@ -18,7 +18,6 @@
package org.apache.paimon.globalindex;
-import org.apache.paimon.CoreOptions;
import org.apache.paimon.Snapshot;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.index.DataEvolutionIndexSourceMeta;
@@ -254,10 +253,7 @@ public final class DataEvolutionGlobalIndexRefreshPlanner {
Pair.of(file.schemaId(), file.writeCols()),
key -> {
TableSchema fileSchema =
schemaLoader.apply(file.schemaId());
- boolean nestedFieldEnabled =
- new CoreOptions(fileSchema.options())
- .dataEvolutionNestedFieldEnabled();
- return fileFields(fileSchema.fields(), file,
nestedFieldEnabled);
+ return fileFields(fileSchema, file);
});
long[] columnSequences = file.columnMaxSequenceNumbers();
long indexedMaxSequence = Long.MIN_VALUE;
diff --git a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
index e9874917a1..a8983256e1 100644
--- a/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
+++ b/paimon-core/src/main/java/org/apache/paimon/io/DataFileMeta.java
@@ -361,6 +361,11 @@ public interface DataFileMeta {
return new Range(firstRowId, firstRowId + rowCount() - 1);
}
+ /**
+ * Columns physically stored in this file. A null value means all table
fields, except that a
+ * normal data-evolution file whose schema enables compact write-column
metadata and has
+ * dedicated BLOB or vector fields contains all non-dedicated fields.
+ */
@Nullable
List<String> writeCols();
@@ -368,8 +373,8 @@ public interface DataFileMeta {
* Maximum sequence number per physical table field after data-evolution
compaction.
*
* <p>Values follow the table-field order selected by {@link #writeCols()}
when it is non-null
- * (system fields are ignored), or the file schema field order otherwise.
A null value means
- * that only the file-level sequence range is available.
+ * (system fields are ignored), or the physical field order implied by a
null write-columns
+ * value otherwise. A null value means that only the file-level sequence
range is available.
*/
@Nullable
long[] columnMaxSequenceNumbers();
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
index fa74011c28..2e734552bd 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/BaseAppendFileStoreWrite.java
@@ -40,6 +40,7 @@ import org.apache.paimon.metrics.MetricRegistry;
import org.apache.paimon.operation.metrics.BlobFetchMetrics;
import org.apache.paimon.reader.RecordReaderIterator;
import org.apache.paimon.statistics.SimpleColStatsCollector;
+import org.apache.paimon.types.DataField;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.CommitIncrement;
import org.apache.paimon.utils.ExceptionUtils;
@@ -59,13 +60,18 @@ import javax.annotation.Nullable;
import java.io.IOException;
import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.function.Function;
import java.util.function.Supplier;
+import java.util.stream.Collectors;
import static org.apache.paimon.format.FileFormat.fileFormat;
+import static org.apache.paimon.types.BlobType.fieldNamesInBlobFile;
+import static org.apache.paimon.types.VectorType.fieldNamesInVectorFile;
import static
org.apache.paimon.utils.StatsCollectorFactories.createStatsFactories;
/** {@link FileStoreWrite} for {@link AppendOnlyFileStore}. */
@@ -86,6 +92,7 @@ public abstract class BaseAppendFileStoreWrite extends
MemoryFileStoreWrite<Inte
private @Nullable BlobFetchMetrics blobFetchMetrics;
private RowType writeType;
private @Nullable List<String> writeCols;
+ private boolean omitAllNonDedicatedWriteCols;
private FileSource fileSource = FileSource.APPEND;
private boolean forceBufferSpill = false;
@@ -116,6 +123,10 @@ public abstract class BaseAppendFileStoreWrite extends
MemoryFileStoreWrite<Inte
this.rowType = rowType;
this.writeType = rowType;
this.writeCols = null;
+ this.omitAllNonDedicatedWriteCols =
+ options.dataEvolutionEnabled()
+ && options.dataEvolutionWriteColsOptimizationEnabled()
+ &&
writesAllNonDedicatedColumns(rowType.getFieldNames(), options);
this.fileFormat = fileFormat(options);
this.pathFactory = pathFactory;
this.blobContext = BlobFileContext.create(rowType, options);
@@ -183,7 +194,8 @@ public abstract class BaseAppendFileStoreWrite extends
MemoryFileStoreWrite<Inte
options.dataEvolutionEnabled(),
rowSidecarFileFormat(),
blobContext,
- fileSource);
+ fileSource,
+ omitAllNonDedicatedWriteCols);
}
public BaseAppendFileStoreWrite withFileSource(FileSource fileSource) {
@@ -209,14 +221,35 @@ public abstract class BaseAppendFileStoreWrite extends
MemoryFileStoreWrite<Inte
if (blobContext != null) {
blobContext = blobContext.withWriteType(writeType);
}
+ this.omitAllNonDedicatedWriteCols =
+ options.dataEvolutionEnabled()
+ && options.dataEvolutionWriteColsOptimizationEnabled()
+ && writesAllNonDedicatedColumns(writeCols, options);
// optimize writeCols to null in following cases:
// writeType contains all columns (without _ROW_ID and
_SEQUENCE_NUMBER)
- if (writeCols.equals(fullNames)) {
+ if (writeCols.equals(fullNames) || omitAllNonDedicatedWriteCols) {
writeCols = null;
}
this.writeCols = writeCols;
}
+ private boolean writesAllNonDedicatedColumns(
+ List<String> writtenColumns, CoreOptions coreOptions) {
+ Set<String> dedicatedFields =
+ new HashSet<>(fieldNamesInBlobFile(rowType,
coreOptions.blobInlineField()));
+ dedicatedFields.addAll(fieldNamesInVectorFile(rowType,
coreOptions.withVectorFormat()));
+ List<String> nonDedicatedFields =
+ rowType.getFields().stream()
+ .map(DataField::name)
+ .filter(name -> !dedicatedFields.contains(name))
+ .collect(Collectors.toList());
+ List<String> writtenNonDedicatedFields =
+ writtenColumns.stream()
+ .filter(name -> !dedicatedFields.contains(name))
+ .collect(Collectors.toList());
+ return writtenNonDedicatedFields.equals(nonDedicatedFields);
+ }
+
private SimpleColStatsCollector.Factory[] statsCollectors() {
return createStatsFactories(options.statsMode(), options,
writeType.getFieldNames());
}
@@ -329,7 +362,7 @@ public abstract class BaseAppendFileStoreWrite extends
MemoryFileStoreWrite<Inte
FileSource.COMPACT,
options.asyncFileWrite(),
options.statsDenseStore(),
- rowType.equals(writeType)
+ rowType.equals(writeType) || omitAllNonDedicatedWriteCols
? null
: options.dataEvolutionNestedFieldEnabled()
? writeType.collectLeafPaths(rowType)
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionFileStoreScan.java
b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionFileStoreScan.java
index 7089d955a8..ee054d445c 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionFileStoreScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/DataEvolutionFileStoreScan.java
@@ -18,7 +18,6 @@
package org.apache.paimon.operation;
-import org.apache.paimon.CoreOptions;
import org.apache.paimon.annotation.VisibleForTesting;
import org.apache.paimon.data.BinaryArray;
import org.apache.paimon.data.InternalArray;
@@ -260,9 +259,7 @@ public class DataEvolutionFileStoreScan extends
AppendOnlyFileStoreScan {
Pair.of(entry.file().schemaId(), entry.file().writeCols()),
pair -> {
TableSchema fileSchema =
scanTableSchema(entry.file().schemaId());
- boolean nestedFieldEnabled =
- new
CoreOptions(fileSchema.options()).dataEvolutionNestedFieldEnabled();
- return fileFieldIds(fileSchema.fields(), entry.file(),
nestedFieldEnabled);
+ return fileFieldIds(fileSchema, entry.file());
});
}
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 56ab8d1885..919100ec1f 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
@@ -324,7 +324,8 @@ public class DataEvolutionSplitRead implements
SplitRead<InternalRow> {
List<RowType> bunchAvailTypes = new ArrayList<>(numBunches);
for (int i = 0; i < numBunches; i++) {
DataFileMeta first = fieldsFiles.get(i).files().get(0);
- bunchDataSchemas[i] =
schemaFetcher.apply(first.schemaId()).project(first.writeCols());
+ bunchDataSchemas[i] =
+
schemaFetcher.apply(first.schemaId()).dataFileSchema(first.writeCols());
bunchAvailTypes.add(rowTypeWithRowTracking(bunchDataSchemas[i].logicalRowType()));
}
DataEvolutionReadPlanner.DataEvolutionReadPlan plan =
@@ -724,7 +725,8 @@ public class DataEvolutionSplitRead implements
SplitRead<InternalRow> {
continue;
}
- TableSchema dataSchema =
schemaFetcher.apply(file.schemaId()).project(file.writeCols());
+ TableSchema dataSchema =
+
schemaFetcher.apply(file.schemaId()).dataFileSchema(file.writeCols());
// columns this file wrote but a newer file already won: their
values here are stale
Set<String> overwrittenCols = new HashSet<>();
for (DataField field : dataSchema.fields()) {
@@ -777,7 +779,7 @@ public class DataEvolutionSplitRead implements
SplitRead<InternalRow> {
Set<Integer> fileFieldIds = new HashSet<>();
for (DataField field :
-
schemaFetcher.apply(file.schemaId()).project(file.writeCols()).fields()) {
+
schemaFetcher.apply(file.schemaId()).dataFileSchema(file.writeCols()).fields())
{
fileFieldIds.add(field.id());
}
Set<String> written = new HashSet<>();
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/EvolutionStatsCache.java
b/paimon-core/src/main/java/org/apache/paimon/operation/EvolutionStatsCache.java
index bf1a2a384e..c7dc0833f2 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/EvolutionStatsCache.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/EvolutionStatsCache.java
@@ -57,7 +57,8 @@ class EvolutionStatsCache {
private static ProjectedFileSchema projectFileSchema(
Function<Long, TableSchema> scanTableSchema, CacheKey key) {
- TableSchema dataFileSchema =
scanTableSchema.apply(key.schemaId).project(key.writeColumns);
+ TableSchema dataFileSchema =
+
scanTableSchema.apply(key.schemaId).dataFileSchema(key.writeColumns);
TableSchema dataFileSchemaWithStats =
dataFileSchema.project(key.valueStatsColumns);
List<DataField> fields = dataFileSchema.fields();
Map<Integer, FileFieldStats> fieldStats = new HashMap<>(fields.size()
* 2);
diff --git
a/paimon-core/src/main/java/org/apache/paimon/operation/commit/RowIdColumnConflictChecker.java
b/paimon-core/src/main/java/org/apache/paimon/operation/commit/RowIdColumnConflictChecker.java
index ce32b5dabb..b81bbfd681 100644
---
a/paimon-core/src/main/java/org/apache/paimon/operation/commit/RowIdColumnConflictChecker.java
+++
b/paimon-core/src/main/java/org/apache/paimon/operation/commit/RowIdColumnConflictChecker.java
@@ -183,9 +183,10 @@ public class RowIdColumnConflictChecker implements
RowIdConflictChecker {
private boolean containsAnyWriteField(Set<Integer> fieldIds, DataFileMeta
file) {
List<String> writeCols = file.writeCols();
- // If write cols == null, it's a full-schema write
if (writeCols == null) {
- return true;
+ Set<Integer> nullWriteFieldIds = new HashSet<>();
+ fieldIdResolver.addAllFieldIds(file.schemaId(), nullWriteFieldIds);
+ return !Collections.disjoint(fieldIds, nullWriteFieldIds);
}
for (String writeCol : writeCols) {
@@ -239,6 +240,7 @@ public class RowIdColumnConflictChecker implements
RowIdConflictChecker {
private final SchemaManager schemaManager;
private final Map<Long, Map<String, Integer>> fieldIdByNameCache = new
HashMap<>();
+ private final Map<Long, List<Integer>> allFieldIdsCache = new
HashMap<>();
private TopLevelFieldIdResolver(SchemaManager schemaManager) {
this.schemaManager = schemaManager;
@@ -246,7 +248,13 @@ public class RowIdColumnConflictChecker implements
RowIdConflictChecker {
@Override
public void addAllFieldIds(long schemaId, Set<Integer> fieldIds) {
- fieldIds.addAll(fieldIdByName(schemaId).values());
+ fieldIds.addAll(
+ allFieldIdsCache.computeIfAbsent(
+ schemaId,
+ id ->
+
schemaManager.schema(id).dataFileSchema(null).fields().stream()
+ .map(DataField::id)
+ .collect(Collectors.toList())));
}
@Override
@@ -271,6 +279,7 @@ public class RowIdColumnConflictChecker implements
RowIdConflictChecker {
private final SchemaManager schemaManager;
private final Map<Long, RowType> rowTypeCache = new HashMap<>();
+ private final Map<Long, List<Integer>> allFieldIdsCache = new
HashMap<>();
private NestedFieldIdResolver(SchemaManager schemaManager) {
this.schemaManager = schemaManager;
@@ -278,9 +287,16 @@ public class RowIdColumnConflictChecker implements
RowIdConflictChecker {
@Override
public void addAllFieldIds(long schemaId, Set<Integer> fieldIds) {
- // A full-schema write touches every leaf field when nested data
evolution is in use, so
- // it conflicts with any partial sub-field write.
- collectLeafIds(rowType(schemaId).getFields(), fieldIds);
+ fieldIds.addAll(
+ allFieldIdsCache.computeIfAbsent(
+ schemaId,
+ id -> {
+ List<Integer> ids = new ArrayList<>();
+ collectLeafIds(
+
schemaManager.schema(id).dataFileSchema(null).fields(),
+ ids);
+ return ids;
+ }));
}
@Override
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/system/FilesTable.java
b/paimon-core/src/main/java/org/apache/paimon/table/system/FilesTable.java
index 3840b4dc54..3602a62168 100644
--- a/paimon-core/src/main/java/org/apache/paimon/table/system/FilesTable.java
+++ b/paimon-core/src/main/java/org/apache/paimon/table/system/FilesTable.java
@@ -410,7 +410,8 @@ public class FilesTable implements ReadonlyTable {
Function<Long, RowDataToObjectArrayConverter> keyConverters,
DataFileMeta file,
SimpleStatsEvolutions simpleStatsEvolutions) {
- StatsLazyGetter statsGetter = new StatsLazyGetter(file,
simpleStatsEvolutions);
+ StatsLazyGetter statsGetter =
+ new StatsLazyGetter(file, simpleStatsEvolutions,
schemaManager);
@SuppressWarnings("unchecked")
Supplier<Object>[] fields =
new Supplier[] {
@@ -478,19 +479,35 @@ public class FilesTable implements ReadonlyTable {
private final DataFileMeta file;
private final SimpleStatsEvolutions simpleStatsEvolutions;
+ private final SchemaManager schemaManager;
private Map<String, Long> lazyNullValueCounts;
private Map<String, Object> lazyLowerValueBounds;
private Map<String, Object> lazyUpperValueBounds;
- private StatsLazyGetter(DataFileMeta file, SimpleStatsEvolutions
simpleStatsEvolutions) {
+ private StatsLazyGetter(
+ DataFileMeta file,
+ SimpleStatsEvolutions simpleStatsEvolutions,
+ SchemaManager schemaManager) {
this.file = file;
this.simpleStatsEvolutions = simpleStatsEvolutions;
+ this.schemaManager = schemaManager;
}
private void initialize() {
+ List<String> writeCols = file.writeCols();
+ if (writeCols == null) {
+ TableSchema fileSchema = schemaManager.schema(file.schemaId());
+ List<DataField> physicalFields =
fileSchema.dataFileSchema(null).fields();
+ if (physicalFields.size() != fileSchema.fields().size()) {
+ writeCols =
+ physicalFields.stream()
+ .map(DataField::name)
+ .collect(Collectors.toList());
+ }
+ }
SimpleStatsEvolution evolution =
- simpleStatsEvolutions.getOrCreate(file.schemaId(),
file.writeCols());
+ simpleStatsEvolutions.getOrCreate(file.schemaId(),
writeCols);
// Create value stats
SimpleStatsEvolution.Result result =
evolution.evolution(file.valueStats(), file.rowCount(),
file.valueStatsCols());
diff --git
a/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java
b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java
index e1fae54ee9..329e35067b 100644
--- a/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java
+++ b/paimon-core/src/main/java/org/apache/paimon/utils/DataEvolutionUtils.java
@@ -18,7 +18,9 @@
package org.apache.paimon.utils;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.SpecialFields;
import org.apache.paimon.table.source.DataSplit;
import org.apache.paimon.types.DataField;
@@ -55,8 +57,32 @@ public class DataEvolutionUtils {
Collection<DataSplit> splits,
Function<Long, List<DataField>> schemaFieldsLoader,
Function<Long, Boolean> nestedFieldEnabledLoader) {
+ return collectWrittenColumnIds(splits, schemaFieldsLoader, null,
nestedFieldEnabledLoader);
+ }
+
+ /** Collect exact written field ids using the physical null-write-cols
schema. */
+ public static Optional<List<Integer>> collectWrittenColumnIds(
+ Collection<DataSplit> splits, Function<Long, TableSchema>
schemaLoader) {
+ Map<Long, TableSchema> schemaCache = new HashMap<>();
+ Function<Long, TableSchema> cachedSchemaLoader =
+ schemaId -> schemaCache.computeIfAbsent(schemaId,
schemaLoader);
+ return collectWrittenColumnIds(
+ splits,
+ schemaId -> cachedSchemaLoader.apply(schemaId).fields(),
+ schemaId ->
cachedSchemaLoader.apply(schemaId).dataFileSchema(null).fields(),
+ schemaId ->
+ new
CoreOptions(cachedSchemaLoader.apply(schemaId).options())
+ .dataEvolutionNestedFieldEnabled());
+ }
+
+ private static Optional<List<Integer>> collectWrittenColumnIds(
+ Collection<DataSplit> splits,
+ Function<Long, List<DataField>> schemaFieldsLoader,
+ @Nullable Function<Long, List<DataField>>
nullWriteSchemaFieldsLoader,
+ Function<Long, Boolean> nestedFieldEnabledLoader) {
Set<Integer> fieldIds = new TreeSet<>();
Map<Long, List<DataField>> schemaFieldsCache = new HashMap<>();
+ Map<Long, List<DataField>> nullWriteSchemaFieldsCache = new
HashMap<>();
Map<Long, Boolean> nestedFieldEnabledCache = new HashMap<>();
Map<Pair<Long, List<String>>, Set<Integer>> fieldIdsCache = new
HashMap<>();
try {
@@ -77,11 +103,31 @@ public class DataEvolutionUtils {
schemaId);
return loaded;
});
+ List<DataField> nullWriteSchemaFields = schemaFields;
+ if (file.writeCols() == null &&
nullWriteSchemaFieldsLoader != null) {
+ nullWriteSchemaFields =
+ nullWriteSchemaFieldsCache.computeIfAbsent(
+ file.schemaId(),
+ schemaId -> {
+ List<DataField> loaded =
+
nullWriteSchemaFieldsLoader.apply(schemaId);
+ checkArgument(
+ loaded != null,
+ "Cannot find schema
%s.",
+ schemaId);
+ return loaded;
+ });
+ }
boolean nestedFieldEnabled =
nestedFieldEnabledCache.computeIfAbsent(
file.schemaId(),
nestedFieldEnabledLoader);
fileFieldIds =
- resolveFileFieldIds(schemaFields, file,
nestedFieldEnabled, true);
+ resolveFileFieldIds(
+ schemaFields,
+ nullWriteSchemaFields,
+ file,
+ nestedFieldEnabled,
+ true);
fieldIdsCache.put(cacheKey, fileFieldIds);
}
fieldIds.addAll(fileFieldIds);
@@ -98,18 +144,30 @@ public class DataEvolutionUtils {
*/
public static Set<Integer> fileFieldIds(
List<DataField> schemaFields, DataFileMeta file, boolean
nestedFieldEnabled) {
- return resolveFileFieldIds(schemaFields, file, nestedFieldEnabled,
false);
+ return resolveFileFieldIds(schemaFields, schemaFields, file,
nestedFieldEnabled, false);
+ }
+
+ public static Set<Integer> fileFieldIds(TableSchema schema, DataFileMeta
file) {
+ boolean nestedFieldEnabled =
+ new
CoreOptions(schema.options()).dataEvolutionNestedFieldEnabled();
+ return resolveFileFieldIds(
+ schema.fields(),
+ schema.dataFileSchema(null).fields(),
+ file,
+ nestedFieldEnabled,
+ false);
}
private static Set<Integer> resolveFileFieldIds(
List<DataField> schemaFields,
+ List<DataField> nullWriteSchemaFields,
DataFileMeta file,
boolean nestedFieldEnabled,
boolean strict) {
List<String> writeCols = file.writeCols();
Set<Integer> ids = new HashSet<>();
if (writeCols == null) {
- for (DataField field : schemaFields) {
+ for (DataField field : nullWriteSchemaFields) {
ids.add(field.id());
}
return ids;
@@ -154,9 +212,47 @@ public class DataEvolutionUtils {
/** Table fields physically present in a file, in their physical write
order. */
public static List<DataField> fileFields(
List<DataField> schemaFields, DataFileMeta file, boolean
nestedFieldEnabled) {
+ return fileFields(schemaFields, schemaFields, file,
nestedFieldEnabled);
+ }
+
+ /** Table fields physically present in a file, including compact
null-write-cols metadata. */
+ public static List<DataField> fileFields(TableSchema schema, DataFileMeta
file) {
+ return fileFields(
+ schema.fields(),
+ schema.dataFileSchema(null).fields(),
+ file,
+ new
CoreOptions(schema.options()).dataEvolutionNestedFieldEnabled());
+ }
+
+ /**
+ * Returns the columns of a partial data file, materializing compact null
write-column metadata
+ * when necessary.
+ *
+ * <p>An empty optional means that a null value still denotes a true
full-schema file. This
+ * distinction is required by callers which only operate on partial
updates.
+ */
+ public static Optional<List<String>> partialFileWriteCols(
+ TableSchema schema, DataFileMeta file) {
+ if (file.writeCols() != null) {
+ return Optional.of(file.writeCols());
+ }
+
+ List<DataField> physicalFields = schema.dataFileSchema(null).fields();
+ if (physicalFields.size() == schema.fields().size()) {
+ return Optional.empty();
+ }
+ return Optional.of(
+
physicalFields.stream().map(DataField::name).collect(Collectors.toList()));
+ }
+
+ private static List<DataField> fileFields(
+ List<DataField> schemaFields,
+ List<DataField> nullWriteSchemaFields,
+ DataFileMeta file,
+ boolean nestedFieldEnabled) {
List<String> writeCols = file.writeCols();
if (writeCols == null) {
- return schemaFields;
+ return nullWriteSchemaFields;
}
if (nestedFieldEnabled) {
diff --git
a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
index 63324e7a84..5f8bd44c0b 100644
---
a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java
@@ -1092,7 +1092,9 @@ public class AppendOnlyWriterTest {
false,
context.options.dataEvolutionEnabled(),
null,
- blobContext);
+ blobContext,
+ FileSource.APPEND,
+ false);
}
private DataFileMeta writeSharedShreddingFile(AppendOnlyWriter writer,
InternalRow... rows)
@@ -1305,7 +1307,9 @@ public class AppendOnlyWriterTest {
false,
options.dataEvolutionEnabled(),
null,
- BlobFileContext.create(writeSchema, options));
+ BlobFileContext.create(writeSchema, options),
+ FileSource.APPEND,
+ false);
writer.setMemoryPool(
new HeapMemorySegmentPool(options.writeBufferSize(),
options.pageSize()));
return Pair.of(writer, compactManager.allFiles());
diff --git
a/paimon-core/src/test/java/org/apache/paimon/append/BlobTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/append/BlobTableTest.java
index ac8acc2a02..ed69827b93 100644
--- a/paimon-core/src/test/java/org/apache/paimon/append/BlobTableTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/append/BlobTableTest.java
@@ -94,6 +94,7 @@ import java.util.stream.Collectors;
import static org.apache.paimon.CoreOptions.FILE_FORMAT_PARQUET;
import static
org.apache.paimon.append.dataevolution.DataEvolutionCompactTask.TaskType.BLOB;
+import static org.apache.paimon.format.blob.BlobFileFormat.isBlobFile;
import static org.assertj.core.api.AssertionsForClassTypes.assertThat;
import static org.assertj.core.api.AssertionsForClassTypes.assertThatThrownBy;
@@ -169,6 +170,93 @@ public class BlobTableTest extends TableTestBase {
assertThat(integer.get()).isEqualTo(1000);
}
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ public void testOmitWriteColsForAllNonDedicatedColumns(boolean
optimizationEnabled)
+ throws Exception {
+ Schema.Builder schemaBuilder = Schema.newBuilder();
+ schemaBuilder.column("f0", DataTypes.INT());
+ schemaBuilder.column("f1", DataTypes.STRING());
+ schemaBuilder.column("f2", DataTypes.BLOB());
+ schemaBuilder.option(CoreOptions.TARGET_FILE_SIZE.key(), "25 MB");
+ schemaBuilder.option(CoreOptions.COMPACTION_MIN_FILE_NUM.key(), "2");
+ schemaBuilder.option(CoreOptions.ROW_TRACKING_ENABLED.key(), "true");
+ schemaBuilder.option(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true");
+ if (optimizationEnabled) {
+ schemaBuilder.option(
+
CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED.key(), "true");
+ }
+ catalog.createTable(identifier(), schemaBuilder.build(), true);
+
+ FileStoreTable table = getTableDefault();
+ writeRows(
+ table,
+ Collections.singletonList(
+ GenericRow.of(
+ 1, BinaryString.fromString("before"), new
BlobData(blobBytes))));
+
+ List<DataFileMeta> initialFiles =
+ table.store().newScan().plan().files().stream()
+ .map(ManifestEntry::file)
+ .collect(Collectors.toList());
+ DataFileMeta initialNormalFile =
+ initialFiles.stream()
+ .filter(file -> !isBlobFile(file.fileName()))
+ .findFirst()
+ .get();
+ assertThat(initialNormalFile.writeCols())
+ .isEqualTo(optimizationEnabled ? null : Arrays.asList("f0",
"f1"));
+ assertThat(
+ initialFiles.stream()
+ .filter(file -> isBlobFile(file.fileName()))
+ .findFirst()
+ .get()
+ .writeCols())
+ .isEqualTo(Collections.singletonList("f2"));
+
+ RowType normalWriteType =
table.schema().logicalRowType().project("f0", "f1");
+ BatchWriteBuilder builder = table.newBatchWriteBuilder();
+ try (BatchTableWrite write =
builder.newWrite().withWriteType(normalWriteType);
+ BatchTableCommit commit = builder.newCommit()) {
+ write.write(GenericRow.of(2, BinaryString.fromString("after")));
+ List<CommitMessage> messages = write.prepareCommit();
+ DataFileMeta updatedNormalFile =
+ ((CommitMessageImpl)
messages.get(0)).newFilesIncrement().newFiles().get(0);
+ assertThat(updatedNormalFile.writeCols())
+ .isEqualTo(optimizationEnabled ? null :
Arrays.asList("f0", "f1"));
+ assignFirstRowId(messages, 0L);
+ commit.commit(messages);
+ }
+
+ List<InternalRow> rows = new ArrayList<>();
+ InternalRowSerializer serializer = new
InternalRowSerializer(table.rowType());
+ readDefault(row -> rows.add(serializer.copy(row)));
+ assertThat(rows.size()).isEqualTo(1);
+ assertThat(rows.get(0).getInt(0)).isEqualTo(2);
+ assertThat(rows.get(0).getString(1).toString()).isEqualTo("after");
+ assertThat(rows.get(0).getBlob(2).toData()).isEqualTo(blobBytes);
+
+ if (optimizationEnabled) {
+ DataEvolutionCompactCoordinator coordinator =
+ new DataEvolutionCompactCoordinator(
+ table, false, false, table.latestSnapshot().get());
+ List<CommitMessage> compactMessages = new ArrayList<>();
+ for (DataEvolutionCompactTask task : coordinator.plan()) {
+ compactMessages.add(task.doCompact(table, commitUser));
+ }
+ assertThat(compactMessages.size()).isGreaterThan(0);
+ commitDefault(compactMessages);
+
+ List<DataFileMeta> compactedNormalFiles =
+ table.store().newScan().plan().files().stream()
+ .map(ManifestEntry::file)
+ .filter(file -> !isBlobFile(file.fileName()))
+ .collect(Collectors.toList());
+ assertThat(compactedNormalFiles.size()).isEqualTo(1);
+ assertThat(compactedNormalFiles.get(0).writeCols()).isNull();
+ }
+ }
+
@Test
public void testArrayBlobField() throws Exception {
Schema.Builder schemaBuilder = Schema.newBuilder();
diff --git
a/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterTest.java
b/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterTest.java
index f03e562f12..7e99a8479f 100644
---
a/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterTest.java
@@ -112,7 +112,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false, // statsDenseStore
- BlobFileContext.create(SCHEMA, options));
+ BlobFileContext.create(SCHEMA, options),
+ false);
}
@Test
@@ -170,7 +171,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, options));
+ BlobFileContext.create(SCHEMA, options),
+ false);
for (int i = 0; i < 11; i++) {
cappedWriter.write(
GenericRow.of(i, BinaryString.fromString("t" + i), new
BlobData(testBlobData)));
@@ -204,7 +206,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, options));
+ BlobFileContext.create(SCHEMA, options),
+ false);
cappedWriter.write(
GenericRow.of(1, BinaryString.fromString("test"), new
BlobData(testBlobData)));
@@ -264,7 +267,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(coreOptions),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, coreOptions));
+ BlobFileContext.create(SCHEMA, coreOptions),
+ false);
List<InternalRow> rows =
Arrays.asList(
@@ -338,7 +342,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false, // statsDenseStore
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
// Create large blob data that will exceed the blob target file size
byte[] largeBlobData = new byte[3 * 1024 * 1024]; // 3 MB blob data
@@ -409,7 +414,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
byte[] blobData = new byte[1024 * 1024];
new Random(321).nextBytes(blobData);
@@ -476,7 +482,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false, // statsDenseStore
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
// Create blob data that will trigger rolling
byte[] blobData = new byte[1024 * 1024]; // 1 MB blob data
@@ -557,7 +564,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false, // statsDenseStore
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
// Create blob data that will trigger rolling
byte[] blobData = new byte[1024 * 1024]; // 1 MB blob data
@@ -779,7 +787,8 @@ public class DedicatedFormatRollingFileWriterTest {
new FileIndexOptions(),
FileSource.APPEND,
false, // statsDenseStore
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
// Write data
for (int i = 0; i < 3; i++) {
diff --git
a/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterVectorTest.java
b/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterVectorTest.java
index bed8ef7b93..9f1234a4da 100644
---
a/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterVectorTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/append/DedicatedFormatRollingFileWriterVectorTest.java
@@ -122,7 +122,8 @@ public class DedicatedFormatRollingFileWriterVectorTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
}
@Test
@@ -176,7 +177,8 @@ public class DedicatedFormatRollingFileWriterVectorTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- null);
+ null,
+ false);
writer.write(GenericRow.of(1, BinaryVector.fromPrimitiveArray(new
float[VECTOR_DIM])));
writer.abort();
@@ -261,7 +263,8 @@ public class DedicatedFormatRollingFileWriterVectorTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())));
+ BlobFileContext.create(SCHEMA, new CoreOptions(new
Options())),
+ false);
List<InternalRow> rows = makeRows(2000, 1);
writer.writeBundle(new SingleUseBundleRecords(rows));
@@ -304,7 +307,8 @@ public class DedicatedFormatRollingFileWriterVectorTest {
new FileIndexOptions(coreOptions),
FileSource.APPEND,
false,
- BlobFileContext.create(SCHEMA, coreOptions));
+ BlobFileContext.create(SCHEMA, coreOptions),
+ false);
List<InternalRow> rows = makeRows(4, 10);
writer.writeBundle(new SingleUseBundleRecords(rows));
@@ -417,7 +421,8 @@ public class DedicatedFormatRollingFileWriterVectorTest {
new FileIndexOptions(),
FileSource.APPEND,
false,
- null);
+ null,
+ false);
// 100k vector-store data would create 1 normal and 3 vector-store
files
int rowNum = 100 * 1000;
diff --git
a/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
index 87023ccb24..c27f34d976 100644
---
a/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/append/VectorStoreTableTest.java
@@ -96,6 +96,43 @@ public class VectorStoreTableTest extends TableTestBase {
assertThat(counter.get()).isEqualTo(rowNum);
}
+ @Test
+ public void testOmitWriteColsForAllNonDedicatedColumns() throws Exception {
+ catalog.createTable(identifier(), schemaDefault(true), true);
+ commitDefault(writeDataDefault(1, 1));
+
+ List<DataFileMeta> files =
+ getTableDefault().store().newScan().plan().files().stream()
+ .map(ManifestEntry::file)
+ .collect(Collectors.toList());
+ assertThat(
+ files.stream()
+ .filter(file ->
!file.fileName().contains(".vector."))
+ .filter(file ->
!file.fileName().endsWith(".blob"))
+ .findFirst()
+ .get()
+ .writeCols())
+ .isNull();
+ assertThat(
+ files.stream()
+ .filter(file ->
file.fileName().endsWith(".blob"))
+ .findFirst()
+ .get()
+ .writeCols())
+ .isEqualTo(Collections.singletonList("f2"));
+ assertThat(
+ files.stream()
+ .filter(file ->
file.fileName().contains(".vector."))
+ .findFirst()
+ .get()
+ .writeCols())
+ .isEqualTo(Collections.singletonList("f3"));
+
+ AtomicInteger count = new AtomicInteger();
+ readDefault(row -> count.incrementAndGet());
+ assertThat(count.get()).isEqualTo(1);
+ }
+
@Test
public void testMultiBatch() throws Exception {
int rowNum = (RANDOM.nextInt(64) + 1) * 2;
@@ -220,6 +257,10 @@ public class VectorStoreTableTest extends TableTestBase {
@Override
protected Schema schemaDefault() {
+ return schemaDefault(false);
+ }
+
+ private Schema schemaDefault(boolean optimizeWriteCols) {
Schema.Builder schemaBuilder = Schema.newBuilder();
schemaBuilder.column("f0", DataTypes.INT());
schemaBuilder.column("f1", DataTypes.STRING());
@@ -233,6 +274,10 @@ public class VectorStoreTableTest extends TableTestBase {
schemaBuilder.option(CoreOptions.VECTOR_FIELD.key(), "f3");
schemaBuilder.option(CoreOptions.VECTOR_FILE_FORMAT.key(), "json");
schemaBuilder.option(CoreOptions.FILE_COMPRESSION.key(), "none");
+ if (optimizeWriteCols) {
+ schemaBuilder.option(
+
CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED.key(), "true");
+ }
return schemaBuilder.build();
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/io/KeyValueFileReadWriteTest.java
b/paimon-core/src/test/java/org/apache/paimon/io/KeyValueFileReadWriteTest.java
index ad8bfb512a..333d61302b 100644
---
a/paimon-core/src/test/java/org/apache/paimon/io/KeyValueFileReadWriteTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/io/KeyValueFileReadWriteTest.java
@@ -385,7 +385,9 @@ public class KeyValueFileReadWriteTest {
false,
options.dataEvolutionEnabled(),
null,
- BlobFileContext.create(schema, options));
+ BlobFileContext.create(schema, options),
+ FileSource.APPEND,
+ false);
appendOnlyWriter.setMemoryPool(
new HeapMemorySegmentPool(options.writeBufferSize(),
options.pageSize()));
appendOnlyWriter.write(
diff --git
a/paimon-core/src/test/java/org/apache/paimon/operation/commit/RowIdColumnConflictCheckerTest.java
b/paimon-core/src/test/java/org/apache/paimon/operation/commit/RowIdColumnConflictCheckerTest.java
index 8c3649900b..de09edcfc3 100644
---
a/paimon-core/src/test/java/org/apache/paimon/operation/commit/RowIdColumnConflictCheckerTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/operation/commit/RowIdColumnConflictCheckerTest.java
@@ -18,6 +18,7 @@
package org.apache.paimon.operation.commit;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.fs.Path;
import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.schema.Schema;
@@ -77,6 +78,25 @@ class RowIdColumnConflictCheckerTest {
.isTrue();
}
+ @Test
+ void testTreatsNullWriteColumnsAsAllNonDedicatedColumns() {
+ RowIdColumnConflictChecker checker = checker(file("current", 0L, 10L,
3L, null));
+
+ assertThat(checker.conflictsWith(file("historical", 0L, 10L, 3L,
Arrays.asList("blob"))))
+ .isFalse();
+ assertThat(checker.conflictsWith(file("historical", 0L, 10L, 3L,
Arrays.asList("value"))))
+ .isTrue();
+
+ checker = checker(file("current", 0L, 10L, 3L, Arrays.asList("blob")));
+ assertThat(checker.conflictsWith(file("historical", 0L, 10L, 3L,
null))).isFalse();
+
+ checker = checker(file("legacy-current", 0L, 10L, 4L, null));
+ assertThat(
+ checker.conflictsWith(
+ file("legacy-historical", 0L, 10L, 4L,
Arrays.asList("blob"))))
+ .isTrue();
+ }
+
@Test
void testMergesOverlappedDeltaRangesAndWriteColumns() {
RowIdColumnConflictChecker checker =
@@ -242,7 +262,41 @@ class RowIdColumnConflictCheckerTest {
Collections.singletonList("id"),
Collections.emptyMap(),
"")));
+ schemas.put(
+ 3L,
+ org.apache.paimon.schema.TableSchema.create(
+ 3L,
+ new Schema(
+ Arrays.asList(
+ new DataField(0, "value",
DataTypes.INT()),
+ new DataField(1, "blob",
DataTypes.BLOB())),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ dataEvolutionOptions(true),
+ "")));
+ schemas.put(
+ 4L,
+ org.apache.paimon.schema.TableSchema.create(
+ 4L,
+ new Schema(
+ Arrays.asList(
+ new DataField(0, "value",
DataTypes.INT()),
+ new DataField(1, "blob",
DataTypes.BLOB())),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.singletonMap(
+
CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true"),
+ "")));
return new TestingSchemaManager(
new Path("/tmp/row-id-column-conflict-checker-test"), schemas);
}
+
+ private static Map<String, String> dataEvolutionOptions(boolean optimized)
{
+ Map<String, String> options = new HashMap<>();
+ options.put(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true");
+ if (optimized) {
+
options.put(CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED.key(),
"true");
+ }
+ return options;
+ }
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/system/FilesTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/system/FilesTableTest.java
index 4c4e3c8859..82f5801071 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/system/FilesTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/system/FilesTableTest.java
@@ -21,6 +21,7 @@ package org.apache.paimon.table.system;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.catalog.Identifier;
import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.BlobData;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.fs.FileIO;
@@ -336,6 +337,43 @@ public class FilesTableTest extends TableTestBase {
.containsExactlyInAnyOrder("{f0=null, f1=a, f2=null}",
"{f0=null, f1=null, f2=1}");
}
+ @Test
+ public void testReadStatsWithOmittedNonDedicatedWriteCols() throws
Exception {
+ String tableName = "DataEvolutionBlobFilesTable";
+ Identifier identifier = identifier(tableName);
+ Schema schema =
+ Schema.newBuilder()
+ .column("f0", DataTypes.INT())
+ .column("blob", DataTypes.BLOB())
+ .column("f1", DataTypes.STRING())
+ .option(CoreOptions.ROW_TRACKING_ENABLED.key(), "true")
+ .option(CoreOptions.DATA_EVOLUTION_ENABLED.key(),
"true")
+ .option(
+
CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED.key(),
+ "true")
+ .build();
+ catalog.createTable(identifier, schema, true);
+
+ FileStoreTable dataEvolutionTable = getTable(identifier);
+ write(
+ dataEvolutionTable,
+ GenericRow.of(1, new BlobData(new byte[] {1, 2, 3}),
BinaryString.fromString("a")));
+
+ FilesTable dataEvolutionFilesTable =
+ (FilesTable)
+ catalog.getTable(
+ identifier(tableName + SYSTEM_TABLE_SPLITTER +
FilesTable.FILES));
+ InternalRow normalFile =
+ read(dataEvolutionFilesTable).stream()
+ .filter(row -> row.isNullAt(19))
+ .findFirst()
+ .get();
+
+ assertThat(normalFile.getString(10).toString()).isEqualTo("{blob=1,
f0=0, f1=0}");
+ assertThat(normalFile.getString(11).toString()).isEqualTo("{blob=null,
f0=1, f1=a}");
+ assertThat(normalFile.getString(12).toString()).isEqualTo("{blob=null,
f0=1, f1=a}");
+ }
+
private void setFirstRowId(List<CommitMessage> commitables, long
firstRowId) {
commitables.forEach(
c -> {
diff --git a/paimon-python/pypaimon/common/options/core_options.py
b/paimon-python/pypaimon/common/options/core_options.py
index 6e92ceb260..0e3ddb9bab 100644
--- a/paimon-python/pypaimon/common/options/core_options.py
+++ b/paimon-python/pypaimon/common/options/core_options.py
@@ -735,6 +735,18 @@ class CoreOptions:
.with_description("Whether to enable data evolution.")
)
+ DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED: ConfigOption[bool] = (
+ ConfigOptions.key("data-evolution.write-cols-optimization.enabled")
+ .boolean_type()
+ .default_value(False)
+ .with_description(
+ "Whether to omit write columns from data file metadata when a "
+ "data evolution file contains all non-dedicated columns. Readers "
+ "always support the omitted metadata, but writing it is disabled "
+ "by default for compatibility with older readers."
+ )
+ )
+
DATA_EVOLUTION_ROW_ID_CONFLICT_REWRITE_MAX_SIZE: ConfigOption[MemorySize]
= (
ConfigOptions.key("data-evolution.row-id-conflict-rewrite.max-size")
.memory_type()
@@ -1408,6 +1420,12 @@ class CoreOptions:
def data_evolution_enabled(self, default=None):
return self.options.get(CoreOptions.DATA_EVOLUTION_ENABLED, default)
+ def data_evolution_write_cols_optimization_enabled(self, default=None):
+ return self.options.get(
+ CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED,
+ default,
+ )
+
def data_evolution_row_id_conflict_rewrite_max_size(self, default=None):
value = self.options.get(
CoreOptions.DATA_EVOLUTION_ROW_ID_CONFLICT_REWRITE_MAX_SIZE,
diff --git a/paimon-python/pypaimon/daft/daft_datasource.py
b/paimon-python/pypaimon/daft/daft_datasource.py
index 663a394352..1745daa32b 100644
--- a/paimon-python/pypaimon/daft/daft_datasource.py
+++ b/paimon-python/pypaimon/daft/daft_datasource.py
@@ -488,6 +488,7 @@ def _blob_native_covering_files(
task_columns: list[str],
blob_column_names: set[str],
partition_keys: list[str],
+ schema_loader=None,
) -> list[DataFileMeta] | None:
"""Return the parquet files that can serve a blob-table split via Daft's
native reader, or ``None`` if the split must use the pypaimon fallback.
@@ -511,7 +512,13 @@ def _blob_native_covering_files(
covering: list[DataFileMeta] = []
for f in files:
name = f.file_name
- write_cols = set(f.write_cols or [])
+ if f.write_cols is None and schema_loader is not None:
+ file_schema = schema_loader(f.schema_id)
+ write_cols = {
+ field.name for field in file_schema.data_file_fields(None)
+ }
+ else:
+ write_cols = set(f.write_cols or [])
carried = write_cols & projected
if name.endswith((".blob", ".video")) or ".vector." in name:
if carried:
@@ -1011,7 +1018,15 @@ class PaimonDataSource(DataSource):
| self._map_blob_column_names
)
return _blob_native_covering_files(
- files, task_columns, blob_column_names, self._table.partition_keys
+ files,
+ task_columns,
+ blob_column_names,
+ self._table.partition_keys,
+ lambda schema_id: (
+ self._table.table_schema
+ if schema_id == self._table.table_schema.id
+ else self._table.schema_manager.get_schema(schema_id)
+ ),
)
@staticmethod
diff --git a/paimon-python/pypaimon/manifest/manifest_file_manager.py
b/paimon-python/pypaimon/manifest/manifest_file_manager.py
index 02c9e7fa86..a3ac5cfe74 100644
--- a/paimon-python/pypaimon/manifest/manifest_file_manager.py
+++ b/paimon-python/pypaimon/manifest/manifest_file_manager.py
@@ -143,10 +143,10 @@ class ManifestFileManager:
schema_id = file_dict['_SCHEMA_ID']
if schema_id == self.table.table_schema.id:
- schema_fields = self.table.table_schema.fields
+ file_schema = self.table.table_schema
else:
- schema_fields =
self.table.schema_manager.get_schema(schema_id).fields
- fields = self._get_value_stats_fields(file_dict, schema_fields)
+ file_schema = self.table.schema_manager.get_schema(schema_id)
+ fields = self._get_value_stats_fields(file_dict, file_schema)
value_dict = dict(file_dict['_VALUE_STATS'])
value_stats = SimpleStats(
min_values=BinaryRow(value_dict['_MIN_VALUES'], fields),
@@ -208,11 +208,23 @@ class ManifestFileManager:
entries.append(entry)
return entries
- def _get_value_stats_fields(self, file_dict: dict, schema_fields: list) ->
List:
+ def _get_value_stats_fields(self, file_dict: dict, file_schema) -> List:
+ # Keep accepting a field list for the focused unit tests and older
+ # internal callers. Manifest reads pass TableSchema so compact
+ # write-column metadata can be resolved from the historical schema.
+ schema_fields = (
+ file_schema.fields
+ if hasattr(file_schema, 'fields')
+ else file_schema
+ )
if file_dict['_VALUE_STATS_COLS'] is None:
if '_WRITE_COLS' in file_dict:
if file_dict['_WRITE_COLS'] is None:
- fields = schema_fields
+ fields = (
+ file_schema.data_file_fields(None)
+ if hasattr(file_schema, 'data_file_fields')
+ else schema_fields
+ )
else:
read_fields = file_dict['_WRITE_COLS']
# writeCols may contain metadata fields (e.g. _ROW_ID,
_SEQUENCE_NUMBER)
diff --git a/paimon-python/pypaimon/ray/row_id_conflict_rewriter.py
b/paimon-python/pypaimon/ray/row_id_conflict_rewriter.py
index 14cfb0ac22..8142665228 100644
--- a/paimon-python/pypaimon/ray/row_id_conflict_rewriter.py
+++ b/paimon-python/pypaimon/ray/row_id_conflict_rewriter.py
@@ -236,6 +236,7 @@ def _rewrite_updates(
candidates = [
item for item in staged
if _is_rewrite_candidate(
+ table,
item,
current_exact_ranges,
latest_snapshot.next_row_id,
@@ -265,7 +266,8 @@ def _rewrite_updates(
rewritten_messages = []
groups: Dict[Tuple[str, ...], List[_StagedFile]] = {}
for candidate in candidates:
- groups.setdefault(tuple(candidate.file.write_cols), []).append(
+ groups.setdefault(tuple(_partial_file_write_cols(
+ table, candidate.file)), []).append(
candidate
)
for columns, files in groups.items():
@@ -398,24 +400,37 @@ def _find_row_id_conflict(error) ->
Optional[RowIdExistenceConflict]:
def _is_rewrite_candidate(
+ table,
item: _StagedFile,
current_exact_ranges,
next_row_id: int,
) -> bool:
file = item.file
+ write_cols = _partial_file_write_cols(table, file)
return (
_is_normal_row_id_file(file)
and file.first_row_id < next_row_id
- and bool(file.write_cols)
+ and bool(write_cols)
and not any(
SpecialFields.is_system_field(name)
- for name in file.write_cols
+ for name in write_cols
)
and _range_key(item.message.partition, item.message.bucket, file)
not in current_exact_ranges
)
+def _partial_file_write_cols(table, file: DataFileMeta):
+ schema = (
+ table.table_schema
+ if file.schema_id == table.table_schema.id
+ else table.schema_manager.get_schema(file.schema_id)
+ )
+ if schema is None:
+ raise RuntimeError(f"Schema {file.schema_id} not found")
+ return schema.partial_file_write_cols(file.write_cols)
+
+
def _ranges_are_still_covered(current_files, candidates) -> bool:
current_ranges = {}
for split, file in current_files:
diff --git a/paimon-python/pypaimon/read/scanner/data_evolution_stats.py
b/paimon-python/pypaimon/read/scanner/data_evolution_stats.py
index c9de93763d..71c564221b 100644
--- a/paimon-python/pypaimon/read/scanner/data_evolution_stats.py
+++ b/paimon-python/pypaimon/read/scanner/data_evolution_stats.py
@@ -178,10 +178,16 @@ class DataEvolutionGroupStatsFilter:
if layout is not None:
return layout
- schema_fields = self.schema_fields(file.schema_id)
+ schema = self.schema_fields(file.schema_id)
+ schema_fields = schema.fields if hasattr(schema, 'fields') else schema
fields_by_name = {field.name: field for field in schema_fields}
- data_fields = self._project_fields(
- schema_fields, fields_by_name, file.write_cols)
+ data_fields = (
+ schema.data_file_fields(None)
+ if file.write_cols is None
+ and hasattr(schema, 'data_file_fields')
+ else self._project_fields(
+ schema_fields, fields_by_name, file.write_cols)
+ )
stats_fields = self._project_fields(
data_fields,
{field.name: field for field in data_fields},
diff --git a/paimon-python/pypaimon/read/scanner/file_scanner.py
b/paimon-python/pypaimon/read/scanner/file_scanner.py
index 7714e7e925..c114fd656e 100755
--- a/paimon-python/pypaimon/read/scanner/file_scanner.py
+++ b/paimon-python/pypaimon/read/scanner/file_scanner.py
@@ -303,6 +303,12 @@ class FileScanner:
return self.table.table_schema.fields
return self.table.schema_manager.get_schema(schema_id).fields
+ def _schema(self, schema_id: int):
+ """Resolve the complete schema including its dedicated-file options."""
+ if schema_id == self.table.table_schema.id:
+ return self.table.table_schema
+ return self.table.schema_manager.get_schema(schema_id)
+
def _deletion_files_map(self, entries: List[ManifestEntry]) -> Dict[tuple,
Dict[str, DeletionFile]]:
if not self.deletion_vectors_enabled:
return {}
@@ -501,7 +507,7 @@ class FileScanner:
group_stats_filter = DataEvolutionGroupStatsFilter(
stats_predicate,
self.table.fields,
- self._schema_fields,
+ self._schema,
)
return entries, DataEvolutionSplitGenerator(
diff --git a/paimon-python/pypaimon/read/split_read.py
b/paimon-python/pypaimon/read/split_read.py
index ad31a0103d..f00976f1f0 100644
--- a/paimon-python/pypaimon/read/split_read.py
+++ b/paimon-python/pypaimon/read/split_read.py
@@ -407,7 +407,7 @@ class SplitRead(ABC):
row_full_fields = self._create_key_value_fields(
file_schema.fields)
else:
- row_full_fields = file_schema.fields
+ row_full_fields = file_schema.data_file_fields(None)
format_reader = FormatRowReader(
self.table.file_io, file_path, read_file_fields,
row_full_fields,
@@ -1363,7 +1363,10 @@ class DataEvolutionSplitRead(SplitRead):
# the file's schema version, not the current table schema.
# The file only contains columns from when it was written.
file_schema = self._resolve_schema(first_file.schema_id)
- field_ids = [field.id for field in file_schema.fields]
+ field_ids = [
+ field.id
+ for field in file_schema.data_file_fields(None)
+ ]
field_ids.append(SpecialFields.ROW_ID.id)
field_ids.append(SpecialFields.SEQUENCE_NUMBER.id)
diff --git a/paimon-python/pypaimon/schema/table_schema.py
b/paimon-python/pypaimon/schema/table_schema.py
index d8586b28d1..604f8c0e3e 100644
--- a/paimon-python/pypaimon/schema/table_schema.py
+++ b/paimon-python/pypaimon/schema/table_schema.py
@@ -23,7 +23,12 @@ from typing import Dict, List, Optional
from pypaimon.common.options.core_options import CoreOptions
from pypaimon.common.file_io import FileIO
from pypaimon.common.json_util import json_field
-from pypaimon.schema.data_types import DataField, current_highest_field_id
+from pypaimon.schema.data_types import (
+ DataField,
+ VectorType,
+ current_highest_field_id,
+ is_blob_file_field,
+)
from pypaimon.schema.schema import Schema
@@ -100,6 +105,58 @@ class TableSchema:
field_map = {f.name: f for f in self.fields}
return [field_map[name] for name in self.bucket_keys]
+ def data_file_fields(
+ self, write_cols: Optional[List[str]]
+ ) -> List[DataField]:
+ """Return fields physically stored in a data file.
+
+ A null ``write_cols`` normally represents the full schema. For a
+ data-evolution schema which enables compact metadata and has dedicated
+ BLOB or vector fields, it represents every non-dedicated field instead.
+ Dedicated files retain explicit write columns.
+ """
+ if write_cols is not None:
+ fields_by_name = {field.name: field for field in self.fields}
+ return [
+ fields_by_name[name]
+ for name in write_cols
+ if name in fields_by_name
+ ]
+
+ core_options = CoreOptions.from_dict(self.options)
+ if (
+ not core_options.data_evolution_enabled(False)
+ or not
core_options.data_evolution_write_cols_optimization_enabled(False)
+ ):
+ return list(self.fields)
+ inline_blob_fields = (
+ core_options.blob_descriptor_fields()
+ | core_options.blob_view_fields()
+ )
+ return [
+ field
+ for field in self.fields
+ if not (
+ is_blob_file_field(field)
+ and field.name not in inline_blob_fields
+ )
+ and not (
+ core_options.with_vector_format()
+ and isinstance(field.type, VectorType)
+ )
+ ]
+
+ def partial_file_write_cols(
+ self, write_cols: Optional[List[str]]
+ ) -> Optional[List[str]]:
+ """Resolve columns for a partial file using compact metadata."""
+ if write_cols is not None:
+ return list(write_cols)
+ physical_fields = self.data_file_fields(None)
+ if len(physical_fields) == len(self.fields):
+ return None
+ return [field.name for field in physical_fields]
+
def to_schema(self) -> Schema:
return Schema(
fields=self.fields,
diff --git a/paimon-python/pypaimon/tests/blob_table_test.py
b/paimon-python/pypaimon/tests/blob_table_test.py
index 1ab02b56ea..c012f317b4 100755
--- a/paimon-python/pypaimon/tests/blob_table_test.py
+++ b/paimon-python/pypaimon/tests/blob_table_test.py
@@ -136,6 +136,103 @@ class DedicatedFormatWriterTest(unittest.TestCase):
blob_writer.close()
+ def test_omit_write_cols_for_all_non_dedicated_columns(self):
+ pa_schema = pa.schema([
+ ('id', pa.int32()),
+ ('blob_data', pa.large_binary()),
+ ('name', pa.string()),
+ ])
+ schema = Schema.from_pyarrow_schema(
+ pa_schema,
+ options={
+ 'row-tracking.enabled': 'true',
+ 'data-evolution.enabled': 'true',
+ 'data-evolution.write-cols-optimization.enabled': 'true',
+ 'metadata.stats-mode': 'full',
+ },
+ )
+ self.catalog.create_table(
+ 'test_db.optimized_write_cols', schema, False)
+ table = self.catalog.get_table('test_db.optimized_write_cols')
+
+ write_builder = table.new_batch_write_builder()
+ writer = write_builder.new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [1],
+ 'name': ['Alice'],
+ 'blob_data': [b'blob_data'],
+ }, schema=pa_schema))
+ commit_messages = writer.prepare_commit()
+ all_files = [
+ file for message in commit_messages for file in message.new_files
+ ]
+ normal_files = [
+ file for file in all_files if file.file_name.endswith('.parquet')
+ ]
+ blob_files = [
+ file for file in all_files if file.file_name.endswith('.blob')
+ ]
+ self.assertEqual(1, len(normal_files))
+ self.assertIsNone(normal_files[0].write_cols)
+ self.assertEqual([['blob_data']], [file.write_cols for file in
blob_files])
+
+ write_builder.new_commit().commit(commit_messages)
+ writer.close()
+
+ from pypaimon.manifest.manifest_file_manager import ManifestFileManager
+ from pypaimon.manifest.manifest_list_manager import ManifestListManager
+ snapshot = table.snapshot_manager().get_latest_snapshot()
+ manifests = ManifestListManager(table).read_all(snapshot)
+ committed_files = ManifestFileManager(table).read_entries_parallel(
+ manifests, drop_stats=False)
+ committed_normal = next(
+ entry.file for entry in committed_files
+ if entry.file.file_name.endswith('.parquet')
+ )
+ self.assertEqual(
+ ['id', 'name'],
+ [field.name for field in
committed_normal.value_stats.min_values.fields],
+ )
+ self.assertEqual(
+ 'Alice', committed_normal.value_stats.min_values.get_field(1)
+ )
+
+ read_builder = table.new_read_builder()
+ result = read_builder.new_read().to_arrow(
+ read_builder.new_scan().plan().splits())
+ self.assertEqual([1], result.column('id').to_pylist())
+ self.assertEqual(['Alice'], result.column('name').to_pylist())
+ self.assertEqual([b'blob_data'],
result.column('blob_data').to_pylist())
+
+ update_builder = table.new_batch_write_builder()
+ table_update = update_builder.new_update().with_update_type(
+ ['id', 'name'])
+ update_messages = table_update.update_by_arrow_with_row_id(
+ pa.Table.from_pydict({
+ '_ROW_ID': pa.array([0], type=pa.int64()),
+ 'id': pa.array([2], type=pa.int32()),
+ 'name': pa.array(['Bob'], type=pa.string()),
+ }))
+ update_files = [
+ file
+ for message in update_messages
+ for file in message.new_files
+ ]
+ self.assertTrue(update_files)
+ self.assertTrue(all(
+ file.file_name.endswith('.parquet')
+ and file.write_cols is None
+ for file in update_files
+ ))
+ update_builder.new_commit().commit(update_messages)
+
+ read_builder = table.new_read_builder()
+ result = read_builder.new_read().to_arrow(
+ read_builder.new_scan().plan().splits())
+ self.assertEqual([2], result.column('id').to_pylist())
+ self.assertEqual(['Bob'], result.column('name').to_pylist())
+ self.assertEqual([b'blob_data'],
result.column('blob_data').to_pylist())
+
def test_split_data_with_pyarrow_6_record_batch_api(self):
from pypaimon.write.writer.dedicated_format_writer import
DedicatedFormatWriter
diff --git a/paimon-python/pypaimon/tests/ray_row_id_conflict_rewriter_test.py
b/paimon-python/pypaimon/tests/ray_row_id_conflict_rewriter_test.py
index 6b2ec0d82b..2273a9c97a 100644
--- a/paimon-python/pypaimon/tests/ray_row_id_conflict_rewriter_test.py
+++ b/paimon-python/pypaimon/tests/ray_row_id_conflict_rewriter_test.py
@@ -21,13 +21,36 @@ from unittest.mock import Mock, patch
from pypaimon.ray.data_evolution_merge_into import _reraise_inner
from pypaimon.ray.row_id_conflict_rewriter import (
+ _partial_file_write_cols,
commit_self_merge_with_compaction_retry,
)
+from pypaimon.schema.data_types import AtomicType, DataField
+from pypaimon.schema.table_schema import TableSchema
from pypaimon.write.commit.conflict_detection import RowIdExistenceConflict
class RayRowIdConflictRewriterTest(unittest.TestCase):
+ def test_resolve_compact_null_write_cols_for_rewrite(self):
+ schema = TableSchema(
+ id=1,
+ fields=[
+ DataField(0, 'id', AtomicType('INT')),
+ DataField(1, 'blob', AtomicType('BLOB')),
+ DataField(2, 'name', AtomicType('STRING')),
+ ],
+ options={
+ 'data-evolution.enabled': 'true',
+ 'data-evolution.write-cols-optimization.enabled': 'true',
+ },
+ )
+ table = Mock(table_schema=schema)
+ file = Mock(schema_id=1, write_cols=None)
+
+ self.assertEqual(
+ ['id', 'name'], _partial_file_write_cols(table, file)
+ )
+
@staticmethod
def _message(snapshot_id, name):
return Mock(check_from_snapshot=snapshot_id, name=name)
diff --git a/paimon-python/pypaimon/tests/table_schema_test.py
b/paimon-python/pypaimon/tests/table_schema_test.py
index 71f426e7ca..85c2b5e425 100644
--- a/paimon-python/pypaimon/tests/table_schema_test.py
+++ b/paimon-python/pypaimon/tests/table_schema_test.py
@@ -116,5 +116,41 @@ class SchemaPrimaryKeyNullabilityTest(unittest.TestCase):
)
+class DataFileFieldsTest(unittest.TestCase):
+
+ @staticmethod
+ def _schema(optimized):
+ options = {'data-evolution.enabled': 'true'}
+ if optimized:
+ options['data-evolution.write-cols-optimization.enabled'] = 'true'
+ return TableSchema(
+ id=0,
+ fields=[
+ DataField(0, 'id', AtomicType('INT')),
+ DataField(1, 'blob', AtomicType('BLOB')),
+ DataField(2, 'name', AtomicType('STRING')),
+ ],
+ options=options,
+ )
+
+ def test_null_write_cols_uses_legacy_meaning_without_option(self):
+ schema = self._schema(False)
+ self.assertEqual(
+ ['id', 'blob', 'name'],
+ [field.name for field in schema.data_file_fields(None)],
+ )
+ self.assertIsNone(schema.partial_file_write_cols(None))
+
+ def test_null_write_cols_resolves_non_dedicated_fields_with_option(self):
+ schema = self._schema(True)
+ self.assertEqual(
+ ['id', 'name'],
+ [field.name for field in schema.data_file_fields(None)],
+ )
+ self.assertEqual(
+ ['id', 'name'], schema.partial_file_write_cols(None)
+ )
+
+
if __name__ == '__main__':
unittest.main()
diff --git a/paimon-python/pypaimon/tests/vector_table_test.py
b/paimon-python/pypaimon/tests/vector_table_test.py
index 6a4515475e..deb7f82d6c 100644
--- a/paimon-python/pypaimon/tests/vector_table_test.py
+++ b/paimon-python/pypaimon/tests/vector_table_test.py
@@ -349,6 +349,7 @@ class VectorTableWriteReadTest(unittest.TestCase):
opts = {
'row-tracking.enabled': 'true',
'data-evolution.enabled': 'true',
+ 'data-evolution.write-cols-optimization.enabled': 'true',
'vector.file.format': 'parquet',
}
s = Schema.from_pyarrow_schema(vector_schema, options=opts)
@@ -366,7 +367,21 @@ class VectorTableWriteReadTest(unittest.TestCase):
},
schema=vector_schema,
))
- wb.new_commit().commit(w.prepare_commit())
+ initial_messages = w.prepare_commit()
+ initial_files = [
+ file for message in initial_messages for file in message.new_files
+ ]
+ normal_file = next(
+ file for file in initial_files
+ if not DataFileMeta.is_vector_file(file.file_name)
+ )
+ vector_file = next(
+ file for file in initial_files
+ if DataFileMeta.is_vector_file(file.file_name)
+ )
+ self.assertIsNone(normal_file.write_cols)
+ self.assertEqual(['embedding'], vector_file.write_cols)
+ wb.new_commit().commit(initial_messages)
w.close()
from pypaimon.snapshot.snapshot import BATCH_COMMIT_IDENTIFIER
diff --git a/paimon-python/pypaimon/tests/write/conflict_detection_test.py
b/paimon-python/pypaimon/tests/write/conflict_detection_test.py
index 36b2dfc053..a82d9b0ee3 100644
--- a/paimon-python/pypaimon/tests/write/conflict_detection_test.py
+++ b/paimon-python/pypaimon/tests/write/conflict_detection_test.py
@@ -25,6 +25,7 @@ from pypaimon.manifest.index_manifest_file import
IndexManifestFile
from pypaimon.manifest.schema.data_file_meta import DataFileMeta
from pypaimon.manifest.schema.manifest_entry import ManifestEntry
from pypaimon.schema.data_types import AtomicType, DataField
+from pypaimon.schema.table_schema import TableSchema
from pypaimon.table.row.generic_row import GenericRow
from pypaimon.write.commit.conflict_detection import (
ConflictDetection,
@@ -511,6 +512,44 @@ class TestRowIdColumnConflictChecker(unittest.TestCase):
write_cols=["col_b"])
self.assertTrue(checker.conflicts_with(committed))
+ def test_null_write_cols_excludes_dedicated_fields(self):
+ schema = TableSchema(
+ id=1,
+ fields=[
+ DataField(1, "value", AtomicType("INT")),
+ DataField(2, "blob", AtomicType("BLOB")),
+ ],
+ highest_field_id=2,
+ options={
+ 'row-tracking.enabled': 'true',
+ 'data-evolution.enabled': 'true',
+ 'data-evolution.write-cols-optimization.enabled': 'true',
+ },
+ )
+ checker = self._make_checker([
+ _make_file(
+ "normal.parquet",
+ row_count=100,
+ first_row_id=0,
+ schema_id=1,
+ write_cols=None,
+ ),
+ ], schema)
+ self.assertFalse(checker.conflicts_with(_make_file(
+ "blob.blob",
+ row_count=100,
+ first_row_id=0,
+ schema_id=1,
+ write_cols=["blob"],
+ )))
+ self.assertTrue(checker.conflicts_with(_make_file(
+ "value.parquet",
+ row_count=100,
+ first_row_id=0,
+ schema_id=1,
+ write_cols=["value"],
+ )))
+
def test_no_conflict_committed_file_no_row_id(self):
delta_files = [
_make_file("d1", row_count=100, first_row_id=0,
write_cols=["col_a"]),
diff --git a/paimon-python/pypaimon/write/commit/conflict_detection.py
b/paimon-python/pypaimon/write/commit/conflict_detection.py
index 871013cf75..2792c84fca 100644
--- a/paimon-python/pypaimon/write/commit/conflict_detection.py
+++ b/paimon-python/pypaimon/write/commit/conflict_detection.py
@@ -94,7 +94,18 @@ class RowIdColumnConflictChecker:
def _contains_any_write_field(self, field_ids, file):
if file.write_cols is None:
- return True
+ schema = self._schema_manager.get_schema(file.schema_id)
+ if schema is None:
+ raise RuntimeError(f"Schema {file.schema_id} not found")
+ data_file_fields = (
+ schema.data_file_fields(None)
+ if hasattr(schema, 'data_file_fields')
+ else schema.fields
+ )
+ return any(
+ field.id in field_ids
+ for field in data_file_fields
+ )
for col_name in file.write_cols:
fid = self._field_id(file, col_name)
if fid is not None and fid in field_ids:
@@ -126,7 +137,12 @@ class RowIdColumnConflictChecker:
if file.write_cols is None:
schema = schema_manager.get_schema(file.schema_id)
if schema is not None:
- for field in schema.fields:
+ data_file_fields = (
+ schema.data_file_fields(None)
+ if hasattr(schema, 'data_file_fields')
+ else schema.fields
+ )
+ for field in data_file_fields:
if not SpecialFields.is_system_field(field.name):
field_ids.add(field.id)
else:
diff --git a/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py
b/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py
index 7eb623271e..bf361d7d62 100644
--- a/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py
+++ b/paimon-python/pypaimon/write/commit/row_id_conflict_rewriter.py
@@ -155,13 +155,14 @@ class RowIdConflictRewriter:
def _is_rewrite_candidate(
self, entry, current_exact_ranges, next_row_id):
file = entry.file
+ write_cols = self._partial_file_write_cols(file)
return (
self._is_normal_row_id_file(file)
and file.first_row_id < next_row_id
- and bool(file.write_cols)
+ and bool(write_cols)
and not any(
SpecialFields.is_system_field(name)
- for name in file.write_cols
+ for name in write_cols
)
and self._range_key(entry) not in current_exact_ranges
)
@@ -209,17 +210,26 @@ class RowIdConflictRewriter:
break
return list(affected.values())
- @staticmethod
- def _group_by_write_columns(candidates):
+ def _group_by_write_columns(self, candidates):
groups = {}
for entry in candidates:
- key = tuple(entry.file.write_cols)
+ key = tuple(self._partial_file_write_cols(entry.file))
groups.setdefault(key, []).append(entry)
return [
(list(columns), entries)
for columns, entries in groups.items()
]
+ def _partial_file_write_cols(self, file):
+ schema = (
+ self.table.table_schema
+ if file.schema_id == self.table.table_schema.id
+ else self.table.schema_manager.get_schema(file.schema_id)
+ )
+ if schema is None:
+ raise RuntimeError(f"Schema {file.schema_id} not found")
+ return schema.partial_file_write_cols(file.write_cols)
+
def _read_staged_update(self, entry, column_names):
read_fields = [self.table.field_dict[name] for name in column_names]
read_fields.append(SpecialFields.ROW_ID)
diff --git a/paimon-python/pypaimon/write/file_store_commit.py
b/paimon-python/pypaimon/write/file_store_commit.py
index 5f8f054aa8..73ecc8e950 100644
--- a/paimon-python/pypaimon/write/file_store_commit.py
+++ b/paimon-python/pypaimon/write/file_store_commit.py
@@ -291,8 +291,11 @@ class FileStoreCommit:
if msg.check_from_snapshot == -1:
continue
for f in msg.new_files:
- if f.write_cols:
- updated_cols.update(f.write_cols)
+ write_cols =
self.table.table_schema.partial_file_write_cols(
+ f.write_cols
+ )
+ if write_cols:
+ updated_cols.update(write_cols)
written_partitions.add(msg.partition)
if updated_cols:
snapshot = self.snapshot_manager.get_latest_snapshot()
diff --git a/paimon-python/pypaimon/write/table_update_by_row_id.py
b/paimon-python/pypaimon/write/table_update_by_row_id.py
index da9b1a8cdc..90c08d20d2 100644
--- a/paimon-python/pypaimon/write/table_update_by_row_id.py
+++ b/paimon-python/pypaimon/write/table_update_by_row_id.py
@@ -855,8 +855,8 @@ class TableUpdateByRowId:
# BlobWriter.prepare_commit preserves write/rolling order, which is
required
# for assigning continuous row-id ranges to rolled blob files.
for file in new_files:
- file.write_cols = file.write_cols or column_names
if DataFileMeta.is_blob_file(file.file_name):
+ file.write_cols = file.write_cols or column_names
if len(file.write_cols) != 1:
raise RuntimeError(
f"Blob update file {file.file_name} should contain "
diff --git a/paimon-python/pypaimon/write/writer/data_vector_writer.py
b/paimon-python/pypaimon/write/writer/data_vector_writer.py
index 0aa4c51810..9294f67d92 100644
--- a/paimon-python/pypaimon/write/writer/data_vector_writer.py
+++ b/paimon-python/pypaimon/write/writer/data_vector_writer.py
@@ -77,7 +77,16 @@ class DataVectorWriter(DataWriter):
self.normal_columns = [
field for field in self.table.table_schema.fields if field.name in
normal_name_set
]
- self.write_cols = self.normal_column_names
+ all_normal_column_names = [
+ col for col in all_column_names if col not in vector_set
+ ]
+ self.write_cols = (
+ None
+ if options.data_evolution_enabled(False)
+ and options.data_evolution_write_cols_optimization_enabled(False)
+ and self.normal_column_names == all_normal_column_names
+ else self.normal_column_names
+ )
self.record_count = 0
self.closed = False
diff --git a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
index 9c80ec281a..3df6db795d 100644
--- a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
+++ b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
@@ -139,7 +139,16 @@ class DedicatedFormatWriter(DataWriter):
self.normal_columns = [
field for field in self.table.table_schema.fields if field.name in
normal_name_set
]
- self.write_cols = self.normal_column_names
+ all_normal_column_names = [
+ col for col in all_column_names if col not in dedicated_set
+ ]
+ self.write_cols = (
+ None
+ if options.data_evolution_enabled(False)
+ and options.data_evolution_write_cols_optimization_enabled(False)
+ and self.normal_column_names == all_normal_column_names
+ else self.normal_column_names
+ )
# State management for blob writer
self.record_count = 0
diff --git
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java
index 88dd1fb403..d9243cf823 100644
---
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java
+++
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/CopyFilesUtil.java
@@ -25,10 +25,14 @@ import org.apache.paimon.fs.SeekableInputStream;
import org.apache.paimon.index.IndexFileMeta;
import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.io.PojoDataFileMeta;
+import org.apache.paimon.schema.TableSchema;
import org.apache.commons.io.IOUtils;
+import javax.annotation.Nullable;
+
import java.io.IOException;
+import java.util.List;
import java.util.Optional;
/** Utils for copy files. */
@@ -49,10 +53,25 @@ public class CopyFilesUtil {
public static DataFileMeta toNewDataFileMeta(
DataFileMeta oldFileMeta, String newFileName, long newSchemaId) {
+ return toNewDataFileMeta(oldFileMeta, newFileName, newSchemaId, null);
+ }
+
+ public static DataFileMeta toNewDataFileMeta(
+ DataFileMeta oldFileMeta,
+ String newFileName,
+ long newSchemaId,
+ @Nullable TableSchema sourceFileSchema) {
String newExternalPath =
externalPathDir(oldFileMeta.externalPath().orElse(null))
.map(dir -> dir + "/" + newFileName)
.orElse(null);
+ List<String> writeCols = oldFileMeta.writeCols();
+ if (sourceFileSchema != null && writeCols == null) {
+ // A null value is interpreted using the schema id. Materialize
its source meaning
+ // before rebinding the file to the target schema, whether it
denotes the full schema
+ // (legacy) or all non-dedicated columns (compact metadata).
+ writeCols = sourceFileSchema.dataFileSchema(null).fieldNames();
+ }
return new PojoDataFileMeta(
newFileName,
oldFileMeta.fileSize(),
@@ -73,7 +92,7 @@ public class CopyFilesUtil {
oldFileMeta.valueStatsCols(),
newExternalPath,
oldFileMeta.firstRowId(),
- oldFileMeta.writeCols(),
+ writeCols,
// Column sequence numbers are positional and cannot be safely
reused after
// changing the schema id. A null value makes readers fall
back conservatively.
null);
diff --git
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/ListDataFilesOperator.java
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/ListDataFilesOperator.java
index faf427165c..366fbc6cc6 100644
---
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/ListDataFilesOperator.java
+++
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/copy/ListDataFilesOperator.java
@@ -26,6 +26,7 @@ import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.io.DataFileMetaSerializer;
import org.apache.paimon.manifest.ManifestEntry;
import org.apache.paimon.partition.PartitionPredicate;
+import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.table.source.ScanMode;
import org.apache.paimon.utils.FileStorePathFactory;
@@ -37,8 +38,10 @@ import javax.annotation.Nullable;
import java.io.IOException;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
+import java.util.Map;
/** List data files. */
public class ListDataFilesOperator extends CopyFilesOperator {
@@ -70,14 +73,24 @@ public class ListDataFilesOperator extends
CopyFilesOperator {
.readFileIterator();
List<CopyFileInfo> dataFiles = new ArrayList<>();
+ Map<Long, TableSchema> sourceSchemas = new HashMap<>();
while (manifestEntries.hasNext()) {
ManifestEntry manifestEntry = manifestEntries.next();
+ DataFileMeta sourceFile = manifestEntry.file();
+ TableSchema sourceFileSchema =
+ sourceSchemas.computeIfAbsent(
+ sourceFile.schemaId(),
+ id ->
+ id == sourceTable.schema().id()
+ ? sourceTable.schema()
+ :
sourceTable.schemaManager().schema(id));
CopyFileInfo dataFile =
pickDataFiles(
manifestEntry,
sourceTable.store().pathFactory(),
targetTable.store().pathFactory(),
- targetTable.schema().id());
+ targetTable.schema().id(),
+ sourceFileSchema);
dataFiles.add(dataFile);
}
return dataFiles;
@@ -87,7 +100,8 @@ public class ListDataFilesOperator extends CopyFilesOperator
{
ManifestEntry manifestEntry,
FileStorePathFactory sourceFileStorePathFactory,
FileStorePathFactory targetFileStorePathFactory,
- long newSchemaId)
+ long newSchemaId,
+ TableSchema sourceFileSchema)
throws IOException {
Path dataFilePath =
sourceFileStorePathFactory
@@ -102,7 +116,7 @@ public class ListDataFilesOperator extends
CopyFilesOperator {
DataFileMeta fileMeta = manifestEntry.file();
DataFileMeta targetFileMeta =
CopyFilesUtil.toNewDataFileMeta(
- fileMeta, targetDataFilePath.getName(), newSchemaId);
+ fileMeta, targetDataFilePath.getName(), newSchemaId,
sourceFileSchema);
return new CopyFileInfo(
dataFilePath.toString(),
targetDataFilePath.toString(),
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
index be2d9f8c1f..047aa9457a 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionCompactMergeConflictRewriter.scala
@@ -24,6 +24,7 @@ import org.apache.paimon.data.BinaryRow
import org.apache.paimon.format.blob.BlobFileFormat.isBlobFile
import org.apache.paimon.io.{CompactIncrement, DataFileMeta, DataIncrement}
import org.apache.paimon.manifest.FileSource
+import org.apache.paimon.schema.TableSchema
import org.apache.paimon.spark.util.ScanPlanHelper
import org.apache.paimon.table.{FileStoreTable, SpecialFields}
import org.apache.paimon.table.sink.{CommitMessage, CommitMessageImpl}
@@ -31,7 +32,7 @@ import org.apache.paimon.table.source.{DataSplit,
IncrementalSplit}
import org.apache.paimon.table.source.snapshot.SnapshotReader
import org.apache.paimon.types.{RowType => PaimonRowType}
import org.apache.paimon.types.VectorType.isVectorStoreFile
-import org.apache.paimon.utils.{Range, RowRangeIndex}
+import org.apache.paimon.utils.{DataEvolutionUtils, Range, RowRangeIndex}
import org.apache.spark.sql.{functions, SparkSession}
import org.apache.spark.sql.PaimonUtils.createDataset
@@ -56,6 +57,7 @@ class DataEvolutionCompactMergeConflictRewriter(
private val partialColumns = new DataEvolutionPartialColumns(table)
private val nestedFieldEnabled =
table.coreOptions().dataEvolutionNestedFieldEnabled()
+ private val fileSchemaCache = mutable.HashMap.empty[Long, TableSchema]
def rewrite(
sparkSession: SparkSession,
@@ -128,7 +130,8 @@ class DataEvolutionCompactMergeConflictRewriter(
// write columns may be dotted sub-field paths (e.g. "nest.a"), so
they cannot be matched
// against top-level field names; collect their union instead,
ordered by the schema.
val updatedFields =
-
updatedWritePaths(files.flatMap(_.file.writeCols().asScala).distinct.toSet)
+ updatedWritePaths(
+ files.flatMap(file =>
partialFileWriteCols(file.file).get).distinct.toSet)
if (updatedFields.isEmpty) {
return JOptional.empty()
}
@@ -351,6 +354,22 @@ class DataEvolutionCompactMergeConflictRewriter(
}
}
+ private def isRegularPartialFile(file: DataFileMeta): Boolean = {
+ isNormalRowIdFile(file) &&
+ file.fileSource().orElse(null) == FileSource.APPEND &&
+ partialFileWriteCols(file).exists(
+ columns => columns.nonEmpty && columns.forall(column =>
!SpecialFields.isSystemField(column)))
+ }
+
+ private def partialFileWriteCols(file: DataFileMeta): Option[Seq[String]] = {
+ val fileSchema = fileSchemaCache.getOrElseUpdate(
+ file.schemaId(),
+ if (file.schemaId() == table.schema().id()) table.schema()
+ else table.schemaManager().schema(file.schemaId()))
+ val columns = DataEvolutionUtils.partialFileWriteCols(fileSchema, file)
+ if (columns.isPresent) Some(columns.get().asScala.toSeq) else None
+ }
+
}
private object DataEvolutionCompactMergeConflictRewriter {
@@ -459,14 +478,6 @@ private object DataEvolutionCompactMergeConflictRewriter {
file.firstRowId() != null && !isBlobFile(file.fileName()) &&
!isVectorStoreFile(file.fileName())
}
- private def isRegularPartialFile(file: DataFileMeta): Boolean = {
- isNormalRowIdFile(file) &&
- file.fileSource().orElse(null) == FileSource.APPEND &&
- file.writeCols() != null &&
- !file.writeCols().isEmpty &&
- file.writeCols().asScala.forall(column =>
!SpecialFields.isSystemField(column))
- }
-
private def quotedColumn(name: String) = {
functions.col("`" + name.replace("`", "``") + "`")
}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
index 3ac2b53a84..ebe2b07108 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/DataEvolutionRowIdConflictRewriter.scala
@@ -23,12 +23,13 @@ import org.apache.paimon.data.BinaryRow
import org.apache.paimon.format.blob.BlobFileFormat.isBlobFile
import org.apache.paimon.io.{DataFileMeta, DataIncrement}
import org.apache.paimon.operation.commit.RowIdExistenceConflictException
+import org.apache.paimon.schema.TableSchema
import org.apache.paimon.spark.util.ScanPlanHelper
import org.apache.paimon.table.{FileStoreTable, SpecialFields}
import org.apache.paimon.table.sink.{CommitMessage, CommitMessageImpl}
import org.apache.paimon.table.source.DataSplit
import org.apache.paimon.types.VectorType.isVectorStoreFile
-import org.apache.paimon.utils.{ExceptionUtils, Range, RetryWaiter}
+import org.apache.paimon.utils.{DataEvolutionUtils, ExceptionUtils, Range,
RetryWaiter}
import org.apache.spark.sql.{Row, SparkSession}
import org.apache.spark.sql.PaimonUtils.createDataset
@@ -41,6 +42,7 @@ import org.slf4j.LoggerFactory
import scala.collection.JavaConverters._
import scala.collection.immutable
+import scala.collection.mutable
/** Rebase staged partial-column files onto current row-id file boundaries. */
private[spark] class DataEvolutionRowIdConflictRewriter(
@@ -51,6 +53,7 @@ private[spark] class DataEvolutionRowIdConflictRewriter(
import DataEvolutionRowIdConflictRewriter._
private val partialColumns = new DataEvolutionPartialColumns(table)
+ private val fileSchemaCache = mutable.HashMap.empty[Long, TableSchema]
def rewrite(
sparkSession: SparkSession,
@@ -109,7 +112,7 @@ private[spark] class DataEvolutionRowIdConflictRewriter(
}
val rewrittenMessages = candidates
- .groupBy(staged => staged.file.writeCols().asScala.toSeq)
+ .groupBy(staged => partialFileWriteCols(staged.file).get)
.toSeq
.flatMap {
case (columnNames, files) =>
@@ -226,6 +229,22 @@ private[spark] class DataEvolutionRowIdConflictRewriter(
})
}
+ private def isRewriteCandidate(file: DataFileMeta, nextRowId: Long): Boolean
= {
+ isNormalRowIdFile(file) &&
+ file.firstRowId() < nextRowId &&
+ partialFileWriteCols(file).exists(
+ columns => columns.nonEmpty && columns.forall(column =>
!SpecialFields.isSystemField(column)))
+ }
+
+ private def partialFileWriteCols(file: DataFileMeta): Option[Seq[String]] = {
+ val fileSchema = fileSchemaCache.getOrElseUpdate(
+ file.schemaId(),
+ if (file.schemaId() == table.schema().id()) table.schema()
+ else table.schemaManager().schema(file.schemaId()))
+ val columns = DataEvolutionUtils.partialFileWriteCols(fileSchema, file)
+ if (columns.isPresent) Some(columns.get().asScala.toSeq) else None
+ }
+
private def withoutCandidates(
message: CommitMessageImpl,
candidates: Set[FileKey]): Option[CommitMessage] = {
@@ -268,14 +287,6 @@ private[spark] object DataEvolutionRowIdConflictRewriter {
case class RewriteResult(commitMessages: Seq[CommitMessage],
rewrittenFileCount: Int)
- private def isRewriteCandidate(file: DataFileMeta, nextRowId: Long): Boolean
= {
- isNormalRowIdFile(file) &&
- file.firstRowId() < nextRowId &&
- Option(file.writeCols()).exists(
- columns =>
- !columns.isEmpty && columns.asScala.forall(column =>
!SpecialFields.isSystemField(column)))
- }
-
private def isNormalRowIdFile(file: DataFileMeta): Boolean = {
file.firstRowId() != null && !isDedicatedFile(file)
}
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/sources/PaimonMicroBatchStream.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/sources/PaimonMicroBatchStream.scala
index 62d787bdea..1730cb6d08 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/sources/PaimonMicroBatchStream.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/sources/PaimonMicroBatchStream.scala
@@ -204,10 +204,7 @@ class PaimonMicroBatchStream(
() =>
DataEvolutionUtils.collectWrittenColumnIds(
admittedSplitSnapshot,
- schemaId => schemaLoader.apply(schemaId).fields(),
- schemaId =>
- new CoreOptions(schemaLoader.apply(schemaId).options())
- .dataEvolutionNestedFieldEnabled()
+ schemaId => schemaLoader.apply(schemaId)
)
)
}
diff --git
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java
index 2f973ffce3..e60c7bf5b3 100644
---
a/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java
+++
b/paimon-spark/paimon-spark-common/src/test/java/org/apache/paimon/spark/copy/CopyFilesUtilTest.java
@@ -18,8 +18,12 @@
package org.apache.paimon.spark.copy;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.io.DataFileMeta;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.stats.SimpleStats;
+import org.apache.paimon.types.DataTypes;
import org.junit.jupiter.api.Test;
@@ -58,4 +62,78 @@ public class CopyFilesUtilTest {
assertThat(copied.writeCols()).containsExactly("a", "b");
assertThat(copied.columnMaxSequenceNumbers()).isNull();
}
+
+ @Test
+ void testMaterializeCompactWriteColsWhenChangingSchemaId() {
+ DataFileMeta source =
+ DataFileMeta.forAppend(
+ "source.parquet",
+ 10L,
+ 2L,
+ SimpleStats.EMPTY_STATS,
+ 1L,
+ 3L,
+ 5L,
+ Collections.emptyList(),
+ null,
+ null,
+ null,
+ null,
+ null,
+ null);
+ TableSchema sourceSchema =
+ TableSchema.create(
+ 5L,
+ Schema.newBuilder()
+ .column("id", DataTypes.INT())
+ .column("blob", DataTypes.BLOB())
+ .column("name", DataTypes.STRING())
+
.option(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true")
+ .option(
+
CoreOptions.DATA_EVOLUTION_WRITE_COLS_OPTIMIZATION_ENABLED
+ .key(),
+ "true")
+ .build());
+
+ DataFileMeta copied =
+ CopyFilesUtil.toNewDataFileMeta(source, "copied.parquet", 9L,
sourceSchema);
+
+ assertThat(copied.schemaId()).isEqualTo(9L);
+ assertThat(copied.writeCols()).containsExactly("id", "name");
+ }
+
+ @Test
+ void testMaterializeLegacyFullSchemaWriteColsWhenChangingSchemaId() {
+ DataFileMeta source =
+ DataFileMeta.forAppend(
+ "source.parquet",
+ 10L,
+ 2L,
+ SimpleStats.EMPTY_STATS,
+ 1L,
+ 3L,
+ 5L,
+ Collections.emptyList(),
+ null,
+ null,
+ null,
+ null,
+ null,
+ null);
+ TableSchema sourceSchema =
+ TableSchema.create(
+ 5L,
+ Schema.newBuilder()
+ .column("id", DataTypes.INT())
+ .column("blob", DataTypes.BLOB())
+ .column("name", DataTypes.STRING())
+
.option(CoreOptions.DATA_EVOLUTION_ENABLED.key(), "true")
+ .build());
+
+ DataFileMeta copied =
+ CopyFilesUtil.toNewDataFileMeta(source, "copied.parquet", 9L,
sourceSchema);
+
+ assertThat(copied.schemaId()).isEqualTo(9L);
+ assertThat(copied.writeCols()).containsExactly("id", "blob", "name");
+ }
}