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 04a34335dc [core][python] Preserve historical field options during 
time travel (#9697)
04a34335dc is described below

commit 04a34335dcb8c8e770eee70ff4f9675e71bc39ee
Author: zhoulii <[email protected]>
AuthorDate: Fri Sep 11 13:54:39 2026 +0800

    [core][python] Preserve historical field options during time travel (#9697)
---
 .../paimon/table/AbstractFileStoreTable.java       |  67 +++++-
 .../table/TimeTravelSchemaEvolutionTest.java       | 260 +++++++++++++++++++++
 paimon-python/pypaimon/table/file_store_table.py   |  39 +++-
 .../table/time_travel_schema_evolution_test.py     | 187 +++++++++++++++
 4 files changed, 542 insertions(+), 11 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java 
b/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
index 2117f4c4dd..e9d9603f5a 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/table/AbstractFileStoreTable.java
@@ -30,7 +30,9 @@ import org.apache.paimon.manifest.IndexManifestEntry;
 import org.apache.paimon.manifest.ManifestEntry;
 import org.apache.paimon.manifest.ManifestFileMeta;
 import org.apache.paimon.operation.FileStoreScan;
+import org.apache.paimon.options.ConfigOption;
 import org.apache.paimon.options.ExpireConfig;
+import org.apache.paimon.options.FallbackKey;
 import org.apache.paimon.options.Options;
 import org.apache.paimon.predicate.Predicate;
 import org.apache.paimon.schema.FileSystemSchemaManager;
@@ -76,11 +78,15 @@ import java.io.FileNotFoundException;
 import java.io.IOException;
 import java.io.UncheckedIOException;
 import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.HashMap;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
+import java.util.Set;
 import java.util.SortedMap;
 import java.util.function.BiConsumer;
 import java.util.function.LongConsumer;
@@ -97,6 +103,9 @@ abstract class AbstractFileStoreTable implements 
FileStoreTable {
     protected final TableSchema tableSchema;
     protected final CatalogEnvironment catalogEnvironment;
 
+    // Track explicit copy() keys, including removals, separately from 
persisted schema options.
+    @Nullable private Set<String> appliedDynamicOptionKeys;
+
     @Nullable protected transient SegmentsCache<Path> manifestCache;
     @Nullable protected transient Cache<Path, Snapshot> snapshotCache;
     @Nullable protected transient Cache<String, Statistics> statsCache;
@@ -363,9 +372,15 @@ abstract class AbstractFileStoreTable implements 
FileStoreTable {
         // copy a new table schema to contain dynamic options
         TableSchema newTableSchema = tableSchema.copy(newOptions.toMap());
 
+        Set<String> mergedDynamicOptionKeys = new 
HashSet<>(dynamicOptions.keySet());
+        if (appliedDynamicOptionKeys != null) {
+            mergedDynamicOptionKeys.addAll(appliedDynamicOptionKeys);
+        }
+
         if (tryTimeTravel) {
             // see if merged options contain time travel option
-            newTableSchema = tryTimeTravel(newOptions).orElse(newTableSchema);
+            newTableSchema =
+                    tryTimeTravel(newOptions, 
mergedDynamicOptionKeys).orElse(newTableSchema);
         }
 
         // validate schema with new options
@@ -380,7 +395,11 @@ abstract class AbstractFileStoreTable implements 
FileStoreTable {
                     () -> schemaManager().listAll(), new 
CoreOptions(newTableSchema.options()));
         }
 
-        return copy(newTableSchema);
+        FileStoreTable copied = copy(newTableSchema);
+        if (copied instanceof AbstractFileStoreTable) {
+            ((AbstractFileStoreTable) copied).appliedDynamicOptionKeys = 
mergedDynamicOptionKeys;
+        }
+        return copied;
     }
 
     @Override
@@ -406,6 +425,7 @@ abstract class AbstractFileStoreTable implements 
FileStoreTable {
                                 fileIO, path, newTableSchema, 
catalogEnvironment)
                         : new PrimaryKeyFileStoreTable(
                                 fileIO, path, newTableSchema, 
catalogEnvironment);
+        copied.appliedDynamicOptionKeys = appliedDynamicOptionKeys;
         if (snapshotCache != null) {
             copied.setSnapshotCache(snapshotCache);
         }
@@ -518,7 +538,7 @@ abstract class AbstractFileStoreTable implements 
FileStoreTable {
         return snapshotExpire;
     }
 
-    private Optional<TableSchema> tryTimeTravel(Options options) {
+    private Optional<TableSchema> tryTimeTravel(Options options, Set<String> 
dynamicOptionKeys) {
         Snapshot snapshot;
         try {
             snapshot =
@@ -530,7 +550,46 @@ abstract class AbstractFileStoreTable implements 
FileStoreTable {
         if (snapshot == null) {
             return Optional.empty();
         }
-        return 
Optional.of(schemaManager().schema(snapshot.schemaId()).copy(options.toMap()));
+        TableSchema historicalSchema = 
schemaManager().schema(snapshot.schemaId());
+        return Optional.of(
+                historicalSchema.copy(
+                        excludeCurrentSchemaFieldOptions(
+                                historicalSchema, options, 
dynamicOptionKeys)));
+    }
+
+    /** Prevents current column declarations from overriding a historical 
schema's field options. */
+    private static Map<String, String> excludeCurrentSchemaFieldOptions(
+            TableSchema historicalSchema, Options options, Set<String> 
dynamicOptionKeys) {
+        // Keep scan and runtime options. Only these directive-managed column 
declarations
+        // must follow the historical schema, since columns may have been 
added or dropped.
+        Map<String, String> historicalOptions = new HashMap<>(options.toMap());
+        for (ConfigOption<String> option :
+                Arrays.asList(
+                        CoreOptions.VECTOR_FIELD,
+                        CoreOptions.BLOB_FIELD,
+                        CoreOptions.BLOB_DESCRIPTOR_FIELD,
+                        CoreOptions.BLOB_VIEW_FIELD)) {
+            // Restore the canonical key and aliases together, or a stale 
alias may take effect
+            // when the historical schema has no canonical value.
+            List<String> keys = new ArrayList<>();
+            keys.add(option.key());
+            for (FallbackKey fallback : option.fallbackKeys()) {
+                keys.add(fallback.getKey());
+            }
+            if (keys.stream().anyMatch(dynamicOptionKeys::contains)) {
+                // Preserve explicit overrides; invalid values must still fail 
schema validation.
+                continue;
+            }
+            for (String key : keys) {
+                String value = historicalSchema.options().get(key);
+                if (value == null) {
+                    historicalOptions.remove(key);
+                } else {
+                    historicalOptions.put(key, value);
+                }
+            }
+        }
+        return historicalOptions;
     }
 
     @Override
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/TimeTravelSchemaEvolutionTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/TimeTravelSchemaEvolutionTest.java
new file mode 100644
index 0000000000..fd567b6415
--- /dev/null
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/TimeTravelSchemaEvolutionTest.java
@@ -0,0 +1,260 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.table;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.Blob;
+import org.apache.paimon.data.BlobDescriptor;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaChange;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.DataTypes;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.nio.file.Files;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Tests for time travel across changes to directive-managed columns. */
+public class TimeTravelSchemaEvolutionTest extends TableTestBase {
+
+    @ParameterizedTest
+    @MethodSource("columnDirectives")
+    public void testTimeTravelAfterAddingColumn(String directive, String 
optionKey)
+            throws Exception {
+        catalog.createTable(identifier(), schemaBuilder(directive, 
true).build(), false);
+        FileStoreTable table = getTableDefault();
+        write(table, GenericRow.of(1, null));
+        table.createTag("before_add", 1);
+        TableSchema historicalSchema = 
table.schemaManager().schema(table.schema().id());
+
+        // A second snapshot shares the old schema, so checking only schema ID 
is insufficient.
+        write(table, GenericRow.of(2, null));
+        catalog.alterTable(
+                identifier(),
+                SchemaChange.addColumn("payload_v2", sourceType(directive), 
directive, null),
+                false);
+        table = getTableDefault();
+        write(table, GenericRow.of(3, null, null));
+        TableSchema latestSchema = table.schema();
+        table = 
table.copy(Collections.singletonMap(CoreOptions.READ_BATCH_SIZE.key(), "32"));
+        FileStoreTable historicalTable =
+                
table.copy(Collections.singletonMap(CoreOptions.SCAN_VERSION.key(), 
"before_add"));
+        assertHistoricalTable(historicalTable, historicalSchema);
+        assertThat(historicalTable.options())
+                .containsEntry(optionKey, "payload")
+                .containsEntry(CoreOptions.READ_BATCH_SIZE.key(), "32")
+                .containsEntry(CoreOptions.SCAN_TAG_NAME.key(), "before_add");
+
+        assertThat(read(table, new int[] {0}))
+                .extracting(row -> row.getInt(0))
+                .containsExactlyInAnyOrder(1, 2, 3);
+        assertThat(getTableDefault().schema()).isEqualTo(latestSchema);
+        
assertThat(table.schemaManager().schema(historicalSchema.id())).isEqualTo(historicalSchema);
+    }
+
+    @Test
+    public void testTimeTravelBeforeFirstVectorColumn() throws Exception {
+        String directive = "__VECTOR_FIELD;3";
+        catalog.createTable(identifier(), schemaBuilder(directive, 
false).build(), false);
+        FileStoreTable table = getTableDefault();
+        write(table, GenericRow.of(1));
+        TableSchema historicalSchema = table.schema();
+        catalog.alterTable(
+                identifier(),
+                SchemaChange.addColumn("payload", sourceType(directive), 
directive, null),
+                false);
+        table = getTableDefault();
+        write(table, GenericRow.of(2, null));
+
+        FileStoreTable historicalTable =
+                
table.copy(Collections.singletonMap(CoreOptions.SCAN_SNAPSHOT_ID.key(), "1"));
+        assertHistoricalTable(historicalTable, historicalSchema);
+        
assertThat(historicalTable.options()).doesNotContainKey(CoreOptions.VECTOR_FIELD.key());
+    }
+
+    @Test
+    public void testTimeTravelAfterDroppingDescriptorColumn() throws Exception 
{
+        String directive = "__BLOB_DESCRIPTOR_FIELD";
+        catalog.createTable(identifier(), schemaBuilder(directive, 
true).build(), false);
+        FileStoreTable table = getTableDefault();
+        byte[] bytes = new byte[] {1, 2, 3};
+        java.nio.file.Path externalFile = tempPath.resolve("payload.bin");
+        Files.write(externalFile, bytes);
+        BlobDescriptor descriptor = new 
BlobDescriptor(externalFile.toString(), 0, bytes.length);
+        write(
+                table,
+                GenericRow.of(1, Blob.fromFile(table.fileIO(), 
descriptor.uri(), 0, bytes.length)));
+        TableSchema historicalSchema = table.schema();
+        catalog.alterTable(identifier(), SchemaChange.dropColumn("payload"), 
false);
+        table = getTableDefault();
+        write(table, GenericRow.of(2));
+        
assertThat(table.options()).doesNotContainKey(CoreOptions.BLOB_DESCRIPTOR_FIELD.key());
+
+        FileStoreTable historicalTable =
+                
table.copy(Collections.singletonMap(CoreOptions.SCAN_SNAPSHOT_ID.key(), "1"));
+        assertHistoricalTable(historicalTable, historicalSchema);
+        assertThat(historicalTable.options())
+                .containsEntry(CoreOptions.BLOB_DESCRIPTOR_FIELD.key(), 
"payload");
+        // The old column must still be read as an inline descriptor, not a 
managed .blob column.
+        List<InternalRow> rows = read(historicalTable);
+        assertThat(rows).hasSize(1);
+        
assertThat(rows.get(0).getBlob(1).toDescriptor()).isEqualTo(descriptor);
+        assertThat(rows.get(0).getBlob(1).toData()).isEqualTo(bytes);
+    }
+
+    @Test
+    public void testTimeTravelWithLegacyDescriptorOption() throws Exception {
+        String legacyKey = "blob.stored-descriptor-fields";
+        Schema schema =
+                schemaBuilder("__BLOB_DESCRIPTOR_FIELD", false)
+                        .column("payload", DataTypes.BLOB())
+                        .option(legacyKey, "payload")
+                        .build();
+        catalog.createTable(identifier(), schema, false);
+        FileStoreTable table = getTableDefault();
+        write(table, GenericRow.of(1, null));
+        TableSchema historicalSchema = table.schema();
+        catalog.alterTable(
+                identifier(),
+                SchemaChange.addColumn(
+                        "payload_v2", DataTypes.BYTES(), 
"__BLOB_DESCRIPTOR_FIELD", null),
+                false);
+        table = getTableDefault();
+        assertThat(table.options())
+                .containsEntry(CoreOptions.BLOB_DESCRIPTOR_FIELD.key(), 
"payload,payload_v2")
+                .doesNotContainKey(legacyKey);
+
+        FileStoreTable historicalTable =
+                
table.copy(Collections.singletonMap(CoreOptions.SCAN_SNAPSHOT_ID.key(), "1"));
+        assertHistoricalTable(historicalTable, historicalSchema);
+        assertThat(historicalTable.options())
+                .containsEntry(legacyKey, "payload")
+                .doesNotContainKey(CoreOptions.BLOB_DESCRIPTOR_FIELD.key());
+        
assertThat(historicalTable.coreOptions().blobDescriptorField()).containsExactly("payload");
+    }
+
+    @Test
+    public void testTimeTravelBeforeLegacyDescriptorOption() throws Exception {
+        catalog.createTable(
+                identifier(), schemaBuilder("__BLOB_DESCRIPTOR_FIELD", 
false).build(), false);
+        FileStoreTable table = getTableDefault();
+        write(table, GenericRow.of(1));
+        TableSchema historicalSchema = table.schema();
+        catalog.alterTable(
+                identifier(), SchemaChange.addColumn("payload", 
DataTypes.BLOB()), false);
+        catalog.alterTable(
+                identifier(),
+                SchemaChange.setOption("blob.stored-descriptor-fields", 
"payload"),
+                false);
+        table = getTableDefault();
+        FileStoreTable historicalTable =
+                
table.copy(Collections.singletonMap(CoreOptions.SCAN_SNAPSHOT_ID.key(), "1"));
+        assertHistoricalTable(historicalTable, historicalSchema);
+        
assertThat(historicalTable.options()).doesNotContainKey("blob.stored-descriptor-fields");
+        
assertThat(historicalTable.coreOptions().blobDescriptorField()).isEmpty();
+    }
+
+    @Test
+    public void testExplicitVectorOptionsArePreserved() throws Exception {
+        String directive = "__VECTOR_FIELD;3";
+        catalog.createTable(identifier(), schemaBuilder(directive, 
true).build(), false);
+        FileStoreTable table = getTableDefault();
+        write(table, GenericRow.of(1, null));
+
+        Map<String, String> queryOptions = new HashMap<>();
+        queryOptions.put(CoreOptions.SCAN_SNAPSHOT_ID.key(), "1");
+        queryOptions.put(CoreOptions.VECTOR_FIELD.key(), "missing");
+        assertThatThrownBy(() -> table.copy(queryOptions))
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessage("Some of the columns specified as vector-field are 
unknown.");
+
+        queryOptions.put(CoreOptions.VECTOR_FIELD.key(), null);
+        FileStoreTable historicalTable = table.copy(queryOptions);
+        
assertThat(historicalTable.options()).doesNotContainKey(CoreOptions.VECTOR_FIELD.key());
+
+        // A later copy must not restore a field option explicitly removed by 
an earlier copy.
+        Map<String, String> readOptions =
+                Collections.singletonMap(CoreOptions.READ_BATCH_SIZE.key(), 
"32");
+        assertThat(historicalTable.copy(readOptions).options())
+                .doesNotContainKey(CoreOptions.VECTOR_FIELD.key());
+        assertThat(
+                        table.copyWithoutTimeTravel(
+                                        Collections.singletonMap(
+                                                
CoreOptions.VECTOR_FIELD.key(), null))
+                                .copy(
+                                        Collections.singletonMap(
+                                                
CoreOptions.SCAN_SNAPSHOT_ID.key(), "1"))
+                                .options())
+                .doesNotContainKey(CoreOptions.VECTOR_FIELD.key());
+    }
+
+    private void assertHistoricalTable(FileStoreTable historicalTable, 
TableSchema historicalSchema)
+            throws Exception {
+        
assertThat(historicalTable.schema().id()).isEqualTo(historicalSchema.id());
+        
assertThat(historicalTable.schema().fields()).isEqualTo(historicalSchema.fields());
+        assertThat(read(historicalTable, new int[] {0}))
+                .extracting(row -> row.getInt(0))
+                .containsExactly(1);
+    }
+
+    private static Schema.Builder schemaBuilder(String directive, boolean 
withPayload) {
+        Schema.Builder builder =
+                Schema.newBuilder()
+                        .column("id", DataTypes.INT())
+                        .option(CoreOptions.FILE_FORMAT.key(), "parquet")
+                        .option(CoreOptions.FILE_COMPRESSION.key(), "none")
+                        .option(CoreOptions.ROW_TRACKING_ENABLED.key(), "true")
+                        .option(CoreOptions.DATA_EVOLUTION_ENABLED.key(), 
"true");
+        if (withPayload) {
+            builder.column("payload", sourceType(directive), directive);
+        }
+        if (directive.startsWith("__VECTOR")) {
+            builder.option(CoreOptions.VECTOR_FILE_FORMAT.key(), "json");
+        }
+        return builder;
+    }
+
+    private static Stream<Arguments> columnDirectives() {
+        return Stream.of(
+                Arguments.of("__VECTOR_FIELD;3", 
CoreOptions.VECTOR_FIELD.key()),
+                Arguments.of("__BLOB_FIELD", CoreOptions.BLOB_FIELD.key()),
+                Arguments.of("__BLOB_DESCRIPTOR_FIELD", 
CoreOptions.BLOB_DESCRIPTOR_FIELD.key()),
+                Arguments.of("__BLOB_VIEW_FIELD", 
CoreOptions.BLOB_VIEW_FIELD.key()));
+    }
+
+    private static DataType sourceType(String directive) {
+        return directive.startsWith("__VECTOR")
+                ? DataTypes.ARRAY(DataTypes.FLOAT())
+                : DataTypes.BYTES();
+    }
+}
diff --git a/paimon-python/pypaimon/table/file_store_table.py 
b/paimon-python/pypaimon/table/file_store_table.py
index 2cc612ae80..d6700c57cb 100644
--- a/paimon-python/pypaimon/table/file_store_table.py
+++ b/paimon-python/pypaimon/table/file_store_table.py
@@ -15,7 +15,7 @@
 # specific language governing permissions and limitations
 # under the License.
 
-from typing import List, Optional
+from typing import List, Optional, Set
 
 from pypaimon.catalog.catalog_environment import CatalogEnvironment
 from pypaimon.common.file_io import FileIO
@@ -527,8 +527,11 @@ class FileStoreTable(Table):
 
         new_table_schema = self.table_schema.copy(new_options=new_options)
 
+        # Cumulative copy() overrides (removals kept as None) vs the on-disk 
schema.
+        applied_options = {**getattr(self, '_applied_dynamic_options', {}), 
**options}
+
         if resolve_time_travel:
-            time_travel_schema = self._try_time_travel(Options(new_options))
+            time_travel_schema = self._try_time_travel(Options(new_options), 
set(applied_options))
             if time_travel_schema is not None:
                 new_table_schema = time_travel_schema
 
@@ -549,12 +552,10 @@ class FileStoreTable(Table):
 
         new_table = FileStoreTable(self.file_io, new_identifier, 
self.table_path,
                                    new_table_schema, catalog_env)
-        # Cumulative copy() overrides (removals kept as None) vs the on-disk 
schema.
-        new_table._applied_dynamic_options = {
-            **getattr(self, '_applied_dynamic_options', {}), **options}
+        new_table._applied_dynamic_options = applied_options
         return new_table
 
-    def _try_time_travel(self, options: Options) -> Optional[TableSchema]:
+    def _try_time_travel(self, options: Options, dynamic_option_keys: 
Set[str]) -> Optional[TableSchema]:
         """
         Try to resolve time travel options and return the corresponding schema.
 
@@ -573,10 +574,34 @@ class FileStoreTable(Table):
             )
             if snapshot is None:
                 return None
-            return 
self.schema_manager.get_schema(snapshot.schema_id).copy(new_options=options.to_map())
+            historical_schema = 
self.schema_manager.get_schema(snapshot.schema_id)
+            return 
historical_schema.copy(new_options=self._exclude_current_schema_field_options(
+                historical_schema, options, dynamic_option_keys))
         except Exception:
             return None
 
+    @staticmethod
+    def _exclude_current_schema_field_options(
+            historical_schema: TableSchema, options: Options, 
dynamic_option_keys: Set[str]) -> dict:
+        # Keep scan and runtime options, but restore column declarations to 
match historical fields.
+        historical_options = dict(options.to_map())
+        for key in (
+                CoreOptions.VECTOR_FIELD.key(),
+                CoreOptions.BLOB_FIELD.key(),
+                CoreOptions.BLOB_DESCRIPTOR_FIELD.key(),
+                CoreOptions.BLOB_VIEW_FIELD.key(),
+                # Restore the legacy key verbatim, not as a canonical 
descriptor option:
+                # Python intentionally ignores it when choosing the read 
layout.
+                'blob.stored-descriptor-fields'):
+            if key in dynamic_option_keys:
+                # Preserve explicit overrides and removals, including those 
from earlier copies.
+                continue
+            if key in historical_schema.options:
+                historical_options[key] = historical_schema.options[key]
+            else:
+                historical_options.pop(key, None)
+        return historical_options
+
     def _create_external_paths(self) -> List[str]:
         from urllib.parse import urlparse
 
diff --git 
a/paimon-python/pypaimon/tests/table/time_travel_schema_evolution_test.py 
b/paimon-python/pypaimon/tests/table/time_travel_schema_evolution_test.py
new file mode 100644
index 0000000000..73fe088c29
--- /dev/null
+++ b/paimon-python/pypaimon/tests/table/time_travel_schema_evolution_test.py
@@ -0,0 +1,187 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import pyarrow as pa
+import pytest
+
+from pypaimon import CatalogFactory, Schema
+from pypaimon.schema.data_types import ArrayType, AtomicType, DataField, 
PyarrowFieldParser
+from pypaimon.schema.schema_change import SchemaChange
+from pypaimon.table.row.blob import BlobDescriptor, BlobViewStruct
+
+
+class TestTimeTravelSchemaEvolution:
+
+    @pytest.fixture(autouse=True)
+    def setup(self, tmp_path):
+        self.root = tmp_path
+        self.catalog = CatalogFactory.create({'warehouse': str(tmp_path / 
'warehouse')})
+        self.catalog.create_database('test', False)
+
+    def _create(self, fields, options=None, name='test.t'):
+        table_options = {
+            'bucket': '-1',
+            'file.format': 'parquet',
+            'file.compression': 'none',
+            'row-tracking.enabled': 'true',
+            'data-evolution.enabled': 'true',
+        }
+        table_options.update(options or {})
+        self.catalog.create_table(name, Schema(fields=fields, 
options=table_options), False)
+        return self.catalog.get_table(name)
+
+    @staticmethod
+    def _write(table, row):
+        builder = table.new_batch_write_builder()
+        writer, commit = builder.new_write(), builder.new_commit()
+        try:
+            writer.write_arrow(pa.Table.from_pylist(
+                [row], 
schema=PyarrowFieldParser.from_paimon_schema(table.fields)))
+            messages = writer.prepare_commit()
+            commit.commit(messages)
+            return messages
+        finally:
+            writer.close()
+            commit.close()
+
+    @staticmethod
+    def _read(table):
+        builder = table.new_read_builder()
+        return 
builder.new_read().to_arrow(builder.new_scan().plan().splits()).to_pylist()
+
+    @pytest.mark.parametrize('key,directive,field_type', [
+        ('vector-field', '__VECTOR_FIELD;3', ArrayType(True, 
AtomicType('FLOAT'))),
+        ('blob-field', '__BLOB_FIELD', AtomicType('BYTES')),
+        ('blob-descriptor-field', '__BLOB_DESCRIPTOR_FIELD', 
AtomicType('BYTES')),
+        ('blob-view-field', '__BLOB_VIEW_FIELD', AtomicType('BYTES')),
+    ])
+    def test_added_field_options_follow_historical_schema(self, key, 
directive, field_type):
+        original = self._create([
+            DataField(0, 'id', AtomicType('INT')),
+            DataField(1, 'payload', field_type, directive),
+        ])
+        historical_options = dict(original.table_schema.options)
+        payload = [1.0, 0.0, 0.0] if key == 'vector-field' else None
+        self._write(original, {'id': 1, 'payload': payload})
+        original.create_tag('before_change', 1)
+        self.catalog.alter_table('test.t', [
+            SchemaChange.add_column('payload_v2', field_type, directive)])
+        current = self.catalog.get_table('test.t')
+        current_options = dict(current.table_schema.options)
+        assert current_options[key] == 'payload,payload_v2'
+        self._write(current, {'id': 2, 'payload': payload, 'payload_v2': 
payload})
+
+        historical = current.copy({'read.batch-size': '7'}).copy({
+            'scan.tag-name': 'before_change'})
+
+        assert historical.field_names == ['id', 'payload']
+        assert historical.table_schema.options[key] == 'payload'
+        assert historical.table_schema.options['read.batch-size'] == '7'
+        assert self._read(historical) == [{'id': 1, 'payload': payload}]
+        assert current.table_schema.options == current_options
+        assert 
current.schema_manager.get_schema(original.table_schema.id).options == 
historical_options
+
+    @pytest.mark.parametrize('key,directive', [
+        ('blob-descriptor-field', '__BLOB_DESCRIPTOR_FIELD'),
+        ('blob-view-field', '__BLOB_VIEW_FIELD'),
+    ])
+    def test_dropped_reference_field_reads_original_payload(self, key, 
directive):
+        payload = b'original-payload'
+        if key == 'blob-descriptor-field':
+            path = self.root / 'payload.bin'
+            path.write_bytes(payload)
+            reference = BlobDescriptor(str(path), 0, len(payload)).serialize()
+        else:
+            upstream = self._create([
+                DataField(0, 'id', AtomicType('INT')),
+                DataField(1, 'payload', AtomicType('BYTES'), '__BLOB_FIELD'),
+            ], name='test.upstream')
+            self._write(upstream, {'id': 1, 'payload': payload})
+            reference = BlobViewStruct('test.upstream', 1, 0).serialize()
+
+        original = self._create([
+            DataField(0, 'id', AtomicType('INT')),
+            DataField(1, 'payload', AtomicType('BYTES'), directive),
+        ])
+        self._write(original, {'id': 1, 'payload': reference})
+        assert self._read(original) == [{'id': 1, 'payload': payload}]
+        self.catalog.alter_table('test.t', 
[SchemaChange.drop_column('payload')])
+        current = self.catalog.get_table('test.t')
+        assert key not in current.table_schema.options
+        self._write(current, {'id': 2})
+
+        historical = current.copy({'scan.snapshot-id': '1'})
+
+        # Non-null references catch a missing decoder that would otherwise 
return raw bytes.
+        assert self._read(historical) == [{'id': 1, 'payload': payload}]
+        assert historical.table_schema.options[key] == 'payload'
+
+    def test_first_vector_field_is_absent_from_historical_options(self):
+        original = self._create([DataField(0, 'id', AtomicType('INT'))])
+        self._write(original, {'id': 1})
+        self.catalog.alter_table('test.t', [SchemaChange.add_column(
+            'embedding', ArrayType(True, AtomicType('FLOAT')), 
'__VECTOR_FIELD;3')])
+        current = self.catalog.get_table('test.t')
+
+        historical = current.copy({'scan.snapshot-id': '1'})
+
+        assert historical.field_names == ['id']
+        assert 'vector-field' not in historical.table_schema.options
+        assert self._read(historical) == [{'id': 1}]
+
+    def test_legacy_descriptor_option_keeps_historical_blob_layout(self):
+        legacy_key = 'blob.stored-descriptor-fields'
+        original = self._create([
+            DataField(0, 'id', AtomicType('INT')),
+            DataField(1, 'payload', AtomicType('BLOB')),
+        ], {legacy_key: 'payload'})
+        payload = b'legacy-blob-payload'
+        messages = self._write(original, {'id': 1, 'payload': payload})
+        assert any(f.file_name.endswith('.blob') for msg in messages for f in 
msg.new_files)
+        self.catalog.alter_table('test.t', [
+            SchemaChange.drop_column('payload'),
+            SchemaChange.add_column('reference', AtomicType('BYTES'), 
'__BLOB_DESCRIPTOR_FIELD'),
+        ])
+        current = self.catalog.get_table('test.t')
+
+        historical = current.copy({'scan.snapshot-id': '1'})
+
+        # Restore the original key, without turning it into an 
inline-descriptor layout switch.
+        assert historical.table_schema.options[legacy_key] == 'payload'
+        assert 'blob-descriptor-field' not in historical.table_schema.options
+        assert not historical.options.blob_descriptor_fields()
+        assert self._read(historical) == [{'id': 1, 'payload': payload}]
+
+    def test_explicit_field_overrides_survive_repeated_copies(self):
+        original = self._create([
+            DataField(0, 'id', AtomicType('INT')),
+            DataField(1, 'payload', ArrayType(True, AtomicType('FLOAT')), 
'__VECTOR_FIELD;3'),
+        ])
+        self._write(original, {'id': 1, 'payload': [1.0, 0.0, 0.0]})
+        self.catalog.alter_table('test.t', [SchemaChange.add_column(
+            'payload_v2', ArrayType(True, AtomicType('FLOAT')), 
'__VECTOR_FIELD;3')])
+        current = self.catalog.get_table('test.t')
+
+        for value in ('payload_v2', None):
+            historical = current.copy({'vector-field': 
value}).copy({'scan.snapshot-id': '1'})
+            repeated = historical.copy({'read.batch-size': '7'})
+            assert repeated.table_schema.options.get('vector-field') == value
+            assert repeated._applied_dynamic_options == {
+                'vector-field': value, 'scan.snapshot-id': '1', 
'read.batch-size': '7'}
+
+        historical = current.copy({'scan.snapshot-id': '1'})
+        assert historical._applied_dynamic_options == {'scan.snapshot-id': '1'}

Reply via email to