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 935867d123 [python] Avoid redundant file metadata requests in Parquet 
reads (#10152)
935867d123 is described below

commit 935867d1230144e0ed9479adbe2455ec0bc7f3de
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Sep 24 10:48:02 2026 +0800

    [python] Avoid redundant file metadata requests in Parquet reads (#10152)
---
 .../filesystem/jindo_file_system_handler.py        | 43 +++++++++++++++++---
 .../pypaimon/read/reader/format_pyarrow_reader.py  | 35 +++++++++++-----
 paimon-python/pypaimon/read/split_read.py          |  3 +-
 .../pypaimon/tests/jindo_file_system_test.py       | 47 +++++++++++++++++++++-
 .../pypaimon/tests/parquet_metadata_cache_test.py  | 34 ++++++++++++++++
 5 files changed, 146 insertions(+), 16 deletions(-)

diff --git a/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py 
b/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
index 73ad09077f..eaad4a89be 100644
--- a/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
+++ b/paimon-python/pypaimon/filesystem/jindo_file_system_handler.py
@@ -16,6 +16,9 @@
 # under the License.
 
 import logging
+import os
+import threading
+from collections import OrderedDict
 
 import pyarrow as pa
 from pyarrow import PythonFile
@@ -50,6 +53,7 @@ _CASE_SENSITIVE_JINDO_CONFIG_KEYS = {
     OssOptions.OSS_ACCESS_KEY_SECRET.key().lower(): 
OssOptions.OSS_ACCESS_KEY_SECRET.key(),
     OssOptions.OSS_SECURITY_TOKEN.key().lower(): 
OssOptions.OSS_SECURITY_TOKEN.key(),
 }
+_KNOWN_FILE_SIZE_CACHE_MAX_ENTRIES = 4096
 
 
 def _jindo_config_value(value) -> str:
@@ -118,8 +122,9 @@ def create_jindo_oss_filesystem(root_uri: str, 
catalog_options: Options):
 
 
 class JindoInputFile:
-    def __init__(self, jindo_stream):
+    def __init__(self, jindo_stream, file_size=None):
         self._stream = jindo_stream
+        self._file_size = file_size
         self._closed = False
 
     @property
@@ -132,13 +137,18 @@ class JindoInputFile:
         if self.closed:
             raise ValueError("I/O operation on closed file")
         if nbytes is None or nbytes < 0:
-            return self._stream.read()
+            if self._file_size is None:
+                return self._stream.read()
+            nbytes = max(0, self._file_size - self.tell())
         return self._stream.read(nbytes)
 
     def seek(self, position: int, whence: int = 0):
         if self.closed:
             raise ValueError("I/O operation on closed file")
-        self._stream.seek(position, whence)
+        if whence == os.SEEK_END and self._file_size is not None:
+            position = self._file_size + position
+            whence = os.SEEK_SET
+        return self._stream.seek(position, whence)
 
     def tell(self) -> int:
         if self.closed:
@@ -210,10 +220,29 @@ class JindoFileSystemHandler(FileSystemHandler):
         self.logger = logging.getLogger(__name__)
         self.root_path = root_path
         self.properties = catalog_options
+        self._known_file_sizes = OrderedDict()
+        self._known_file_sizes_lock = threading.Lock()
 
         config = build_jindo_config(catalog_options)
         self._jindo_fs = jfs.connect(self.root_path, "root", config)
 
+    def register_file_size(self, path: str, file_size: int):
+        """Register an immutable file's size supplied by Paimon metadata."""
+        normalized = self._normalize_path(path)
+        with self._known_file_sizes_lock:
+            self._known_file_sizes[normalized] = file_size
+            self._known_file_sizes.move_to_end(normalized)
+            while (len(self._known_file_sizes)
+                   > _KNOWN_FILE_SIZE_CACHE_MAX_ENTRIES):
+                self._known_file_sizes.popitem(last=False)
+
+    def _known_file_size(self, path: str):
+        with self._known_file_sizes_lock:
+            file_size = self._known_file_sizes.get(path)
+            if file_size is not None:
+                self._known_file_sizes.move_to_end(path)
+            return file_size
+
     def __eq__(self, other):
         if isinstance(other, JindoFileSystemHandler):
             return self.root_path == other.root_path
@@ -325,12 +354,16 @@ class JindoFileSystemHandler(FileSystemHandler):
     def open_input_stream(self, path: str):
         normalized = self._normalize_path(path)
         jindo_stream = self._jindo_fs.open(normalized, "rb")
-        return PythonFile(JindoInputFile(jindo_stream), mode="r")
+        return PythonFile(
+            JindoInputFile(jindo_stream, self._known_file_size(normalized)),
+            mode="r")
 
     def open_input_file(self, path: str):
         normalized = self._normalize_path(path)
         jindo_stream = self._jindo_fs.open(normalized, "rb")
-        return PythonFile(JindoInputFile(jindo_stream), mode="r")
+        return PythonFile(
+            JindoInputFile(jindo_stream, self._known_file_size(normalized)),
+            mode="r")
 
     def open_output_stream(self, path: str, metadata):
         normalized = self._normalize_path(path)
diff --git a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py 
b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
index ed6dd37215..572f90a54b 100644
--- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
+++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
@@ -32,6 +32,7 @@ from pyarrow import RecordBatch
 from pypaimon.common.file_io import FileIO
 from pypaimon.common.options.config import CatalogOptions
 from pypaimon.common.options.core_options import CoreOptions
+from pypaimon.filesystem.pyarrow_file_io import _pyarrow_lt_7
 from pypaimon.data.map_shared_shredding import (
     assemble_normal_map_selected_keys,
     assemble_shared_shredding_map,
@@ -241,7 +242,7 @@ def _estimate_file_format_dataset_size(dataset, 
file_format: str) -> Optional[in
 
 
 def _estimate_file_format_cache_entry_size(
-        key: Tuple[Any, str, str],
+        key: Tuple[Any, ...],
         dataset,
         file_format: str) -> Optional[int]:
     metadata_size = _estimate_file_format_dataset_size(dataset, file_format)
@@ -255,9 +256,7 @@ def _estimate_file_format_cache_entry_size(
     visible_size = (
         metadata_size
         + sys.getsizeof(key)
-        + sys.getsizeof(key[0])
-        + sys.getsizeof(key[1])
-        + sys.getsizeof(key[2])
+        + sum(sys.getsizeof(part) for part in key)
         + sys.getsizeof(dataset)
         + sys.getsizeof((dataset, metadata_size))
         + _FILE_FORMAT_METADATA_CACHE_CONTAINER_OVERHEAD
@@ -274,22 +273,39 @@ def _file_format_metadata_cache_max_size(file_io: FileIO) 
-> int:
 
 
 def _file_format_dataset(file_io: FileIO, file_format: str, file_path: str,
-                         cache_max_size: int):
+                         cache_max_size: int,
+                         file_size: Optional[int] = None):
     file_path_for_pyarrow = file_io.to_filesystem_path(file_path)
     filesystem = file_io.filesystem
+    known_size = file_size if file_size is not None and file_size > 0 else None
+
+    # PyArrow's Python FileSystemHandler API opens files by path, so its
+    # FileInfo overload cannot forward file_size to handlers such as Jindo.
+    # Let capable handlers consume the size from Paimon's immutable file
+    # metadata before PyArrow opens the file.
+    handler = getattr(filesystem, "handler", None)
+    register_file_size = getattr(handler, "register_file_size", None)
+    if known_size is not None and register_file_size is not None:
+        register_file_size(file_path_for_pyarrow, known_size)
 
     def load():
         if file_format == 'parquet':
             parquet_format = ds.ParquetFileFormat()
+            fragment_options = {}
+            if known_size is not None and not _pyarrow_lt_7():
+                fragment_options["file_size"] = known_size
             fragment = parquet_format.make_fragment(
-                file_path_for_pyarrow, filesystem=filesystem)
+                file_path_for_pyarrow, filesystem=filesystem,
+                **fragment_options)
             # Reuse this fragment's footer for schema discovery and scanning.
             return ds.FileSystemDataset(
                 [fragment], fragment.physical_schema, parquet_format, 
filesystem)
         return ds.dataset(
             file_path_for_pyarrow, format=file_format, filesystem=filesystem)
 
-    key = (_FilesystemIdentity(filesystem), file_format, file_path_for_pyarrow)
+    key = (
+        _FilesystemIdentity(filesystem), file_format, file_path_for_pyarrow,
+        known_size)
     if cache_max_size <= 0:
         _reset_file_format_dataset_cache()
         return load()
@@ -347,7 +363,8 @@ class FormatPyArrowReader(RecordBatchReader):
                  predicate_field_names: Optional[Set[str]] = None,
                  row_indices: Optional[List[int]] = None,
                  row_ranges: Optional[List[Tuple[int, int]]] = None,
-                 row_group_cache: Optional[_DecodedRowGroupCache] = None):
+                 row_group_cache: Optional[_DecodedRowGroupCache] = None,
+                 file_size: Optional[int] = None):
         from pypaimon.filesystem.resolving_file_io import ResolvingFileIO
         if isinstance(file_io, ResolvingFileIO):
             file_io = file_io._get_fileio(file_path)
@@ -359,7 +376,7 @@ class FormatPyArrowReader(RecordBatchReader):
         self._row_group_cache_path = file_path_for_pyarrow
         cache_max_size = _file_format_metadata_cache_max_size(file_io)
         self.dataset = _file_format_dataset(
-            file_io, file_format, file_path, cache_max_size)
+            file_io, file_format, file_path, cache_max_size, file_size)
         self._range_slicer = None
         self._selected_parquet_row_groups = None
         self._exhausted = False
diff --git a/paimon-python/pypaimon/read/split_read.py 
b/paimon-python/pypaimon/read/split_read.py
index 1932852563..2a044c4045 100644
--- a/paimon-python/pypaimon/read/split_read.py
+++ b/paimon-python/pypaimon/read/split_read.py
@@ -456,7 +456,8 @@ class SplitRead(ABC):
                 nested_name_paths=ordered_nested_paths,
                 predicate_field_names=predicate_fields,
                 row_ranges=parquet_row_ranges,
-                row_group_cache=self._parquet_row_group_cache)
+                row_group_cache=self._parquet_row_group_cache,
+                file_size=file.file_size)
         elif file_format == CoreOptions.FILE_FORMAT_ROW:
             if has_nested:
                 raise NotImplementedError(
diff --git a/paimon-python/pypaimon/tests/jindo_file_system_test.py 
b/paimon-python/pypaimon/tests/jindo_file_system_test.py
index b828fd8b83..8c62113ca9 100644
--- a/paimon-python/pypaimon/tests/jindo_file_system_test.py
+++ b/paimon-python/pypaimon/tests/jindo_file_system_test.py
@@ -27,7 +27,11 @@ from pyarrow.fs import PyFileSystem
 from pypaimon.common.options import Options
 from pypaimon.common.options.config import OssOptions
 from pypaimon.filesystem import jindo_file_system_handler as jindo_module
-from pypaimon.filesystem.jindo_file_system_handler import 
JindoFileSystemHandler, JINDO_AVAILABLE
+from pypaimon.filesystem.jindo_file_system_handler import (
+    JindoFileSystemHandler,
+    JindoInputFile,
+    JINDO_AVAILABLE,
+)
 
 
 class _RecordingConfig:
@@ -38,6 +42,47 @@ class _RecordingConfig:
         self.values[key] = value
 
 
+class _RecordingInputStream:
+    def __init__(self, position=0):
+        self.position = position
+        self.seeks = []
+        self.reads = []
+        self.closed = False
+
+    def seek(self, position, whence=0):
+        self.seeks.append((position, whence))
+        self.position = position
+        return position
+
+    def tell(self):
+        return self.position
+
+    def read(self, size=None):
+        self.reads.append(size)
+        self.position += size or 0
+        return b"x" * (size or 0)
+
+    def close(self):
+        self.closed = True
+
+
+class JindoInputFileTest(unittest.TestCase):
+
+    def test_known_size_avoids_backend_seek_from_end(self):
+        stream = _RecordingInputStream()
+        input_file = JindoInputFile(stream, file_size=100)
+
+        self.assertEqual(95, input_file.seek(-5, os.SEEK_END))
+        self.assertEqual([(95, os.SEEK_SET)], stream.seeks)
+
+    def test_known_size_bounds_unlimited_read(self):
+        stream = _RecordingInputStream(position=40)
+        input_file = JindoInputFile(stream, file_size=100)
+
+        self.assertEqual(b"x" * 60, input_file.read())
+        self.assertEqual([60], stream.reads)
+
+
 class JindoConfigTest(unittest.TestCase):
 
     def test_forwards_native_options_to_connect(self):
diff --git a/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py 
b/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
index 12cead96fa..711e73191e 100644
--- a/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
+++ b/paimon-python/pypaimon/tests/parquet_metadata_cache_test.py
@@ -95,6 +95,9 @@ class _CountingFileSystemHandler(pafs.FSSpecHandler):
         self.calls.append("open_input_file")
         return super().open_input_file(path)
 
+    def register_file_size(self, path, file_size):
+        self.calls.append(("register_file_size", path, file_size))
+
 
 class FileFormatMetadataCacheTest(unittest.TestCase):
     def setUp(self):
@@ -232,6 +235,37 @@ class FileFormatMetadataCacheTest(unittest.TestCase):
                 self.assertEqual({"value": list(range(10))}, results[0])
                 self.assertEqual(results[0], results[1])
 
+    def test_forwards_known_file_size(self):
+        parquet_format = unittest.mock.Mock()
+        fragment = unittest.mock.Mock(physical_schema=pa.schema([]))
+        parquet_format.make_fragment.return_value = fragment
+        with patch.object(
+                reader_module.ds, "ParquetFileFormat",
+                return_value=parquet_format), patch.object(
+                    reader_module.ds, "FileSystemDataset",
+                    return_value=unittest.mock.sentinel.dataset):
+            dataset = reader_module._file_format_dataset(
+                self.file_io, "parquet", self.paths[0], 0, 123)
+        self.assertIs(unittest.mock.sentinel.dataset, dataset)
+        expected_options = (
+            {} if reader_module._pyarrow_lt_7()
+            else {"file_size": 123})
+        parquet_format.make_fragment.assert_called_once_with(
+            self.paths[0], filesystem=self.file_io.filesystem,
+            **expected_options)
+
+        handler = _CountingFileSystemHandler()
+        self.file_io.filesystem = pafs.PyFileSystem(handler)
+        file_size = os.path.getsize(self.paths[0])
+
+        dataset = reader_module._file_format_dataset(
+            self.file_io, "parquet", self.paths[0], 0, file_size)
+
+        self.assertEqual({"value": list(range(10))},
+                         dataset.to_table().to_pydict())
+        self.assertIn(
+            ("register_file_size", self.paths[0], file_size), handler.calls)
+
     def test_fragment_metadata_is_reused_without_io(self):
         handler = _CountingFileSystemHandler()
         self.file_io.filesystem = pafs.PyFileSystem(handler)

Reply via email to