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 16f0d42161 [python] Add schema short-circuit to SplitRead and 
FileScanner read paths (#8217)
16f0d42161 is described below

commit 16f0d421616ee78f1767fcdddbcb51e527e0664c
Author: gjl <[email protected]>
AuthorDate: Sat Jun 20 13:09:34 2026 +0800

    [python] Add schema short-circuit to SplitRead and FileScanner read paths 
(#8217)
---
 .../pypaimon/manifest/manifest_file_manager.py     |  6 +-
 .../pypaimon/read/scanner/file_scanner.py          | 13 ++-
 paimon-python/pypaimon/read/split_read.py          | 17 ++--
 .../pypaimon/tests/partition_predicate_test.py     |  2 +-
 .../scanner/file_scanner_schema_fields_test.py     | 94 ++++++++++++++++++++++
 5 files changed, 121 insertions(+), 11 deletions(-)

diff --git a/paimon-python/pypaimon/manifest/manifest_file_manager.py 
b/paimon-python/pypaimon/manifest/manifest_file_manager.py
index 2d97516829..af710f94d6 100644
--- a/paimon-python/pypaimon/manifest/manifest_file_manager.py
+++ b/paimon-python/pypaimon/manifest/manifest_file_manager.py
@@ -124,7 +124,11 @@ class ManifestFileManager:
                 null_counts=key_dict['_NULL_COUNTS'],
             )
 
-            schema_fields = 
self.table.schema_manager.get_schema(file_dict['_SCHEMA_ID']).fields
+            schema_id = file_dict['_SCHEMA_ID']
+            if schema_id == self.table.table_schema.id:
+                schema_fields = self.table.table_schema.fields
+            else:
+                schema_fields = 
self.table.schema_manager.get_schema(schema_id).fields
             fields = self._get_value_stats_fields(file_dict, schema_fields)
             value_dict = dict(file_dict['_VALUE_STATS'])
             value_stats = SimpleStats(
diff --git a/paimon-python/pypaimon/read/scanner/file_scanner.py 
b/paimon-python/pypaimon/read/scanner/file_scanner.py
index 650bbe8ca4..e5f4e332a1 100755
--- a/paimon-python/pypaimon/read/scanner/file_scanner.py
+++ b/paimon-python/pypaimon/read/scanner/file_scanner.py
@@ -227,14 +227,19 @@ class FileScanner:
         # bucket set per ``total_buckets`` value.
         self._bucket_selector = self._init_bucket_selector()
 
-        def schema_fields_func(schema_id: int):
-            return self.table.schema_manager.get_schema(schema_id).fields
-
         self.simple_stats_evolutions = SimpleStatsEvolutions(
-            schema_fields_func,
+            self._schema_fields,
             self.table.table_schema.id
         )
 
+    def _schema_fields(self, schema_id: int):
+        """Resolve schema fields, short-circuiting current table schema id to 
avoid
+        filesystem access (REST catalog would get 403).
+        """
+        if schema_id == self.table.table_schema.id:
+            return self.table.table_schema.fields
+        return self.table.schema_manager.get_schema(schema_id).fields
+
     def _deletion_files_map(self, entries: List[ManifestEntry]) -> Dict[tuple, 
Dict[str, DeletionFile]]:
         if not self.deletion_vectors_enabled:
             return {}
diff --git a/paimon-python/pypaimon/read/split_read.py 
b/paimon-python/pypaimon/read/split_read.py
index 103d1e8510..d6a74353f3 100644
--- a/paimon-python/pypaimon/read/split_read.py
+++ b/paimon-python/pypaimon/read/split_read.py
@@ -167,6 +167,14 @@ class SplitRead(ABC):
     def _nested_path_by_name(self) -> Optional[Dict[str, List[str]]]:
         return self._cached_nested_path_by_name
 
+    def _resolve_schema(self, schema_id: int):
+        """Resolve schema, short-circuiting current table schema id to avoid
+        filesystem access (REST catalog would get 403).
+        """
+        if schema_id == self.table.table_schema.id:
+            return self.table.table_schema
+        return self.table.schema_manager.get_schema(schema_id)
+
     def _push_down_predicate(self) -> Optional[Predicate]:
         if self.predicate is None:
             return None
@@ -304,8 +312,7 @@ class SplitRead(ABC):
             if has_nested:
                 raise NotImplementedError(
                     "Nested-field projection is not supported on ROW files")
-            file_schema = self.table.schema_manager.get_schema(
-                file.schema_id)
+            file_schema = self._resolve_schema(file.schema_id)
             if file.write_cols:
                 field_map = {f.name: f for f in file_schema.fields}
                 row_full_fields = [field_map[n] for n in file.write_cols
@@ -389,7 +396,7 @@ class SplitRead(ABC):
         key = (schema_id, tuple(read_fields))
         if key not in self.schema_id_2_fields:
             nested_path_by_name = self._nested_path_by_name()
-            schema = self.table.schema_manager.get_schema(schema_id)
+            schema = self._resolve_schema(schema_id)
             schema_fields = (
                 SpecialFields.row_type_with_row_tracking(schema.fields)
                 if self.row_tracking_enabled else schema.fields
@@ -461,7 +468,7 @@ class SplitRead(ABC):
         nested-projection reads."""
         if self._nested_path_by_name() is not None:
             return None
-        file_schema = self.table.schema_manager.get_schema(file.schema_id)
+        file_schema = self._resolve_schema(file.schema_id)
         if file_schema is None:
             return None
         return self._final_data_fields_from(
@@ -1091,7 +1098,7 @@ class DataEvolutionSplitRead(SplitRead):
                 # For regular files without write_cols, derive field IDs from
                 # the file's schema version, not the current table schema.
                 # The file only contains columns from when it was written.
-                file_schema = 
self.table.schema_manager.get_schema(first_file.schema_id)
+                file_schema = self._resolve_schema(first_file.schema_id)
                 field_ids = [field.id for field in file_schema.fields]
                 field_ids.append(SpecialFields.ROW_ID.id)
                 field_ids.append(SpecialFields.SEQUENCE_NUMBER.id)
diff --git a/paimon-python/pypaimon/tests/partition_predicate_test.py 
b/paimon-python/pypaimon/tests/partition_predicate_test.py
index dbc50d0585..fdea8b17e6 100644
--- a/paimon-python/pypaimon/tests/partition_predicate_test.py
+++ b/paimon-python/pypaimon/tests/partition_predicate_test.py
@@ -65,7 +65,7 @@ def _mock_scanner_table():
     table.options.data_evolution_enabled.return_value = False
     table.options.deletion_vectors_enabled.return_value = False
     table.options.scan_manifest_parallelism.return_value = 1
-    table.table_schema = Mock(id=0)
+    table.table_schema = Mock(id=0, fields=TABLE_FIELDS)
     table.schema_manager = Mock()
     table.schema_manager.get_schema.return_value = Mock(fields=TABLE_FIELDS)
     return table
diff --git 
a/paimon-python/pypaimon/tests/scanner/file_scanner_schema_fields_test.py 
b/paimon-python/pypaimon/tests/scanner/file_scanner_schema_fields_test.py
new file mode 100644
index 0000000000..b75e7505d0
--- /dev/null
+++ b/paimon-python/pypaimon/tests/scanner/file_scanner_schema_fields_test.py
@@ -0,0 +1,94 @@
+# 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 unittest
+from unittest.mock import MagicMock
+
+from pypaimon.read.scanner.file_scanner import FileScanner
+from pypaimon.schema.data_types import AtomicType, DataField
+
+
+class FileScannerSchemaFieldsShortCircuitTest(unittest.TestCase):
+    """Test _schema_fields short-circuits current schema id (REST catalog 403 
fix)."""
+
+    def _make_scanner(self, current_schema_id, current_fields,
+                      historical_fields_map=None):
+        """Build a FileScanner bypassing __init__ to isolate _schema_fields."""
+        scanner = FileScanner.__new__(FileScanner)
+        scanner.table = MagicMock()
+        scanner.table.table_schema.id = current_schema_id
+        scanner.table.table_schema.fields = current_fields
+        scanner.table.schema_manager = MagicMock()
+
+        def get_schema(schema_id):
+            if historical_fields_map is None or schema_id not in 
historical_fields_map:
+                raise AssertionError(
+                    f"schema_manager.get_schema({schema_id}) was called "
+                    "but no historical schema was registered for it"
+                )
+            historical = MagicMock()
+            historical.fields = historical_fields_map[schema_id]
+            return historical
+
+        scanner.table.schema_manager.get_schema.side_effect = get_schema
+        return scanner
+
+    def test_short_circuits_current_schema_id(self):
+        """Current schema id returns in-memory fields without filesystem 
access."""
+        current_fields = [DataField(0, "a", AtomicType("INT"))]
+        scanner = self._make_scanner(
+            current_schema_id=5, current_fields=current_fields
+        )
+
+        result = scanner._schema_fields(5)
+
+        self.assertIs(result, current_fields)
+        scanner.table.schema_manager.get_schema.assert_not_called()
+
+    def test_delegates_for_historical_schema(self):
+        """Historical schema id delegates to schema_manager.get_schema()."""
+        current_fields = [DataField(0, "a", AtomicType("INT"))]
+        historical_fields = [
+            DataField(0, "a", AtomicType("INT")),
+            DataField(1, "b", AtomicType("STRING")),
+        ]
+        scanner = self._make_scanner(
+            current_schema_id=5,
+            current_fields=current_fields,
+            historical_fields_map={3: historical_fields},
+        )
+
+        result = scanner._schema_fields(3)
+
+        self.assertEqual(result, historical_fields)
+        scanner.table.schema_manager.get_schema.assert_called_once_with(3)
+
+    def test_short_circuit_works_for_zero_schema_id(self):
+        """Schema id == 0 still short-circuits (guards against truthiness 
bugs)."""
+        current_fields = [DataField(0, "x", AtomicType("INT"))]
+        scanner = self._make_scanner(
+            current_schema_id=0, current_fields=current_fields
+        )
+
+        result = scanner._schema_fields(0)
+
+        self.assertIs(result, current_fields)
+        scanner.table.schema_manager.get_schema.assert_not_called()
+
+
+if __name__ == "__main__":
+    unittest.main()

Reply via email to